Batch vs Stream Processing
批处理与流处理怎么选:有界与无界输入、事件时间与处理时间、窗口与 watermark、迟到数据、重放和回填、micro-batch,以及 Lambda 与 Kappa 两种混合做法;附一个交易风控加日报的算例、翻车表和面试答法。
Batch processing 是把一段时间的数据先收集起来,再一次性处理。它更适合大规模、可延迟的计算。
Stream processing 则是数据一到就处理,强调 near real-time 响应。

本章回答的是选哪一种、混合时怎么分工。流处理内部的四种拓扑(consumer group、changelog 状态、checkpoint 窗口、多消费者留存日志)和它们的故障恢复,在 流处理架构 里逐一展开。
真正的区别:输入有没有尽头
Google 2015 年的 Dataflow Model 论文(VLDB)建议别再用 batch / stream 描述引擎,而是描述数据:有界(bounded) 数据有尽头,作业读完就结束;无界(unbounded) 数据一直在来,作业永远不会「读完」。由此派生出三个真正要做的决定:
- 按哪个时钟算。 事件时间(event time)是事件发生的时间,处理时间(processing time)是它被处理的时间。手机离线 20 分钟再上传,两者就差 20 分钟。批处理天然按事件时间:昨天的分区里就是昨天的数据。流处理必须显式选。
- 什么时候算「够了」。 批处理的答案是「输入读完了」。流处理没有读完这一天,要靠 watermark 声明「事件时间 T 之前的数据应该都到了」,然后关闭窗口出结果。Flink 的 allowed lateness 默认是 0:落在 watermark 之后的元素直接丢弃,除非你配了允许迟到的时长或者旁路输出。
- 错了怎么重算。 批处理重跑一遍输入即可,前提是输入不可变、输出按分区覆盖写。流处理要从保留的日志重放,Kafka topic 的
retention.ms默认 604800000,也就是 7 天:超过 7 天的错误,流里已经没有原始数据可重放了。
对比
| 维度 | Batch | Stream |
|---|---|---|
| 输入 | 有界:一个小时、一天的分区 | 无界:持续到达的事件 |
| 端到端延迟 | 调度周期加运行时长,通常分钟到小时 | 秒级甚至亚秒级 |
| 正确性 | 输入完整后再算,天然处理迟到 | 靠 watermark 和允许迟到,完整性和延迟互相换 |
| 状态 | 作业内临时状态,跑完即弃 | 常驻状态,要 checkpoint 或 changelog |
| 重算 | 重跑作业,覆盖输出分区 | 从日志 offset 重放,受留存期限制 |
| 资源 | 按需启动,跑完释放 | 常驻进程,峰值容量要一直留着 |
| 典型失败 | 作业超时、上游分区没到齐就开跑 | 消费积压、watermark 卡住、重启后重复写 |
Micro-batch 位于中间:Spark Structured Streaming 不指定 trigger 时默认就是 micro-batch 模式,上一批处理完立即开始下一批。它的 API 看起来是流,延迟通常是秒级,内部仍是一批一批的有界计算。
两种混合做法
- Lambda architecture:同一份数据走两条路,流处理层给出低延迟的近似结果,批处理层定期用全量数据算出准确结果并覆盖。代价是同一套业务逻辑要在两个引擎里各写一遍,而且两边的结果很容易对不上。
- Kappa architecture:只保留流处理一条路,需要重算时,用新版本的作业从日志头部重放,追上后切流量,再下掉旧作业。代价是日志留存期要覆盖重算范围,重放期间要有两倍的算力。
面试里说「hybrid」时,要说清是哪一种,以及当批和流的结果不一致时以谁为准。一个常见且合理的答案是:面向用户的实时视图用流,财务、计费、对外报表以批为准,流的结果标注为「预估」。
有约束的设计问题:风控和日报
一家支付公司每天 2 亿笔交易,平均约 200,000,000 / 86,400 ≈ 2,315 笔/秒,峰值按均值 5 倍估约 11,600 笔/秒(倍数是假设)。两个需求:
- 风控:同一张卡 10 分钟内在 3 个以上国家消费就拦截,判断必须在授权返回前完成,授权的延迟预算里留给风控的只有几十毫秒。
- 商户日报:每个商户每天的交易额和手续费,第二天早上 8 点前给到,金额必须和清算一致。
风控只能是流:它是按卡号分区的 keyed state,每张卡保留最近 10 分钟的国家集合,状态量大约是「10 分钟内活跃卡数 × 几个国家」,可以放进算子本地状态。窗口按事件时间开,迟到事件不能直接丢,否则漏拦截;但授权已经返回的交易拦不回来,所以迟到事件只用于事后告警。
日报应该是批:凌晨等清算文件和所有迟到事件到齐后,读前一天的分区算一遍,按日期分区覆盖写。它和风控共用同一份原始事件,但不共用计算:日报用流来做,要处理跨午夜的迟到、作业重启时不重复累加、和清算对账,全部比批复杂,而换来的「更早看到」并没有人需要。
商户后台如果想看「今天到目前为止的交易额」,才加一条流作业给出实时预估,第二天被批结果覆盖,这就是上一节的 Lambda 式分工。
常见翻车
| 翻车 | 现象 | 修法 |
|---|---|---|
| 流作业按处理时间开窗 | 平时结果正常,一重放或一积压,所有事件落进同一个窗口 | 按事件时间开窗,重放才能得到同样的结果 |
| 默认 allowed lateness 为 0 | 离线上传的移动端事件被静默丢掉,统计偏低 | 配允许迟到加旁路输出,迟到数据单独计数 |
| 某个分区没有新数据,watermark 不前进 | 所有窗口都不关闭,结果停止输出 | 配置空闲分区检测,按分区监控 watermark 延迟 |
| 需要重算的范围超过日志留存 | 发现 bug 时原始事件已过期,无法重放 | 原始事件另存对象存储,批作业可以从那里回填 |
| 批作业在上游分区到齐之前启动 | 日报少了最后一小时 | 以上游写完的标记或分区就绪信号触发,而不是固定时间 |
| 流作业重启后重复写 sink | 计数翻倍 | sink 幂等(按键 upsert)或事务写入,见流处理架构一章 |
| 批和流算同一指标但口径不同 | 实时看板和日报差 3%,没人知道信谁 | 明确哪个是权威值,共享同一份口径定义 |
面试时这样回答
- 先问延迟需求来自谁:是用户的下一次操作(风控、推荐)还是人第二天看(报表)。前者才需要流。
- 说清时间语义:用事件时间,说出 watermark 和允许迟到,以及迟到数据去哪。
- 说清重算:批靠重跑分区;流靠从日志重放,说出留存期,并说明原始数据另存一份用于回填。
- 说清状态与正确性:流的状态怎么恢复、sink 怎么防重复,引到流处理架构的 checkpoint 或 changelog。
- 混合时说清谁是权威:Lambda 还是 Kappa,结果不一致时以批为准还是以流为准。
相关章节:流处理架构、Message Brokers、Event Sourcing、Latency vs Throughput。