批处理框架和流处理框架不是二选一的对立关系,而是可以按数据时效需求灵活组合,用一套技术栈覆盖离线与实时两类场景。这是当前数据架构领域已经形成的行业共识,也是Lambda架构和Kappa架构在实际落地中反复验证过的路径。
为什么批处理和流处理要搭配着用
先看两边的性格差异,批处理框架的代表是MapReduce、Spark SQL、Hive这类工具,它们适合处理海量历史数据,吞吐量高,但延迟通常在分钟级甚至小时级,流处理框架则以Flink、Storm、Kafka Streams为主,主打毫秒级或秒级延迟,但面对超大状态和复杂ETL逻辑时,资源消耗和运维复杂度会明显上升。
把两者放到同一张桌上,你会发现它们其实是互补的。
- 批处理的强项在于全量计算的准确性,它能对过去一整天的数据做完整的重算,修正流处理中可能出现的乱序或迟到数据问题。
- 流处理的强项在于实时反馈的速度,它能在交易发生的同时完成风控判断,或者在大促看板上实时滚动成交额。
- 真实业务往往同时需要这两类能力,比如日志分析,白天用流处理实时监控异常,晚上用批处理生成完整报表。
- 只靠单一阵营,要么牺牲时效,要么牺牲数据完整性。
表格对比更直观:
| 维度 | 批处理框架 | 流处理框架 |
|---|---|---|
| 延迟 | 分钟级到小时级 | 毫秒级到秒级 |
| 数据范围 | 全量历史数据 | 窗口内增量数据 |
| 计算模式 | 一次性触发或定时调度 | 持续运行、事件驱动 |
| 消费场景 | 报表、数据仓库、机器学习训练 | 实时大屏、风控、实时推荐 |
| 资源特征 | 批式高峰占用,可错峰 | 常驻资源,持续消耗 |
行业共识认为,纯批处理无法支撑实时决策,纯流处理又难以在故障恢复时做全量回溯,组合使用不是把两套系统简单堆在一起,而是让它们各管一段,数据在不同的时效层级上各取所需。
实时数仓架构怎么选
这是很多开发团队一开始就会卡住的问题,有人上来就追Flink,认为把一切做成流就是未来;也有人守着Hive不肯放,觉得稳定压倒一切,选择实时数仓架构的关键不在于技术栈的新旧,而在于你愿意为多少延迟买单、愿意付出多少运维成本。

批流分离:先用离线链路把底打牢
最稳妥的做法,是把离线批处理作为数据仓库的主链路,每天凌晨调度Hive或Spark任务,将前一天的业务数据清洗、加工、分层,写入数仓的明细层和汇总层,这条路走得很成熟,数据质量可控,回溯也方便。
在此基础上,单独搭建一条轻量级的实时链路,用Flink或Kafka Streams只处理当天的高时效数据,比如实时订单金额、实时用户数、实时异常告警,这套组合的典型特征是离线链路负责全量,实时链路负责增量,两边通过统一的数据格式和存储层打通。
这种方式适合大多数中小团队,它不要求在架构上一步到位,也不必让流任务承接复杂的数据治理逻辑。
批流一体:让一套引擎搞定两类计算
如果团队技术储备够、业务对时效的敏感度又极高,就可以考虑批流一体,用Flink的DataStream API和Table API统一处理有界数据和无界数据,离线任务和实时任务跑在同一个引擎上,复用同一套SQL逻辑、同一套UDF函数。
这样一来,离线链路和实时链路不再分家,数据的口径天然一致,对于数据时效要求苛刻的场景,比如实时个性化推荐的特征计算、实时反欺诈的规则引擎,这种做法的收益非常明显,缺点是Flink本身的学习曲线较陡,状态后端、检查点机制、水位线这些概念对新人并不友好。
Lambda架构和Kappa架构区别
这两套架构的争论,本质上就是在回答批处理和流处理怎么组合的问题。
Lambda架构在2011年前后由Storm的作者Nathan Marz提出,核心思路是并行维护两条链路,一条走批处理,负责完整准确的历史数据,产出最终的视图;另一条走实时计算,负责最新的增量数据,提供低延迟的临时视图,查询时合并两边的结果。
Kappa架构则更激进一些,它主张放弃独立的批处理链路,所有数据都走流式处理,用Kafka这样的分布式日志系统做统一的数据缓冲,需要全量重算时,直接把历史数据重新输入流处理引擎跑一遍即可,业内专家指出,Kappa架构在理论上更简洁,但对消息队列的存储能力和流引擎的重放性能要求比较高,实践中真正完全放弃批处理的团队并不算多。
两套链路各自的适用场景
选择时需要结合团队实际情况判断。
- Lambda架构适合的场景:数据量庞大且口径经常调整,比如金融行业的历史交易明细、用户行为日志的长期留存分析,这类场景需要频繁回溯重算,离线链路的存在让数据修复变得简单直接。
- Kappa架构适合的场景:数据模式相对稳定,业务高度依赖实时反馈,比如短视频推荐里需要秒级更新的用户兴趣标签、电商大促过程中的实时库存扣减。
- 一个参考做法是,先按Lambda架构把离线数据和实时数据分开建设,等Flink团队成熟后,再用流式重放替代部分离线计算,这样风险最低。

