我们做数据治理平台时,数据质量检核一开始只有批式方案——T+1 跑 Spark 任务,第二天早上出质量报告。这个方案对离线数仓够用,但业务侧很快提出了新诉求:
"订单数据进 Kafka 之后,能不能在进数仓之前就发现质量问题?等第二天才发现字段缺失、金额异常,下游报表已经错了一整天了。"
这个诉求指向的是流式数据质量检核——数据在流上,质量问题就要在流上发现。
我们基于 Flink + 自研规则引擎做了一套实时校验方案,上线运行半年,日均校验 2 亿+ 条记录,平均告警延迟 3 秒。这篇文章记录完整的架构设计、核心实现和踩过的坑。
很多人以为流式检核就是把批式的 SQL 规则搬到 Flink 上跑,实际上两者有本质差异。
维度 | 批式检核 | 流式检核 |
|---|---|---|
数据范围 | 全量历史数据 | 增量流式数据 |
检核粒度 | 表级(整表扫一遍) | 记录级 + 窗口级 |
时效性 | T+1 | 秒级~分钟级 |
状态管理 | 无状态(每次全量算) | 有状态(需要维护中间状态) |
规则类型 | 静态规则为主 | 静态 + 时序 + 跨流关联 |
告警方式 | 日报/周报 | 实时告警 |
关键差异在两点:
第一,状态管理。 批式检核每次都是全量扫描,不需要记住上次的结果。流式检核是增量的,很多规则需要跨记录、跨窗口的状态。比如"过去 5 分钟订单金额的均值是否异常",这个"过去 5 分钟的均值"就是一个需要持续维护的状态。
第二,告警语义。 批式检核的告警是"这张表有问题",流式检核的告警要精确到"哪一条记录、在什么时间、违反了哪条规则"。粒度更细,对告警的实时性和准确性要求也更高。
复制┌─────────────────────────────────────────────────────────┐
│ 数据源层 │
│ Kafka(订单/用户行为/日志) │ CDC(MySQL Binlog) │
├─────────────────────────────────────────────────────────┤
│ 检核引擎层 │
│ Flink 消费 → 规则匹配 → 状态聚合 → 违规判定 → 告警输出 │
├─────────────────────────────────────────────────────────┤
│ 规则管理层 │
│ 规则配置(MySQL) → 规则解析 → 规则下发(动态加载) │
├─────────────────────────────────────────────────────────┤
│ 结果与告警层 │
│ 违规明细(ClickHouse) → 告警(钉钉/邮件) → 质量看板 │
└─────────────────────────────────────────────────────────┘核心设计原则:规则配置和检核执行解耦。 规则存在 MySQL 里,Flink 任务启动时加载一次,运行期间通过广播流(Broadcast State)动态更新,新增/修改规则不需要重启 Flink 任务。
我们把流式质量规则分成四类,覆盖绝大多数场景:
规则类型 | 说明 | 示例 |
|---|---|---|
记录级规则 | 单条记录即可判定,无状态 | 字段非空、值域合法、正则匹配 |
窗口级规则 | 需要窗口内聚合后判定 | 5 分钟内订单量骤降、金额均值异常 |
跨流规则 | 需要关联多个数据流 | 订单的用户 ID 必须在用户流中存在 |
时序规则 | 需要跨窗口的状态 | 连续 3 个窗口质量下降 |
规则配置用 JSON 表达,兼顾可读性和可解析性:
json复制{
"rule_id": "order_amount_range",
"rule_name": "订单金额范围校验",
"rule_type": "RECORD",
"source": "kafka_orders",
"condition": {
"type": "AND",
"children": [
{"type": "NOT_NULL", "field": "order_amount"},
{"type": "RANGE", "field": "order_amount", "min": 0, "max": 100000}
]
},
"severity": "HIGH",
"enabled": true
}窗口级规则的配置:
json复制{
"rule_id": "order_volume_drop",
"rule_name": "订单量骤降检测",
"rule_type": "WINDOW",
"source": "kafka_orders",
"window": {"type": "TUMBLING", "size": "5m"},
"metric": {"type": "COUNT", "field": "order_id"},
"anomaly": {
"type": "THRESHOLD",
"operator": "LT",
"value": 1000,
"comment": "5分钟订单量低于1000触发告警"
},
"severity": "CRITICAL",
"enabled": true
}规则引擎的核心是把 JSON 规则解析成 Flink 可以执行的算子。记录级规则最简单,直接映射到 filter 或 map 算子:
java复制// 记录级规则:解析 condition 树,生成过滤函数
public class RecordRuleEvaluator implements Serializable {
private final ConditionNode condition;
public RecordRuleEvaluator(ConditionNode condition) {
this.condition = condition;
}
// 判断单条记录是否违规
public boolean isViolated(Map<String, Object> record) {
return !evaluate(condition, record);
}
private boolean evaluate(ConditionNode node, Map<String, Object> record) {
switch (node.getType()) {
case "NOT_NULL":
return record.get(node.getField()) != null;
case "RANGE":
Object val = record.get(node.getField());
if (val == null) return false;
double d = ((Number) val).doubleValue();
return d >= node.getMin() && d <= node.getMax();
case "REGEX":
Object s = record.get(node.getField());
return s != null && s.toString().matches(node.getPattern());
case "AND":
return node.getChildren().stream().allMatch(c -> evaluate(c, record));
case "OR":
return node.getChildren().stream().anyMatch(c -> evaluate(c, record));
default:
return true;
}
}
}记录级规则最简单,Flink 里就是一个 flatMap + 侧输出流:
java复制DataStream<OrderEvent> source = env
.addSource(createKafkaSource())
.map(this::parseJson)
.filter(Objects::nonNull);
// 主输出:合规数据,继续往下游走
// 侧输出:违规数据,进入告警流
OutputTag<ViolationRecord> violationTag =
new OutputTag<ViolationRecord>("violation") {};
SingleOutputStreamOperator<OrderEvent> checked = source
.process(new ProcessFunction<OrderEvent, OrderEvent>() {
@Override
public void processElement(OrderEvent value, Context ctx,
Collector<OrderEvent> out) {
List<ViolationRecord> violations = new ArrayList<>();
for (RecordRule rule : recordRules) {
if (rule.getEvaluator().isViolated(value.toMap())) {
violations.add(new ViolationRecord(rule, value));
}
}
if (violations.isEmpty()) {
out.collect(value); // 合规,放行
} else {
for (ViolationRecord v : violations) {
ctx.output(violationTag, v); // 违规,进告警流
}
}
}
});
// 告警流:写 ClickHouse + 推送告警
DataStream<ViolationRecord> violationStream = checked.getSideOutput(violationTag);
violationStream.addSink(new ClickHouseSink());
violationStream.addSink(new AlertSink());窗口级规则的核心是窗口聚合 + 异常判定。用 Flink 的窗口算子做聚合,然后对聚合结果做异常检测:
java复制// 窗口级规则:5分钟滚动窗口,统计订单量,判断是否低于阈值
DataStream<OrderEvent> source = ...;
source
.keyBy(e -> e.getRuleId()) // 按规则分组
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new OrderCountAggregator(), new WindowResultFunction())
.process(new AnomalyDetector(windowRules))
.addSink(new AlertSink());聚合函数:
java复制public class OrderCountAggregator
implements AggregateFunction<OrderEvent, Long, Long> {
@Override
public Long createAccumulator() { return 0L; }
@Override
public Long add(OrderEvent value, Long acc) { return acc + 1; }
@Override
public Long getResult(Long acc) { return acc; }
@Override
public Long merge(Long a, Long b) { return a + b; }
}异常检测函数:
java复制public class AnomalyDetector
extends ProcessWindowFunction<Long, Alert, String, TimeWindow> {
private final List<WindowRule> rules;
@Override
public void process(String ruleId, Context context,
Iterable<Long> counts, Collector<Alert> out) {
long count = counts.iterator().next();
WindowRule rule = findRule(ruleId);
if (rule.getAnomaly().isTriggered(count)) {
out.collect(new Alert(ruleId, count, context.window()));
}
}
}跨流规则(比如"订单的用户 ID 必须在用户流中存在")需要关联两个数据流。用 Flink 的 connect + CoProcessFunction 实现:
java复制// 用户流:维护一个用户 ID 的状态集合
DataStream<UserEvent> userStream = ...;
DataStream<OrderEvent> orderStream = ...;
orderStream
.connect(userStream)
.keyBy(o -> o.getUserId(), u -> u.getUserId())
.process(new CrossStreamValidator())
.addSink(new AlertSink());java复制public class CrossStreamValidator
extends CoProcessFunction<OrderEvent, UserEvent, Alert> {
// 用户 ID 状态:存最近 N 天的用户集合
private ValueState<Boolean> userExists;
@Override
public void open(Configuration parameters) {
ValueStateDescriptor<Boolean> desc =
new ValueStateDescriptor<>("userExists", Boolean.class);
userExists = getRuntimeContext().getState(desc);
}
// 处理订单流:检查用户是否存在
@Override
public void processElement1(OrderEvent order, Context ctx, Collector<Alert> out) {
Boolean exists = userExists.value();
if (exists == null || !exists) {
out.collect(new Alert("user_not_found", order.getUserId()));
}
}
// 处理用户流:更新用户存在状态
@Override
public void processElement2(UserEvent user, Context ctx, Collector<Alert> out) {
userExists.update(true);
}
}规则需要动态更新(新增规则、修改阈值、下线规则),不能让业务方每次改规则都重启 Flink 任务。用 Broadcast State 实现:
java复制// 规则广播流:从 MySQL 定期拉取规则变更,广播给所有并行子任务
BroadcastStream<RuleUpdate> ruleBroadcast = env
.addSource(new RuleSource()) // 定期查 MySQL,发现变更就发一条
.broadcast(ruleStateDescriptor);
DataStream<OrderEvent> source = ...;
source
.connect(ruleBroadcast)
.process(new BroadcastRuleProcessor())
.addSink(new AlertSink());java复制public class BroadcastRuleProcessor
extends BroadcastProcessFunction<OrderEvent, RuleUpdate, Alert> {
private final MapStateDescriptor<String, Rule> ruleStateDescriptor =
new MapStateDescriptor<>("rules", String.class, Rule.class);
// 处理数据流:用当前广播状态中的规则做检核
@Override
public void processElement(OrderEvent value, ReadOnlyContext ctx,
Collector<Alert> out) {
ReadOnlyBroadcastState<String, Rule> rules =
ctx.getBroadcastState(ruleStateDescriptor);
for (Map.Entry<String, Rule> entry : rules.immutableEntries()) {
Rule rule = entry.getValue();
if (rule.getEvaluator().isViolated(value.toMap())) {
out.collect(new Alert(rule, value));
}
}
}
// 处理广播流:更新规则状态
@Override
public void processBroadcastElement(RuleUpdate update, Context ctx,
Collector<Alert> out) {
BroadcastState<String, Rule> rules =
ctx.getBroadcastState(ruleStateDescriptor);
if (update.getAction() == RuleUpdate.Action.ADD) {
rules.put(update.getRule().getId(), update.getRule());
} else if (update.getAction() == RuleUpdate.Action.DELETE) {
rules.remove(update.getRule().getId());
}
}
}跨流检核的 userExists 状态,如果用户流是持续增长的(每天新增几十万用户),状态会无限膨胀,最终 OOM。
解决方案:状态加 TTL。 给状态设置过期时间,超时自动清理:
java复制StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Time.days(30))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
ValueStateDescriptor<Boolean> desc =
new ValueStateDescriptor<>("userExists", Boolean.class);
desc.enableTimeToLive(ttlConfig);窗口级规则如果阈值设置不当,会触发告警风暴。比如"订单量低于 1000 告警",如果业务本身就有低谷期(凌晨订单量本来就低),那每天凌晨都会误报。
解决方案:引入基线对比。 不是跟固定阈值比,而是跟历史同期比:
java复制// 不是"低于1000告警",而是"低于过去7天同时段均值的50%告警"
double baseline = getHistoricalBaseline(windowStart, ruleId); // 查历史基线
double threshold = baseline * 0.5;
if (count < threshold) {
out.collect(new Alert(ruleId, count, baseline, window));
}Kafka 数据乱序是常态,如果直接按处理时间(Processing Time)开窗,窗口统计会不准。
解决方案:用事件时间(Event Time)+ Watermark。
java复制source
.assignTimestampsAndWatermarks(
WatermarkStrategy
.<OrderEvent>forBoundedOutOfOrderness(Duration.ofSeconds(30))
.withTimestampAssigner((e, ts) -> e.getEventTime())
)
.keyBy(e -> e.getRuleId())
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
...同一个质量问题,可能因为 Flink 的 checkpoint 重放、或者规则在多个窗口重复触发,导致重复告警。
解决方案:告警去重。 在告警输出前做一次去重,按"规则 ID + 时间窗口 + 关键字段"做幂等:
java复制public class AlertDeduplicator {
// 用 Redis 做去重,key = ruleId + windowStart,TTL = 10分钟
public boolean shouldAlert(String ruleId, long windowStart) {
String key = ruleId + ":" + windowStart;
return redis.setnx(key, "1", 600); // 设置成功返回 true,表示首次告警
}
}上线运行半年后的实际数据:
指标 | 数值 |
|---|---|
日均校验记录数 | 2.3 亿条 |
峰值吞吐 | 8500 条/秒 |
平均告警延迟 | 3.2 秒 |
P99 告警延迟 | 8.5 秒 |
在线规则数 | 156 条 |
Flink 并行度 | 24(8 个 TaskManager × 3 slot) |
状态后端 | RocksDB(增量 checkpoint) |
关键调优点:
流式数据质量检核的核心,不是把批式规则搬到流上,而是重新设计规则模型和状态管理。
几个关键设计决策:
流式检核不是批式检核的替代,而是补充。批式做全量兜底(T+1 出完整质量报告),流式做实时预警(秒级发现异常),两者结合才是完整的数据质量体系。
本文基于作者团队在生产环境的流式数据质量检核实践。技术栈为 Flink 1.16 + Kafka + ClickHouse,规则引擎为自研。不同版本 API 可能有差异,请以官方文档为准。
原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。
如有侵权,请联系 cloudcommunity@tencent.com 删除。