引言
过去十年,实时与批处理被当成两条路:批处理算得全但慢,流处理算得快但"不够严谨"。为了两头兼顾,很多团队搭建了 Lambda 架构——一套流式链路出实时结果,一套批式链路出最终结果,再用一个"合并层"缝合。代价是两套代码、两套口径、两倍运维。流批一体的目标,是用一套引擎、一套语义、一套数据同时覆盖批与流,让口径天然一致、代码只写一遍。
流批一体的本质:不是消灭批或消灭流,而是让"同一份数据、同一套计算逻辑"既能跑批也能跑流——差异只在于触发方式与时间语义。
本文从 Lambda/Kappa 演进讲起,拆解流批同语义、实时对账、统一计算层选型与湖仓底座,最后给出务实的迁移路径。
一、从批处理到流处理的演进
1.1 流批分离的历史
# 传统架构
# 批: 日/小时级 ETL → 数仓 → 报表(T+1)
# 流: Kafka → 流引擎 → 实时看板(秒级)
# 痛点: 两套代码、两套口径、两倍运维
# 最痛: 实时结果与批结果"对不上"
1.2 Lambda 架构
Lambda 用"并行两条链路 + 合并"解决实时与准确的矛盾:
# Lambda 架构
# 流式层(Batch Layer): 定期全量重算 → 准确
# 实时层(Speed Layer): 增量计算 → 快速但不完整
# 合并层(Serving Layer): 合并两条结果
# 代价: 两套代码逻辑必须一致(很难) + 合并复杂度
1.3 Kappa 架构
Kappa 只保留一条流式链路,用**重放(Replay)**实现批的语义:需要"批结果"时,从历史数据重放流作业即可。
# Kappa 架构
# 所有数据进 Kafka(保留期足够长)
# 只有一条流作业
# 回算/补数 = 从指定 offset 重放
# 优势: 单一代码、单一口径
# 挑战: 大保留期存储成本、历史重放性能
二、Lambda 与 Kappa 选型
2.1 两者对比
| 维度 | Lambda | Kappa |
|---|---|---|
| 计算链路 | 流 + 批 两条 | 一条流式 |
| 代码一致性 | 两套逻辑难对齐 | 天然一致 |
| 回算能力 | 批链负责 | 重放实现 |
| 复杂度 | 高(合并层) | 低 |
| 适用 | 历史包袱重 | 新建/可重放 |
2.2 什么时候仍需要 Lambda
- 历史数据量巨大,无法全部常驻 Kafka/事件流。
- 流式引擎重放海量历史性能不可接受。
- 已有成熟的批式数仓,迁移成本过高。
2.3 走向一体的判断
# 判断条件
# [ ] 业务是否同时需要秒级 + 精确结果
# [ ] 数据能否在事件流中保留足够久
# [ ] 回算是否只需覆盖近期窗口
# [ ] 团队是否能驾驭统一引擎
# 大多数新项目: 用流批一体(湖仓 + 统一引擎)更划算
三、流批同语义
3.1 为什么需要同语义
流批"对不上"的根因,是同一指标用了两套语义:批用 SQL 算、流用代码算,窗口、去重、迟到处理规则都不同。同语义要求:批 SQL 和流 SQL 是同一份逻辑。
3.2 Flink SQL / Table API 的流批统一
Flink 的 Table/SQL API 让同一段 SQL 既能在批模式跑,也能在流模式跑:
-- 同一段 SQL: 批模式与流模式都可执行
SELECT
user_id,
COUNT(DISTINCT order_id) AS order_cnt,
SUM(amount) AS gmv
FROM shop.orders
GROUP BY user_id;
# 批模式: 处理完所有数据后输出(有界流)
# 流模式: 数据到达即增量更新(无界流)
# 优化器: 同一计划生成, 差异只在流式算子
3.3 窗口语义一致
流批统一最考验窗口:批的"一天"和流的"一天"必须是同一个事件时间窗口。
-- 流批一致的 TUMBLE 窗口
SELECT
TUMBLE_START(event_time, INTERVAL '1' DAY) AS day,
SUM(amount) AS gmv
FROM shop.orders
GROUP BY TUMBLE(event_time, INTERVAL '1' DAY);
# 一致性关键
# 都用事件时间(而非处理时间)
# 迟到数据策略一致(批=全量, 流=水印+侧输出)
# 去重语义一致(主键去重)
四、实时数仓与批处理对账
4.1 为什么必须对账
流批一体最大的风险是"看似一体、实际漂移":流式聚合与批式重算结果不一致。对账就是用批式结果校验流式结果,把漂移暴露出来。
# 对账对象
# 行数 / 主键唯一性
# 金额类指标(SUM/AVG)
# 状态类指标(当前状态分布)
# 窗口边界口径
4.2 对账方法
-- 对账: 流式表 vs 日批表的 GMV
SELECT
'stream' AS source, COALESCE(SUM(amount), 0) AS gmv
FROM stream_dws_daily_gmv WHERE dt = '2026-09-27'
UNION ALL
SELECT
'batch' AS source, COALESCE(SUM(amount), 0) AS gmv
FROM batch_dws_daily_gmv WHERE dt = '2026-09-27';
-- 差异非零 → 告警 → 定位漂移环节
4.3 对账工具化
# 对账工程化
# 每日定时对账 + 差异告警
# 差异阈值分级(可容忍/需调查/事故)
# 漂移根因分类: 迟到数据/口径差异/去重差异
# 对账结果纳入数据质量监控
五、统一计算层选型
5.1 引擎能力对比
| 引擎 | 批 | 流 | 同语义 | 生态 |
|---|---|---|---|---|
| Flink | ✅ | ✅ | 强(SQL/Table) | Kafka/湖仓 |
| Spark | ✅ | 微批 | 中等 | 批生态成熟 |
| Doris/StarRocks | ✅ | 部分 | 有限 | 实时分析 |
| 湖仓 + 引擎 | 组合 | 组合 | 视实现 | 开放 |
5.2 Flink vs Spark 的选择
# Flink: 真流式 + SQL 同语义, 流批一体首选
# Spark: 批生态强 + 微批近似流, 存量批任务平滑
# 实践: 新实时链路用 Flink, 存量批保留 Spark, 逐渐归一
# 关键: 选一个能"一套逻辑两模式"的引擎
5.3 数据格式统一
流批一体的前提是数据格式统一,否则流式写入的 Parquet 与批式格式不一致,对账无从谈起:
# 统一要求
# 统一表格式(Iceberg/Delta/Hudi)
# 统一序列化(Avro/JSON Schema)
# 统一目录(湖仓 Catalog)
# 统一分区与主键定义
六、流批一体的数据架构
6.1 湖仓作为统一底座
流批一体的物理底座通常是湖仓:对象存储 + 开放表格式 + 统一目录。流和批都写同一张 Iceberg 表,查询引擎统一读取。
# 架构形态
# 接入: Kafka(事件) + CDC
# 统一表: Iceberg 分层(Bronze/Silver/Gold)
# 流写入: Flink 两阶段提交
# 批写入: Spark 批任务
# 查询: Trino / Doris 统一读
6.2 一张表的双模写入
-- 流式 upsert 与批式 merge 写同一张表
MERGE INTO lake.dws.daily_order_revenue t
USING (SELECT order_id, amount, dt FROM lake.ods.orders_cdc) s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET t.amount = s.amount
WHEN NOT MATCHED THEN INSERT (order_id, amount, dt)
VALUES (s.order_id, s.amount, s.dt);
快照隔离让流批并发写同一张表互不阻塞——这是湖仓表格式给流批一体最大的礼物。
6.3 增量读与回算
# 增量读: 从快照差异读新增 → 下游增量消费
# 回算: 指定历史快照重放 → 修正口径
# 时间旅行: 审计任意时点数据
# 流批一体 = 流提供实时, 批/重放提供纠错
七、从流批分离到一体的演进路径
7.1 现状评估
# 迁移前评估
# 盘点: 现有批管道/流管道清单
# 口径: 哪些指标流批不一致(高风险)
# 数据: 哪些已入湖/可重放
# 依赖: 下游消费方与 SLA
7.2 分阶段迁移
| 阶段 | 动作 | 成果 |
|---|---|---|
| 阶段一 | 湖仓底座就绪,统一格式 | 存储统一 |
| 阶段二 | 核心指标切 Flink SQL 双模式 | 新链路一体 |
| 阶段三 | 实时链路接入湖仓 + 对账 | 实时+准确 |
| 阶段四 | 存量批任务逐个归并 | 代码收敛 |
| 阶段五 | 下线 Lambda 合并层 | 架构简化 |
7.3 风险与回退
# 风险: 流式引擎重放性能不足
# 缓解: 保留近期快照 + 用批引擎补历史
# 风险: 对账差异长期不收敛
# 缓解: 先统一口径, 再统一实现
# 回退策略: 每条链路双跑期, 达标后再切换
八、工程实践与避坑
8.1 实践清单
# [ ] 一套引擎、一份 SQL, 流批同语义
# [ ] 事件时间 + 一致窗口
# [ ] 湖仓单底座, 流批写同一张表
# [ ] 每日对账 + 差异告警
# [ ] 迟到数据策略明确
# [ ] 回算能力就绪(快照/重放)
# [ ] 流批依赖同一目录与序列化
8.2 常见陷阱
- 用处理时间当窗口:流批对账永远对不上。
- 流批写不同表:失去统一底座的收益。
- 只做实时不做对账:口径漂移静默恶化。
- Kappa 硬扛海量历史:保留期与重放成本失控。
- 双跑期太短:差异未收敛就切换,事故后很难回退。
8.3 案例:电商 GMV 流批一体
某电商把 GMV 从 Lambda 迁移到"Flink SQL + Iceberg":
| 阶段 | 动作 | 结果 |
|---|---|---|
| 迁移 | 核心指标统一为 Flink SQL | 一套代码 |
| 底座 | 流批写同一 Iceberg 表 | 口径一致 |
| 对账 | 每日对账 + 差异告警 | 漂移 24h 内暴露 |
| 收尾 | 下线批链路与合并层 | 运维减半 |
总结
| 决策点 | 推荐 | 理由 |
|---|---|---|
| 架构 | 流批一体优先 | 单一语义、单一代码 |
| 引擎 | Flink SQL/Table API | 真流式 + 同语义 |
| 底座 | 湖仓(Iceberg) | 统一存储 + 快照隔离 |
| 写入 | 流批双模写同一张表 | 口径天然一致 |
| 对账 | 每日批式校验 | 暴露漂移 |
| 回算 | 快照 + 重放 | 纠错能力 |
流批一体的终局是**“一份数据、一套逻辑、任意模式”**:逻辑不因运行模式而变,数据不因链路而分裂。它不是某个引擎的专利,而是一种架构纪律——统一语义、统一底座、统一对账。从核心指标试点,让对账先跑起来,再逐步把存量链路归并进来,流批分离的负担终将被一点点卸掉。
参考与延伸阅读
- Flink 官方文档:Table/SQL API 流批统一与动态表
- Jay Kreps:Questioning the Lambda Architecture(Kappa 起源)
- Apache Iceberg 官方文档:增量读取与快照
- 实时数据仓库 — Kafka+Flink+ClickHouse 落地
- 流处理进阶 — 窗口与水印语义
- Apache Flink 流处理 — Flink 实战基础
- Apache Spark 批处理 — 批式引擎衔接
- 湖仓一体架构 — 湖仓统一底座
- Apache Iceberg 深度剖析 — 表格式机制
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。