考察点
这是实时数仓的入门架构题,考的是你对两种架构的取舍有没有实战体感。面试官想听到的不是定义背诵,而是:Lambda 双链路口径不一致怎么发生的、Kappa 重算为什么难、以及你们实际怎么选的。追问大概率落到流批一体(Flink SQL、湖仓一体)如何缓解这个问题。
参考答案
Lambda:一套业务,两套链路
Lambda 分三层:批处理层(Batch Layer)存全量历史,用 Hive/Spark 跑 T+1,结果准确可重算;速度层(Speed Layer)用 Flink/Storm 消费 Kafka 实时流,产出最近几小时的增量结果,牺牲一点准确换时效;服务层(Serving Layer)把批视图和实时视图合并对外——查询最近数据时,用实时结果补齐批处理还没覆盖的窗口。
它的问题实践中每个人都踩过:
- 口径漂移。同一个指标写两遍,离线用 Hive SQL、实时用 Flink,两边对「有效订单」的过滤条件、时区处理、去重逻辑稍有出入,批层覆盖掉速度层结果时数字会跳一下,业务方立刻来问。保持两套代码逻辑等价是持续性投入,不是一次对齐就完事。
- 双倍成本。存储两份、计算两份、运维两个技术栈,招人都要两拨技能。
- 链路复杂。服务层合并逻辑(以批为准、实时补差)本身就容易出 bug。
但 Lambda 有不可替代的优点:批层是全量、可重算的真理来源。任何口径变更、脏数据修复,重跑批任务就行,历史永远是对的。
Kappa:一切皆是流
Kappa 的出发点很直接:既然速度层能干,能不能不要批层?架构上只留一条流链路——所有数据进 Kafka,Flink 消费处理直接出结果。需要重算时,把 Kafka 里的历史数据从头回放(replay)一遍,用同一份代码重新跑出结果覆盖旧的。没有双链路,口径天然一致,代码只维护一份。
代价也很硬:
- 回放能力依赖消息保留期。Kafka 保留 7 天,重算三个月历史就无从谈起,要延长保留或把历史归档到廉价存储(对象存储 + 归档 topic),成本和组织复杂度都上来了。
- 大状态难搞。流上算长周期指标(累计 90 天留存),状态要一直攒着,RocksDB state backend 撑到 TB 级时 checkpoint 时长和恢复时间都很难看。批处理跑全量反而轻松。
- 乱序和 join 的复杂度全在流上解决。双流 join 的状态管理、维表变更的时点语义(维度回滚),都比批处理难写难调。
怎么选
| 考量 | 偏 Lambda | 偏 Kappa |
|---|---|---|
| 实时需求占比 | 少数场景要实时,大盘是 T+1 | 绝大多数据都要秒/分钟级 |
| 历史重算频率 | 口径常变、经常回刷 | 口径稳定,重算少 |
| 数据量与状态 | 明细量大、长周期指标多 | 链路短、状态可控 |
| 团队 | 有离线团队,加实时小组 | 团队以流计算为主 |
流批一体:现实的中间路线
纯 Lambda 和纯 Kappa 都少见了,主流是用一套引擎抹平两条链路:Flink SQL 统一批流 API,同一套 SQL 既能跑 Kafka 流也能跑 Hive/Iceberg 批;湖仓一体(Paimon、Hudi、Iceberg)让流写入的表同时支持批读批算,重算时不用回放 Kafka,直接对湖表跑批。这本质是把 Lambda 的「双链路」压缩成「一套代码、两种执行模式」,口径不一致和维护成本都大幅下降。面试里聊到这个层度,说明你跟上了最近几年的实践演进。
可能的追问
- Lambda 里批层和速度层结果合并时怎么保证不重复不丢? 关键是时间边界的幂等:批层按分区覆盖(T 日批结果覆盖 T 日整分区),实时只负责批还没跑出来的窗口,查询时按时间路由,交界处以批为准,实时结果设置 TTL 到期失效。
- Kappa 回放历史时正在跑的实时作业怎么办? 一般用双跑切换:新作业从归档起点回放追到当前位点,与线上作业结果对平后切流量,老作业再停,避免停流重算造成数据断档。
- Flink 流批一体真能做到一套 SQL 吗? 计算逻辑可以,但 sink 和语义细节仍有差异(流上是持续输出、批是一次产出,upsert 语义、分区提交时机不同),工程上会沉淀公共视图层屏蔽差异,业务 SQL 复用率能做到八成以上。