我要提问
ARTICLE / 004 · 大数据开发

原创文章

资深开发者执笔的深度技术长文,从原理到工程落地,逐层拆解。

Flink 实时数仓从采集到可视化的全链路

Flink 实时数仓从采集到可视化的全链路

离线数仓用 T+1 的批处理回答"昨天发生了什么",实时数仓用流处理回答"现在正在发生什么"。Flink 凭借精准的流式语义、丰富的状态管理与 SQL 抽象,成为实时数仓的事实引擎。但"实时数仓"远不止跑一个 Flink 任务——它是一条从采集、清洗、聚合、落库到可视化的完整链路,每一环都有坑。本文以一个电商实时大屏为例,走完全链路。

一、实时数仓的分层架构

和离线数仓类似,实时数仓也分层:ODS(原始日志)→ DWD(明细宽表)→ DWS(轻度聚合)→ ADS(应用层指标)。区别在于每一层都是流,而不是表。数据在 Kafka 中以 topic 形式分层流转,Flink 任务负责层间加工,最终结果写入 ClickHouse / Doris / Redis 供查询。

实时数仓的灵魂是"分层 + 物化"。分层让复杂逻辑解耦,物化让查询变快。流的物化就是把中间结果写回 Kafka,让下游可以复用,避免每个应用都从 ODS 重算。
  • ODS 层:原始业务日志、Binlog 直接入 Kafka,保留全量字段,不做加工。
  • DWD 层:清洗、脱敏、维度补全后的明细宽表,是后续所有聚合的基础。
  • DWS 层:按主题轻度聚合(如按店铺、按分钟聚合的成交额)。
  • ADS 层:面向应用的最终指标,直接喂给大屏或报表。

二、数据采集:Binlog 与日志双通道

采集层要解决"把业务库与日志系统的变更实时送进 Kafka"。数据库变更走 CDC(Change Data Capture):用 Canal/Debezium 订阅 MySQL Binlog,把 INSERT/UPDATE/DELETE 转成 Kafka 消息。日志走 Filebeat/Fluent 采集应用日志文件。两条通道都写入 ODS topic,统一用 Flink 消费。

-- DWD 层:从 ODS 订单流清洗出有效订单
CREATE TABLE dwd_order (
  order_id   STRING,
  user_id    BIGINT,
  shop_id    BIGINT,
  amount     DECIMAL(18,2),
  status     STRING,
  event_time TIMESTAMP(3),
  WATERMARK FOR event_time AS event_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'dwd_order',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'json'
);

三、维表关联:把流变成宽表

明细流里通常只有外键(user_id、shop_id),要补全维度信息(用户名、店铺名、类目)需要关联维表。维表存在 MySQL/HBase/Redis 中。Flink 提供 Lookup Join,对每条流数据去维表查一次。但 Lookup Join 是同步阻塞的,高 QPS 下会成为瓶颈。

3.1 维表关联的三种方式

  • Lookup Join:每条数据查一次外部维表,延迟低但压力大,适合维表大、更新频繁的场景,需配合缓存。
  • Regular Join:把维表也作为流接入 Flink,用状态保存全量维表,自动响应维表变更。适合维表较小(百万级以内)的场景。
  • Temporal Table Join:按事件时间关联维表的版本快照,保证可重放、结果确定,适合对一致性要求高的场景。
-- Lookup Join:订单流关联店铺维表
SELECT o.order_id, o.amount, s.shop_name, s.category
FROM dwd_order AS o
LEFT JOIN shop_dim FOR SYSTEM_TIME AS OF o.event_time AS s
ON o.shop_id = s.shop_id;

Lookup Join 加缓存是工程标配。Flink 支持 lookup.cache 配置:PARTIAL 缓存命中则不查外部,未命中才查并回填缓存,可设 TTL。缓存能将外部查询 QPS 降低一个数量级。

四、聚合与水位线:DWS 层的核心

