首页
学习
活动
专区
圈层
工具
发布
社区首页 >专栏 >流式数据质量检核怎么做:Flink + 规则引擎的实时校验实践

流式数据质量检核怎么做:Flink + 规则引擎的实时校验实践

原创
作者头像
KylenReview
发布2026-09-04 10:16:41
发布2026-09-04 10:16:41
290
举报
文章被收录于专栏:数据治理数据治理

我们做数据治理平台时,数据质量检核一开始只有批式方案——T+1 跑 Spark 任务,第二天早上出质量报告。这个方案对离线数仓够用,但业务侧很快提出了新诉求:

"订单数据进 Kafka 之后,能不能在进数仓之前就发现质量问题?等第二天才发现字段缺失、金额异常,下游报表已经错了一整天了。"

这个诉求指向的是流式数据质量检核——数据在流上,质量问题就要在流上发现。

我们基于 Flink + 自研规则引擎做了一套实时校验方案,上线运行半年,日均校验 2 亿+ 条记录,平均告警延迟 3 秒。这篇文章记录完整的架构设计、核心实现和踩过的坑。


一、流式检核和批式检核,到底差在哪

很多人以为流式检核就是把批式的 SQL 规则搬到 Flink 上跑,实际上两者有本质差异。

维度

批式检核

流式检核

数据范围

全量历史数据

增量流式数据

检核粒度

表级(整表扫一遍)

记录级 + 窗口级

时效性

T+1

秒级~分钟级

状态管理

无状态(每次全量算)

有状态(需要维护中间状态)

规则类型

静态规则为主

静态 + 时序 + 跨流关联

告警方式

日报/周报

实时告警

关键差异在两点:

第一,状态管理。 批式检核每次都是全量扫描,不需要记住上次的结果。流式检核是增量的,很多规则需要跨记录、跨窗口的状态。比如"过去 5 分钟订单金额的均值是否异常",这个"过去 5 分钟的均值"就是一个需要持续维护的状态。

第二,告警语义。 批式检核的告警是"这张表有问题",流式检核的告警要精确到"哪一条记录、在什么时间、违反了哪条规则"。粒度更细,对告警的实时性和准确性要求也更高。


二、整体架构

代码语言:javascript
复制
复制┌─────────────────────────────────────────────────────────┐
│                      数据源层                            │
│   Kafka(订单/用户行为/日志)  │  CDC(MySQL Binlog)    │
├─────────────────────────────────────────────────────────┤
│                      检核引擎层                          │
│  Flink 消费 → 规则匹配 → 状态聚合 → 违规判定 → 告警输出  │
├─────────────────────────────────────────────────────────┤
│                      规则管理层                          │
│  规则配置(MySQL) → 规则解析 → 规则下发(动态加载)      │
├─────────────────────────────────────────────────────────┤
│                      结果与告警层                        │
│  违规明细(ClickHouse) → 告警(钉钉/邮件) → 质量看板    │
└─────────────────────────────────────────────────────────┘

核心设计原则:规则配置和检核执行解耦。 规则存在 MySQL 里,Flink 任务启动时加载一次,运行期间通过广播流(Broadcast State)动态更新,新增/修改规则不需要重启 Flink 任务。


三、规则引擎设计

3.1 规则分类

我们把流式质量规则分成四类,覆盖绝大多数场景:

规则类型

说明

示例

记录级规则

单条记录即可判定,无状态

字段非空、值域合法、正则匹配

窗口级规则

需要窗口内聚合后判定

5 分钟内订单量骤降、金额均值异常

跨流规则

需要关联多个数据流

订单的用户 ID 必须在用户流中存在

时序规则

需要跨窗口的状态

连续 3 个窗口质量下降

3.2 规则 DSL 设计

规则配置用 JSON 表达,兼顾可读性和可解析性:

代码语言:javascript
复制
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
}

窗口级规则的配置:

代码语言:javascript
复制
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
}

3.3 规则解析与执行

规则引擎的核心是把 JSON 规则解析成 Flink 可以执行的算子。记录级规则最简单,直接映射到 filtermap 算子:

代码语言:javascript
复制
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;
        }
    }
}

四、核心实现

4.1 记录级检核:无状态,直接过滤

记录级规则最简单,Flink 里就是一个 flatMap + 侧输出流:

代码语言:javascript
复制
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());

4.2 窗口级检核:有状态,需要聚合

窗口级规则的核心是窗口聚合 + 异常判定。用 Flink 的窗口算子做聚合,然后对聚合结果做异常检测:

代码语言:javascript
复制
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());

聚合函数:

