流批一体:从 Lambda/Kappa 架构到统一计算层

深入解析流批一体(Stream-Batch Unification)架构:流批分离的历史与痛点、Lambda 架构双链路的一致性问题、Kappa 架构的统一路径、流批同语义(Flink SQL/Table API)与窗口语义一致性、实时数仓与批处理对账方法、统一计算层选型(Flink/Spark/湖仓底座)、一张表的双模写入与增量回算、从流批分离到一体的分阶段演进路径与工程避坑清单。

引言

过去十年,实时与批处理被当成两条路:批处理算得全但慢,流处理算得快但"不够严谨"。为了两头兼顾,很多团队搭建了 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 两者对比

维度LambdaKappa
计算链路流 + 批 两条一条流式
代码一致性两套逻辑难对齐天然一致
回算能力批链负责重放实现
复杂度高(合并层)低
适用历史包袱重新建/可重放

2.2 什么时候仍需要 Lambda

  • 历史数据量巨大,无法全部常驻 Kafka/事件流。
  • 流式引擎重放海量历史性能不可接受。
  • 已有成熟的批式数仓,迁移成本过高。

2.3 走向一体的判断

# 判断条件
# [ ] 业务是否同时需要秒级 + 精确结果
# [ ] 数据能否在事件流中保留足够久
# [ ] 回算是否只需覆盖近期窗口
# [ ] 团队是否能驾驭统一引擎
# 大多数新项目: 用流批一体(湖仓 + 统一引擎)更划算

三、流批同语义

3.1 为什么需要同语义

流批"对不上"的根因,是同一指标用了两套语义:批用 SQL 算、流用代码算,窗口、去重、迟到处理规则都不同。同语义要求:批 SQL 和流 SQL 是同一份逻辑。

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✅部分有限实时分析
湖仓 + 引擎组合组合视实现开放
# 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)统一存储 + 快照隔离
写入流批双模写同一张表口径天然一致
对账每日批式校验暴露漂移
回算快照 + 重放纠错能力

流批一体的终局是**“一份数据、一套逻辑、任意模式”**:逻辑不因运行模式而变,数据不因链路而分裂。它不是某个引擎的专利,而是一种架构纪律——统一语义、统一底座、统一对账。从核心指标试点,让对账先跑起来,再逐步把存量链路归并进来,流批分离的负担终将被一点点卸掉。


参考与延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「data-engineering」更多文章

  1. 数据平台成本与 FinOps:存储、计算、弹性与降本实践
  2. 数据网格 Data Mesh:领域数据产品、自助平台与联邦治理
  3. 列存与向量化查询引擎:Doris、ClickHouse、StarRocks 与 Trino 的性能本质