DWS 层做时间窗口聚合:每分钟、每小时的成交额、UV 等。Flink SQL 用 TUMBLE/HOP/SESSION 窗口函数表达。关键在于水位线(Watermark)——它决定窗口何时关闭、何时输出结果。水位线 = 当前最大事件时间 − 允许延迟,水位线越过窗口结束时间,窗口就触发计算。

-- DWS 层:按店铺每分钟聚合成交额与订单数
SELECT
  shop_id,
  TUMBLE_START(event_time, INTERVAL '1' MINUTE) AS win_start,
  SUM(amount)   AS gmv,
  COUNT(order_id) AS order_cnt
FROM dwd_order
GROUP BY shop_id, TUMBLE(event_time, INTERVAL '1' MINUTE);
水位线是实时数仓的"时间守门人"。设得太紧,乱序数据被丢弃,指标偏低;设得太松,窗口迟迟不关,实时性变差。mrgr.cn 的经验是:日志类数据用 5-10 秒延迟,Binlog 类数据用 1-3 秒,再配合 allowedLateness 兜底迟到数据。

4.1 状态 TTL 与膨胀治理

窗口聚合、Regular Join 都会产生状态。如果不设 TTL,状态会无限膨胀最终撑爆 RocksDB。必须为状态设置生存周期:table.exec.state.ttl。TTL 要大于窗口长度 + 允许延迟,否则窗口还没关闭状态就被清掉了。但也不能太大,否则状态膨胀。这是实时数仓最常踩的坑之一。

  • TTL < 窗口长度:窗口状态被提前清理,结果丢数据。
  • TTL 过大:状态膨胀,Checkpoint 变慢,OOM 风险。
  • 经验值:TTL = 窗口长度 + allowedLateness + 安全余量(如 1 小时)。

五、落库与查询:ClickHouse 的取舍

ADS 层结果要写入支持高并发查询的引擎。ClickHouse 是实时大屏的常见选择:列存 + 向量化执行,单表聚合查询极快。但 ClickHouse 不擅长高频写入与精确更新,所以 Flink 写入时要"批量 + 异步",并选用 MergeTree 系列表引擎。对需要更新的指标(如订单状态变更),用 ReplacingMergeTree 或 CollapsingMergeTree。

-- Flink 写 ClickHouse(JDBC 连接器,批量攒批)
INSERT INTO ads_shop_minute (shop_id, win_start, gmv, order_cnt)
SELECT shop_id, win_start, gmv, order_cnt FROM dws_shop_minute;

-- ClickHouse 建表:按天分区,按 shop_id 排序
CREATE TABLE ads_shop_minute (
  shop_id    UInt64,
  win_start  DateTime,
  gmv        Decimal(18,2),
  order_cnt  UInt32
) ENGINE = ReplacingMergeTree()
PARTITION BY toYYYYMMDD(win_start)
ORDER BY (shop_id, win_start);

六、Exactly-Once 与可视化保障

实时数仓要做到端到端 Exactly-Once:Flink Checkpoint + Kafka 事务 + 下游幂等写入。Flink 内部通过 Checkpoint 保证状态一致;Sink 端通过两阶段提交(Kafka 事务 sink)或幂等写入(ClickHouse 用 ReplacingMergeTree 去重)保证不重不丢。可视化层用 Grafana / 自研大屏直查 ClickHouse,刷新间隔 5-10 秒。

  • Flink Checkpoint 间隔:秒级到分钟级,间隔越短恢复越快但开销越大。
  • Sink 幂等:用主键去重,避免重启后重复写入放大指标。
  • 大屏查询:预聚合 + 物化视图,避免实时查大表。
  • 监控:消费延迟、水位线延迟、Checkpoint 失败率必须告警。

结语

实时数仓的难点不在 Flink SQL 本身,而在于把"采集、维表、水位线、状态、落库"串成一条稳定可运维的链路。在 mrgr.cn 的实践中,分层物化让下游复用成为可能,水位线与 TTL 调优是日常运维的主线,Exactly-Once 是数据准确性的底线。下一篇我们会拆解 Kubernetes 故障排查,欢迎持续关注 mrgr.cn 大数据组。

返回文章列表