批处理与流式 — 按什么标准来选
一句话总结
批处理定期处理边界明确的数据集合,流处理持续处理无休止到达的事件;两者真正的差别不在速度,而在于由谁来划定边界。
为什么需要它
‘实时更好,所以用流处理吧’这样的决定经常出现,而代价通常会在六个月后找上门。批处理失败后重新运行即可;流处理则是带着状态持续运行的系统,重新处理的设计困难得多。
它如何工作
批处理的核心机制是水位线。保存已经处理到哪里,下次运行时只获取之后的数据。
SELECT * FROM orders
WHERE ordered_at > (SELECT last_ordered_at FROM etl_watermark WHERE job_name = 'orders_archive')
AND ordered_at < :batch_end;
这个简单模式中隐藏着多个陷阱。
- **是否包含边界值。**混淆
>和>=,会导致每次运行都重复或遗漏一条记录。 - **相同时间戳。**同一时刻写入多行时,只用
>会让部分数据永远遗漏。把时间与主键一起比较,或把水位线稍微往回移,再通过 upsert 吸收重复,通常更安全。 - **迟到的数据。**如果源系统很晚才提交事务,时间早于水位线的行可能随后才出现。此时,基于时间的水位线会永远错过该行。
在流处理中,这个问题更加明显。事件时间与处理时间不同,因此必须判断‘现在是否可以关闭这个窗口’,所以水位线要与允许延迟一起定义。对于迟到事件,也必须通过策略决定是丢弃还是重新打开窗口。
传递保证也分为三种:至多一次(可能丢失)、至少一次(可能重复)和恰好一次。大多数实际系统提供至少一次,并选择在消费端吸收重复,使结果看起来像只处理了一次。因此,下一个模块中的幂等性至关重要。
根据什么选择
‘实时更好’不是判断标准。考察以下四点,答案通常就会明确。
| 问题 | 适合批处理 | 适合流处理 |
|---|---|---|
| 何时需要结果 | 小时、天级 | 秒、分钟级 |
| 如何处理迟到数据 | 自然纳入下一批 | 需要设计水位线与重新处理 |
| 能否重新计算全部数据 | 容易 | 困难(恢复状态) |
| 运维负担 | 失败后重新运行即可 | 必须始终在线 |
最重要的是能否回退。批处理只需修复逻辑并重新运行昨天的数据。流处理若要做同样的事,就必须回退 offset、重置状态,并应对下游重复。
因此,实际工作中常见的答案是两者都要。用流处理快速给出近似值,再用批处理以准确值覆盖(Lambda 架构)。如今,使用同一套代码同时处理两者的方案(Kappa、Flink、Beam)越来越多,但流处理的运维复杂度依旧更高。
关于时间的三件事
流处理中一半的偏差都来自时间定义。
- 事件时间——实际发生的时刻。分析始终应以此为准。
- 采集时间——进入系统的时刻。
- 处理时间——执行计算的时刻。重新处理时它会改变,因此不能用作基准。
移动应用可能处于飞行模式,几小时后才发送事件。按事件时间聚合时,迟到数据会进入已经关闭的窗口(window)。水位线是‘认为这个时刻之前的数据不会再到达’的声明;越过该线才到达的数据会被丢弃或走单独路径处理。
워터마크 = 지금까지 본 최대 이벤트 시각 − 허용 지연(예: 10분)
增加允许延迟会提高准确度,但结果也会相应推迟。这是唯一一个在准确度与延迟之间进行权衡的旋钮。
批处理也可以增量运行
批处理并不需要每次读取全部数据。记录最后处理的位置,只读取之后的数据即可。不过,边界要有重叠。
-- 워터마크를 그대로 쓰면 경계에 걸친 것을 놓친다
where updated_at >= :last_watermark - interval '10 minutes'
and updated_at < :now
重叠区间会再次读取,但只要写入是幂等的,就不会造成任何问题。这里同样适用具备幂等性后,重叠读取的成本就等于零这一原则。
实际工作中常见的情况
批处理优于流处理的场景比想象中更多。如果源数据每天只更新一次,那么只有管道实时运行毫无价值。如果报表使用者每天早晨只查看一次,凌晨批处理就足够了。首先应问的是:消费者实际上多久查看一次?
反过来,批处理的代价也很明确。周期越长,一次失败导致的延迟就越大;单次处理量也更大,资源使用会形成尖峰。因此,常见做法是把大型批任务拆成小块处理。这样既可仅重新运行失败的数据块,也能缩短锁定时间。
下一项实验要做什么
把订单数据导出为 CSV 并写入归档表,记录水位线并执行增量写入,然后验证同一任务运行两次后结果是否保持不变。