LabHub
学习 学习路径 课程

数据流水线

批处理与流式 — 按什么标准来选

在 LabHub 中继续学习

一句话总结

批处理定期处理边界明确的数据集合,流处理持续处理无休止到达的事件;两者真正的差别不在速度,而在于由谁来划定边界

概念图: 由谁来划定边界 · 水位线 · 是否包含边界值。 · 相同时间戳。

为什么需要它

‘实时更好,所以用流处理吧’这样的决定经常出现,而代价通常会在六个月后找上门。批处理失败后重新运行即可;流处理则是带着状态持续运行的系统,重新处理的设计困难得多。

它如何工作

批处理的核心机制是水位线。保存已经处理到哪里,下次运行时只获取之后的数据。

SELECT * FROM orders
WHERE ordered_at > (SELECT last_ordered_at FROM etl_watermark WHERE job_name = 'orders_archive')
  AND ordered_at < :batch_end;

这个简单模式中隐藏着多个陷阱。

在流处理中,这个问题更加明显。事件时间与处理时间不同,因此必须判断‘现在是否可以关闭这个窗口’,所以水位线要与允许延迟一起定义。对于迟到事件,也必须通过策略决定是丢弃还是重新打开窗口。

传递保证也分为三种:至多一次(可能丢失)、至少一次(可能重复)和恰好一次。大多数实际系统提供至少一次,并选择在消费端吸收重复,使结果看起来像只处理了一次。因此,下一个模块中的幂等性至关重要。

根据什么选择

‘实时更好’不是判断标准。考察以下四点,答案通常就会明确。

问题 适合批处理 适合流处理
何时需要结果 小时、天级 秒、分钟级
如何处理迟到数据 自然纳入下一批 需要设计水位线与重新处理
能否重新计算全部数据 容易 困难(恢复状态)
运维负担 失败后重新运行即可 必须始终在线

最重要的是能否回退。批处理只需修复逻辑并重新运行昨天的数据。流处理若要做同样的事,就必须回退 offset、重置状态,并应对下游重复。

因此,实际工作中常见的答案是两者都要。用流处理快速给出近似值,再用批处理以准确值覆盖(Lambda 架构)。如今,使用同一套代码同时处理两者的方案(Kappa、Flink、Beam)越来越多,但流处理的运维复杂度依旧更高。

关于时间的三件事

流处理中一半的偏差都来自时间定义。

移动应用可能处于飞行模式,几小时后才发送事件。按事件时间聚合时,迟到数据会进入已经关闭的窗口(window)。水位线是‘认为这个时刻之前的数据不会再到达’的声明;越过该线才到达的数据会被丢弃或走单独路径处理。

워터마크 = 지금까지 본 최대 이벤트 시각 − 허용 지연(예: 10분)

增加允许延迟会提高准确度,但结果也会相应推迟。这是唯一一个在准确度与延迟之间进行权衡的旋钮

批处理也可以增量运行

批处理并不需要每次读取全部数据。记录最后处理的位置,只读取之后的数据即可。不过,边界要有重叠

-- 워터마크를 그대로 쓰면 경계에 걸친 것을 놓친다
where updated_at >= :last_watermark - interval '10 minutes'
  and updated_at <  :now

重叠区间会再次读取,但只要写入是幂等的,就不会造成任何问题。这里同样适用具备幂等性后,重叠读取的成本就等于零这一原则。

实际工作中常见的情况

批处理优于流处理的场景比想象中更多。如果源数据每天只更新一次,那么只有管道实时运行毫无价值。如果报表使用者每天早晨只查看一次,凌晨批处理就足够了。首先应问的是:消费者实际上多久查看一次?

反过来,批处理的代价也很明确。周期越长,一次失败导致的延迟就越大;单次处理量也更大,资源使用会形成尖峰。因此,常见做法是把大型批任务拆成小块处理。这样既可仅重新运行失败的数据块,也能缩短锁定时间。

下一项实验要做什么

把订单数据导出为 CSV 并写入归档表,记录水位线并执行增量写入,然后验证同一任务运行两次后结果是否保持不变。