流批一体架构实践

流批一体架构实践:Lambda 与 Kappa 架构、统一计算引擎 Flink/Spark、数据湖仓一体化、精确一次语义、实践案例与踩坑清单

传统数据处理被劈成两半:**批量(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 既能跑离线任务也能跑实时任务。

维度LambdaKappa流批一体
引擎数211
代码/口径两套,易不一致一套一套,天然一致
离线能力强弱(依赖重放)强(批模式)
实时能力弱强强(流模式)
运维重中轻

2. 统一计算引擎

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精确一次,不丢不重最强保证,代价最高

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. 落地建议

  1. 先统一口径,再统一引擎:指标定义放一处,引擎只是执行者
  2. 从 Kappa 起步,批作为重放:减少一套链路
  3. 湖格式尽早引入:事务与时间旅行是流批共存的基石
  4. 精确一次按需:非计费场景用 At-Least-Once + 幂等,省下复杂度
  5. 监控全面:lag、checkpoint 时长、小文件数量、状态大小都要看
  6. 演练故障恢复: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/ 的幂等保障,可以构建出既快又准、可回溯可对账的数据处理体系。

继续阅读

探索更多技术文章

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

全部文章 返回首页

「distributed-systems」更多文章

  1. Serverless 架构实践
  2. 事件溯源与 CQRS 架构
  3. 分布式时钟与逻辑时钟