传统数据处理被劈成两半:**批量(Batch)**追求"离线算准",流式(Streaming)追求"实时算出"。两套引擎、两套代码、两套口径,结果常常"实时对不上离线"。流批一体(Stream-Batch Fusion)用同一套引擎、同一套 SQL、同一套口径同时处理离线与实时,让"用批的准确、用流的时效"。本文从 Lambda/Kappa 演进讲起,覆盖统一引擎、湖仓一体化、精确一次语义与实战踩坑。
一句话:流批一体不是消灭批或流,而是让一套逻辑两种跑法,口径永远一致。
1. 从 Lambda 到流批一体
1.1 Lambda 架构
Lambda 同时维护批与流两条链路,最后合并结果:
数据 ──► 批路径:全量重算(离线数仓)──► 批结果 ─┐
└─► 流路径:增量计算(实时)──► 流结果 ─┴─► 合并查询层
痛点一目了然:两套代码两套口径,“批是对、流是快"的说法在合并层经常打架;运维成本翻倍。
1.2 Kappa 架构
Kappa 主张"一切皆流”,只用流引擎:
数据 ──► 消息队列/Kafka(保留完整历史)──► 流引擎重放计算 ──► 结果存储
离线需求 = 从 Kafka 起点重新流式计算(重放),不再有独立批链路
Kappa 消除了双代码问题,但依赖"消息队列保存全量历史",长周期重放的存储与计算成本不低。
1.3 流批一体的本质
流批一体是对 Kappa 的工程化深化:同一套引擎提供批(Bounded)与流(Unbounded)两种执行模式,同一份 SQL 既能跑离线任务也能跑实时任务。
| 维度 | Lambda | Kappa | 流批一体 |
|---|---|---|---|
| 引擎数 | 2 | 1 | 1 |
| 代码/口径 | 两套,易不一致 | 一套 | 一套,天然一致 |
| 离线能力 | 强 | 弱(依赖重放) | 强(批模式) |
| 实时能力 | 弱 | 强 | 强(流模式) |
| 运维 | 重 | 中 | 轻 |
2. 统一计算引擎
2.1 Flink 的批流一体
Flink 把批看作"有界流"(Bounded Stream),流看作"无界流",运行时统一处理。同一份 SQL 在批模式与流模式下执行,结果语义一致:
-- 这份 SQL 既可用于离线日批,也可用于实时滚动窗口
SELECT
user_id,
COUNT(*) AS order_cnt,
SUM(amount) AS amount_total
FROM orders
GROUP BY user_id, TUMBLE(ts, INTERVAL '5' MINUTE)
批模式:读取有界输入,全量计算后输出
流模式:读取无界输入,窗口触发增量输出
两种模式产出同一口径
2.2 Spark 的批流一体
Spark Structured Streaming 同样用"微批 + 增量处理"统一 API:批任务是"一条 SQL 一次性执行",流任务是"同一 SQL 每个微批执行"。Spark 的优势是批生态成熟、易与数据湖集成。
2.3 引擎选型
| 需求 | 推荐引擎 | 理由 |
|---|---|---|
| 低延迟实时(秒级) | Flink | 真正的流处理,延迟低 |
| 大规模离线批 | Spark | 批处理生态与稳定性成熟 |
| 批流一体统一 SQL | 均可 | Flink/Spark 都提供统一 API |
| 与数据湖强绑定 | Spark + 湖格式 | 湖表读写支持最全 |
2.4 查询统一:物化视图
批流一体不仅统一计算引擎,还统一查询视角:同一张逻辑表,离线跑全量、实时跑增量,通过物化视图对外暴露一致结果。物化视图由引擎自动维护,批模式全量构建、流模式增量更新,查询侧无需区分数据来自实时还是离线:
物化视图 orders_daily:
批模式:每日凌晨对全量数据重算构建
流模式:滚动窗口对增量数据更新维护
查询侧:select * from orders_daily,看到同一口径
物化视图解决了"实时表与离线表各查各的"问题,是流批一体在查询层的自然延伸,也让 BI 工具只面对一套表结构。
3. 数据湖仓一体化
3.1 湖仓(Lakehouse)解决什么
传统数仓成本高、扩展慢;数据湖便宜但缺事务与表语义。湖仓把事务、Schema、ACID 带到数据湖上:
| 能力 | 数据湖 | 湖仓 |
|---|---|---|
| 文件格式 | Parquet/ORC | 同左 |
| 事务 | 无(多个 writer 互相覆盖) | ACID 增量提交 |
| 表语义 | 弱 | 类数仓表 + Schema 演进 |
| 实时写入 | 困难 | 流式写入 + 小文件治理 |
3.2 三种主流湖格式
| 格式 | 特色 | 典型场景 |
|---|---|---|
| Apache Iceberg | 快照隔离、隐藏分区、时间旅行 | 数仓演进、多引擎共享 |
| Apache Hudi | 增量 upsert、MOR 读写优化 | 实时更新、近实时 |
| Delta Lake | 与 Spark 深度绑定、Vacuum 治理 | Spark 为主的湖仓 |
-- Iceberg 时间旅行:查询任意历史快照
SELECT * FROM orders_history
FOR SYSTEM_TIME AS OF '2026-09-28 12:00:00';
3.3 CDC 入湖
流批一体的典型数据入口是 CDC(Change Data Capture):数据库 binlog 实时入湖,湖上既跑实时指标又跑离线分析,源头口径统一:
MySQL binlog ──► Canal/Debezium ──► Kafka ──► Flink/Spark ──► Iceberg 湖表
│
├─► 实时维度表(流)
└─► 离线 ODS 分区(批)
4. 精确一次语义
4.1 三种交付语义
| 语义 | 含义 | 结果 |
|---|---|---|
| At-Most-Once | 最多一次,丢数据 | 吞吐最高,可能丢 |
| At-Least-Once | 至少一次,可能重复 | 简单,下游需去重 |
| Exactly-Once | 精确一次,不丢不重 | 最强保证,代价最高 |
4.2 Flink 精确一次
Flink 用**检查点(Checkpoint)+ 两阶段提交(两阶段提交)**实现端到端精确一次:
Barrier 机制:
算子收到 barrier 后快照状态
所有 barrier 对齐 → 快照完成,视为一致点
故障恢复:从最近一致点重放,配合 Kafka 幂等写入
// 端到端精确一次的关键配置
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(60_000); // 每 60s 一个 checkpoint
env.getCheckpointConfig()
.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
// Sink 侧:两阶段提交,事务提交前不落最终结果
FlinkKafkaProducer<String> kafkaSink = new FlinkKafkaProducer<>(
"sink-topic",
serializer,
FlinkKafkaProducer.Semantic.EXACTLY_ONCE);
4.3 精确一次与幂等的取舍
精确一次依赖"检查点 + 状态后端 + 事务型 Sink",性能与复杂度都不低。实践中常用折中:At-Least-Once + 下游幂等。对账单、计费等"可去重"的场景,At-Least-Once 更简单务实。幂等方案可参考 https://plumephp.com/distributed-idempotency-reliability/。
4.4 精确一次的适用判断
| 场景 | 语义选择 | 理由 |
|---|---|---|
| 支付、计费、扣减 | Exactly-Once | 重复即直接损失 |
| 统计指标、报表 | At-Least-Once + 幂等 | 可去重,成本低 |
| 日志、行为数据 | At-Most-Once | 少量丢失可接受 |
| 大促实时大屏 | At-Least-Once + 对账 | 追求低延迟,结果可复核 |
判断原则:按"重复的代价"而非"技术的酷炫"选择语义。业务侧能去重就用幂等替代精确一次,把复杂度留给真正需要它的场景。
5. 实践案例
5.1 实时数仓分层
ODS 层:CDC 入湖的原始数据(批流共用)
DWD 层:清洗明细(流式宽表)
DWS 层:聚合指标(滚动窗口,批流同口径)
ADS 层:应用侧结果(实时报表 + 离线回溯)
5.2 大促实时大屏
实时大屏用流式计算 GMV、订单量等指标;活动结束后,同一套 SQL 跑离线批任务,产出与实时完全一致的"最终成绩单"——口径一致性直接避免"实时和报表对不上"的公关事故。
5.3 湖表近实时分析
Iceberg 支持流式写入小文件 + 周期性 Compaction。数据分钟级可见,同时保留离线全量分析的确定性。
5.4 与消息队列的关系
流批一体的数据源几乎都是消息队列:Kafka 既承担缓冲,也是重放的"时间轴"。对消息队列的选型、分区与消费语义,可参考 https://plumephp.com/message-queue-deep-dive/。
5.5 案例:用户行为分析漏斗
用户点击流以事件形式进入 Kafka,流引擎实时计算转化漏斗;离线侧用同一份 SQL 重跑全量行为数据,产出"历史漏斗基线"。实时漏斗与离线基线口径完全一致,运营看到"当前转化率 vs 历史基线"直接可比:
点击 → 浏览 → 加购 → 支付(转化漏斗)
实时链路:流式窗口计算今日漏斗,分钟级刷新
离线链路:批任务回溯历史漏斗基线,每日更新
同一份 SQL、同一套口径,差异只剩时间维度
该案例的典型收益:实时结果与离线报表永远对得上,避免"大屏一个数、报表另一个数"的经典事故。
6. 踩坑清单
6.1 Watermark 与乱序
流处理必须处理乱序。Watermark 表示"这个时间之前的数据已到齐",设置不当会造成窗口少算:
Watermark = 观察到的最大事件时间 - 允许乱序延迟
允许乱序延迟过小 → 迟到数据被丢(少算)
允许乱序延迟过大 → 结果迟迟不输出(延迟高)
6.2 状态管理
有状态算子的状态无限增长会拖垮作业:
- 用 TTL 清理过期状态(如 7 天)
- 大状态用 RocksDB 状态后端,注意磁盘与性能
- 状态 Schema 变更要配合迁移,避免状态不兼容
6.3 小文件问题
流式持续写入湖表会产生海量小文件,查询性能直线下降:
治理:定期 Compaction 合并小文件
策略:写入后按大小/时间触发合并任务
6.4 背压
流作业上游快、下游慢形成背压。检查点超时、消息积压、延迟上涨是背压的典型信号。背压治理要明确"积压可容忍的阈值、触发降级的条件",与限流熔断机制联动:入口限流挡住超量数据,积压达到阈值触发降级或告警扩容。相关机制可参考 https://plumephp.com/rate-limiting-circuit-breaker/。
6.5 多份口径漂移
流批一体最大的风险仍是口径漂移:同是 GMV,实时口径把"取消订单"算进去了,离线口径不算。口径要收敛为一处定义(统一指标平台),批流只是同一口径的两种执行。
6.6 事务与对账
湖仓表级事务(如 https://plumephp.com/distributed-transactions/ 视角下的跨表一致性)配合对账机制,防止"看起来写成功了实际不一致"。对账是流批链路最后的兜底防线。
7. 落地建议
- 先统一口径,再统一引擎:指标定义放一处,引擎只是执行者
- 从 Kappa 起步,批作为重放:减少一套链路
- 湖格式尽早引入:事务与时间旅行是流批共存的基石
- 精确一次按需:非计费场景用 At-Least-Once + 幂等,省下复杂度
- 监控全面:lag、checkpoint 时长、小文件数量、状态大小都要看
- 演练故障恢复:checkpoint 恢复、湖表回滚、重放演练常态化
总结
| 主题 | 关键内容 |
|---|---|
| 架构演进 | Lambda → Kappa → 流批一体(一套口径两种跑法) |
| 统一引擎 | Flink 有界/无界统一、Spark 微批统一 |
| 湖仓一体化 | Iceberg/Hudi/Delta、ACID 与时间旅行、CDC 入湖 |
| 精确一次 | Checkpoint + 两阶段提交,或用幂等折中 |
| 实践案例 | 实时数仓、大促大屏、近实时湖表 |
| 踩坑 | Watermark、状态 TTL、小文件、口径漂移、对账 |
流批一体的核心不是某项具体技术,而是**“口径收敛"的工程理念**:把"准"与"快"统一到同一份逻辑,让离线与实时互为备份、彼此印证。它大幅降低了数据链路的双维护成本,也让"实时结果可被离线复现"成为质量红线。配合 https://plumephp.com/message-queue-deep-dive/ 的数据通道设计和 https://plumephp.com/distributed-idempotency-reliability/ 的幂等保障,可以构建出既快又准、可回溯可对账的数据处理体系。
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。