代码语言:javascript
复制
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; }
}

异常检测函数:

代码语言:javascript
复制
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()));
        }
    }
}

4.3 跨流检核:双流 JOIN

跨流规则(比如"订单的用户 ID 必须在用户流中存在")需要关联两个数据流。用 Flink 的 connect + CoProcessFunction 实现:

代码语言:javascript
复制
java复制// 用户流:维护一个用户 ID 的状态集合
DataStream<UserEvent> userStream = ...;
DataStream<OrderEvent> orderStream = ...;

orderStream
    .connect(userStream)
    .keyBy(o -> o.getUserId(), u -> u.getUserId())
    .process(new CrossStreamValidator())
    .addSink(new AlertSink());
代码语言:javascript
复制
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);
    }
}

4.4 规则动态下发:广播流

规则需要动态更新(新增规则、修改阈值、下线规则),不能让业务方每次改规则都重启 Flink 任务。用 Broadcast State 实现:

代码语言:javascript
复制
java复制// 规则广播流:从 MySQL 定期拉取规则变更,广播给所有并行子任务
BroadcastStream<RuleUpdate> ruleBroadcast = env
    .addSource(new RuleSource())  // 定期查 MySQL,发现变更就发一条
    .broadcast(ruleStateDescriptor);

DataStream<OrderEvent> source = ...;

source
    .connect(ruleBroadcast)
    .process(new BroadcastRuleProcessor())
    .addSink(new AlertSink());
代码语言:javascript
复制
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());
        }
    }
}

五、踩过的坑

坑一:状态无限增长导致 OOM

跨流检核的 userExists 状态,如果用户流是持续增长的(每天新增几十万用户),状态会无限膨胀,最终 OOM。

解决方案:状态加 TTL。 给状态设置过期时间,超时自动清理:

代码语言:javascript
复制
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 告警",如果业务本身就有低谷期(凌晨订单量本来就低),那每天凌晨都会误报。

解决方案:引入基线对比。 不是跟固定阈值比,而是跟历史同期比:

代码语言:javascript
复制
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。

代码语言:javascript
复制
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 + 时间窗口 + 关键字段"做幂等:

代码语言:javascript
复制
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)

关键调优点:

  • RocksDB 状态后端:状态量大,用 RocksDB 而不是内存状态,避免 OOM。
  • 增量 checkpoint:开启增量 checkpoint,checkpoint 耗时从 40 秒降到 5 秒。
  • 背压处理:告警写入 ClickHouse 偶发慢查询导致背压,通过增加 sink 并行度和批量写入解决。

七、总结

流式数据质量检核的核心,不是把批式规则搬到流上,而是重新设计规则模型和状态管理

几个关键设计决策:

  1. 规则分类:记录级(无状态)、窗口级(有状态)、跨流(双流 JOIN)、时序(跨窗口),不同类别用不同的 Flink 算子实现。
  2. 规则动态下发:用 Broadcast State 实现规则热更新,业务方改规则不用重启任务。
  3. 状态管理:跨流检核的状态必须加 TTL,否则无限增长 OOM。
  4. 告警质量:窗口级规则要引入基线对比避免误报,告警要做去重避免重复推送。

流式检核不是批式检核的替代,而是补充。批式做全量兜底(T+1 出完整质量报告),流式做实时预警(秒级发现异常),两者结合才是完整的数据质量体系。


本文基于作者团队在生产环境的流式数据质量检核实践。技术栈为 Flink 1.16 + Kafka + ClickHouse,规则引擎为自研。不同版本 API 可能有差异,请以官方文档为准。

原创声明:本文系作者授权腾讯云开发者社区发表,未经许可,不得转载。

如有侵权,请联系 cloudcommunity@tencent.com 删除。

目录
  • 一、流式检核和批式检核,到底差在哪
  • 二、整体架构
  • 三、规则引擎设计
    • 3.1 规则分类
    • 3.2 规则 DSL 设计
    • 3.3 规则解析与执行
  • 四、核心实现
    • 4.1 记录级检核:无状态,直接过滤
    • 4.2 窗口级检核:有状态,需要聚合
    • 4.3 跨流检核:双流 JOIN
    • 4.4 规则动态下发:广播流
  • 五、踩过的坑
    • 坑一:状态无限增长导致 OOM
    • 坑二:窗口级规则的告警风暴
    • 坑三:乱序数据导致窗口计算不准
    • 坑四:告警去重缺失导致重复推送
  • 六、性能数据
  • 七、总结
问题归档专栏文章快讯文章归档关键词归档开发者手册归档开发者手册 Section 归档