离线数仓用 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 大数据组。