数据时效要求高的场景有哪些
组合使用的核心驱动力就是数据时效,哪些场景必须使用批处理和流处理的组合,工程上其实有清晰的分层标准。
- 小时级或天级时效:离线报表、财务结算、数据仓库分层建设,用批处理即可。
- 分钟级时效:运营看板、监控预警、流量波动分析,可以引入微批处理,比如Spark Streaming。
- 秒级甚至毫秒级时效:风控拦截、实时竞价、异常检测,必须使用真正的流处理引擎,如Flink或Kafka Streams。
一个比较典型的真实案例是电商平台的用户增长看板,运营人员进入后台查看大促实时成交额,背后是Flink在做秒级聚合;同时分析师早上上班看到的日活、转化率、复购率报表,则是凌晨用Spark批处理产出的结果,两条链路并行不悖,数据最终都落到同一张宽表里。
实操落地:批处理和流处理组合的具体步骤
下面给出一种经过验证的落地路径,便于按图索骥。
第一步:梳理数据时效分级,打开数据字典,把当前所有的数据需求列出,按实时、准实时、离线三个等级打标,实时数据对应毫秒到秒级响应,准实时数据对应分钟级延迟,离线数据对应小时级或天级周期。
第二步:选型匹配框架,离线链路建议用Spark SQL或Trino,配合Hive或Iceberg存储;实时链路用Flink,配合Kafka作为数据管道;准实时链路可以用Flink的微批模式或Spark Structured Streaming,按需选择。
第三步:统一存储与元数据,这是组合使用的核心关键,离线作业和实时作业需要读写同一份数据资产,建议引入Hudi或Iceberg这类数据湖表格式,支持流式写入和批量读取,这样批处理产出的结果可以直接被流处理引用,流处理的输出也能无缝融入离线数仓。
第四步:设置数据质量校验任务,用批处理的结果校验流处理累计值,比如拿前一天的离线交易总额与当日实时总额的存档做交叉比对,差值超过阈值就触发告警,这能有效规避流处理乱序数据带来的误差。

运维和技术栈层面的注意事项
组合使用意味着技术栈变宽,有几个运维细节容易被忽略。
- 资源隔离:批处理任务一般集中在夜间跑,流任务全天候运行,如果共用Yarn队列,夜间批处理的高峰容易把流任务的资源挤掉,导致实时链路延迟飙升,建议为实时任务单独划分资源队列,并设置优先级。
- 状态管理:Flink的状态后端要单独规划,使用RocksDB时需为每个作业预留足够的本地磁盘空间,否则checkpoint失败会产生雪崩效应。
- 数据延迟补偿:流处理总会存在丢数据或迟到数据的情况,设计时要预留一个补偿通道,常用的做法是流处理产出的结果写入Kafka,批处理在次日对同一批数据做全量重算,用重算结果覆盖昨天的临时数据。
- 监控指标:除了任务本身的运行状态,还要关注数据产出延迟,可以给每个数据表添加一个“数据新鲜度”指标,记录最后一条数据的写入时间,延迟超过阈值则自动提醒。
Q&A
批处理和流处理框架组合使用会不会增加太多运维负担?
确实会,多一套Flink集群意味着需要额外关注checkpoint、反压、状态大小等问题,为降低负担,建议先在同一个数据湖存储层上做统一,让批任务和流任务共用同一份表数据,减少双份数据的存储和同步费用,离线任务与实时任务尽量复用同一个调度平台,如Airflow或DolphinScheduler,避免两套系统割裂。
什么情况下可以只用流处理框架,完全放弃批处理?
当业务的数据量可控且历史回溯需求极少时可以考虑,比如轻量级物联网设备上报数据的实时展示,保留最近几天数据即可满足需求,但涉及财务结算、审计合规或需要跨年度对比分析的场景,离线批处理的重算能力仍然不可替代。
Kafka在批流组合中一般扮演什么角色?
Kafka主要作为统一的数据管道,上游业务数据先写入Kafka,实时链路直接从Kafka消费计算,离线链路则由调度任务按批次从Kafka拉取落盘,用Kafka做缓冲的好处是,流处理和批处理消费的是同一份数据镜像,不会因为两边读取时机不同而产生数据差异,同时也方便对历史消息做Replay回溯计算。