Batch vs Stream Processing

批处理与流处理怎么选:有界与无界输入、事件时间与处理时间、窗口与 watermark、迟到数据、重放和回填、micro-batch,以及 Lambda 与 Kappa 两种混合做法;附一个交易风控加日报的算例、翻车表和面试答法。

Batch processing 是把一段时间的数据先收集起来,再一次性处理。它更适合大规模、可延迟的计算。

Stream processing 则是数据一到就处理,强调 near real-time 响应。

batch-vs-stream

本章回答的是选哪一种、混合时怎么分工。流处理内部的四种拓扑(consumer group、changelog 状态、checkpoint 窗口、多消费者留存日志)和它们的故障恢复,在 流处理架构 里逐一展开。

真正的区别:输入有没有尽头

Google 2015 年的 Dataflow Model 论文(VLDB)建议别再用 batch / stream 描述引擎,而是描述数据:有界(bounded) 数据有尽头,作业读完就结束;无界(unbounded) 数据一直在来,作业永远不会「读完」。由此派生出三个真正要做的决定:

  1. 按哪个时钟算。 事件时间(event time)是事件发生的时间,处理时间(processing time)是它被处理的时间。手机离线 20 分钟再上传,两者就差 20 分钟。批处理天然按事件时间:昨天的分区里就是昨天的数据。流处理必须显式选。
  2. 什么时候算「够了」。 批处理的答案是「输入读完了」。流处理没有读完这一天,要靠 watermark 声明「事件时间 T 之前的数据应该都到了」,然后关闭窗口出结果。Flink 的 allowed lateness 默认是 0:落在 watermark 之后的元素直接丢弃,除非你配了允许迟到的时长或者旁路输出。
  3. 错了怎么重算。 批处理重跑一遍输入即可,前提是输入不可变、输出按分区覆盖写。流处理要从保留的日志重放,Kafka topic 的 retention.ms 默认 604800000,也就是 7 天:超过 7 天的错误,流里已经没有原始数据可重放了。

对比

维度BatchStream
输入有界:一个小时、一天的分区无界:持续到达的事件
端到端延迟调度周期加运行时长,通常分钟到小时秒级甚至亚秒级
正确性输入完整后再算,天然处理迟到靠 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%,没人知道信谁明确哪个是权威值,共享同一份口径定义

面试时这样回答

  1. 先问延迟需求来自谁:是用户的下一次操作(风控、推荐)还是人第二天看(报表)。前者才需要流。
  2. 说清时间语义:用事件时间,说出 watermark 和允许迟到,以及迟到数据去哪。
  3. 说清重算:批靠重跑分区;流靠从日志重放,说出留存期,并说明原始数据另存一份用于回填。
  4. 说清状态与正确性:流的状态怎么恢复、sink 怎么防重复,引到流处理架构的 checkpoint 或 changelog。
  5. 混合时说清谁是权威:Lambda 还是 Kappa,结果不一致时以批为准还是以流为准。

相关章节:流处理架构、Message Brokers、Event Sourcing、Latency vs Throughput。

一手证据