流处理不能只问“数据到了没有”。一条记录有业务发生时间、进入系统时间和实际执行时间;乱序让三者分离,持续流又没有天然结尾。watermark负责声明事件时间推进到哪里,checkpoint负责保存恢复点,backpressure负责在下游变慢时限制上游。三个机制解决三个问题。

分布式系统(31):Ray 的 Future、对象血缘与重试边界

事件时间把结果归到业务时间轴

processing time取算子处理记录时的墙钟,延迟低,却会被排队、重启和机器速度改变。event time取记录携带的业务时间戳;同一批输入重放时仍能落入同一窗口,更适合账单、监控和会话分析。

sequenceDiagram
  participant S as Source
  participant O as Operator
  S->>O: event(ts=2)
  S->>O: event(ts=8)
  S->>O: watermark(6)
  S->>O: event(ts=4, late)

事件时间没有免费确定性。系统不能知道网络里是否还藏着更早记录,只能根据source策略发出watermark。Flink文档把 Watermark(t) 表述为:该流声称时间已经达到t,后续不应再有时间戳不大于t的记录。它是进度假设,不是物理事实。

watermark决定何时停止等待

窗口 [0,5) 收到 ts=2 后不能立即断言完整。watermark推进到6时,算子可以触发窗口结果。随后才到达的 ts=4 已经迟于当前事件时间,如何处理由lateness策略决定。

flowchart LR
  E2[ts=2] --> W0[window 0..5 sum=1]
  WM[watermark=6] --> Fire[首次触发]
  E4[迟到 ts=4] --> P{允许迟到?}
  P -->|否| Side[丢弃或side output]
  P -->|是| Update[更新窗口并再次发出]

本地模型使用整数时间,并把 [0,5) 的最大时间戳记为4。allowed lateness为0时,watermark 6已越过清理阈值,ts=4 被放进too-late side output,窗口结果保持1;允许3个时间单位时,清理阈值为7,watermark 6尚未达到它,窗口更新为2并发出late-update。正确性不只是“有没有迟到”,还包括sink能否处理更新、撤回或追加结果。

并行输入还需要取最慢进度。一个多输入算子的当前watermark通常取各输入watermark的最小值,否则快分区可能让窗口提前关闭,随后把慢分区的正常记录误判为迟到。空闲分区需要显式idle处理,否则会长期拖住全局进度。

状态让无界输入变成可计算问题

窗口和聚合需要记住尚未闭合的数据。keyed state与key分区共同移动,使同一key的更新由对应并行实例处理。状态可能包括窗口累加值、去重集合、定时器和模型参数。

状态大小不会因为流没有结尾而自动受限。窗口清理、TTL、迟到保留期和key基数决定其上界。把watermark设得很保守可少丢迟到事件,却延迟触发并延长状态寿命;设得激进则降低延迟,但增加迟到修正或丢弃。

checkpoint把输入位置与算子状态对齐

Flink checkpoint在source注入barrier。aligned模式下,多输入算子等到同一checkpoint的所有barrier后形成一致状态切面,再继续传播barrier。完成的checkpoint包含输入位置与各算子状态;unaligned模式还会保存in-flight数据。失败恢复时,状态回到该切面,source从对应位置重放。

flowchart LR
  S1[Source 1] -->|records + barrier k| O[Stateful operator]
  S2[Source 2] -->|records + barrier k| O
  O -->|snapshot state k| B[(State backend)]
  O -->|barrier k| K[Sink]

aligned checkpoint会暂时阻塞已收到barrier的输入,避免把checkpoint后的记录混进快照;这会在反压下放大对齐时间。unaligned checkpoint收到第一个barrier便开始快照,无需等待barrier alignment,并把被checkpoint越过的in-flight数据纳入快照,代价是更多I/O与更大快照。它们是不同的快照实现,不改变“状态与输入位置一致”的目标。

checkpoint完成不表示任意外部sink只执行一次。恢复会重放checkpoint后的记录;算子managed state回到旧值,checkpoint之后已发往普通HTTP或非事务数据库的效果不会自动撤回。

sequenceDiagram
  participant C as Checkpoint k
  participant O as Operator
  participant X as External sink
  C-->>O: state count=1, offset=1
  O->>X: effect(event 8) succeeds
  Note over O: crash before checkpoint k+1
  C-->>O: restore count=1, offset=1
  O->>X: replay effect(event 8)

端到端exactly-once至少要求source可重放,算子状态由一致checkpoint恢复,sink还能把写入与checkpoint完成协调起来,例如两阶段提交、事务批次或稳定幂等键。少一环就只能声明更窄的保证。

backpressure是流量反馈,不是丢弃策略

下游处理速度低于上游生产速度时,buffer逐渐占满。Flink把没有可用output buffer的时间计入back pressure,压力沿数据流反方向传播,让上游少生产或少读取。

flowchart RL
  Sink[慢sink 1条/步] -->|无空闲buffer| Op[operator]
  Op -->|反压| Source[快source 3条/步]
  Source -.数据方向.-> Op
  Op -.数据方向.-> Sink

本地模型给出18条固定输入,算子队列容量为4;前六步每步有3条进入source待发送区,sink每步只消费1条。五个步骤观察到反压,source侧最多积压9条;停止新增输入后,18条最终全部进入队列并被消费,丢弃数为0。这个结果只说明保留未接纳输入时反压可以限速,不证明无限过载可被消化;持续输入率高于处理率会增加等待、checkpoint对齐压力,最终还会碰到source自身的保留上限。

本地实验:乱序、恢复与反压

1
2
3
4
5
mkdir -p examples/distributed-systems/.build/stream32/tmp
export TMPDIR="$PWD/examples/distributed-systems/.build/stream32/tmp"
export TMP="$TMPDIR" TEMP="$TMPDIR" PYTHONDONTWRITEBYTECODE=1
python3 -B examples/distributed-systems/stream32/check.py \
--output examples/distributed-systems/.build/stream32/observations.json

Python 3.12.3正式运行覆盖两个lateness策略、一条checkpoint恢复时间线和一个有界buffer。输出分别记录checkpoint位点1、恢复位点1和重放后的最终位点2;managed count最终仍为2,外部sink却收到两次 effect-for-event-8。状态exactly-once与外部效果exactly-once因此被单独记录。

完整输出见结构化观察,研究和运行证据见实验证据,边界见验证说明。

安全性、活性与工程选择

事件时间窗口的安全性取决于时间戳提取、watermark策略和迟到更新语义一致。checkpoint安全性要求barrier切面、状态快照和source位置匹配;sink保证还要增加提交协议。活性要求source持续推进watermark、checkpoint最终完成、下游处理能力足够,不能让某个非idle输入永远停住最小watermark。

需求 关键状态 常见代价
结果按业务时间稳定 event timestamp、watermark 等待乱序数据
接收更多迟到数据 allowed lateness、更新sink 更长状态寿命、结果修订
故障后状态一致 source offset、operator checkpoint 快照I/O、barrier对齐
外部效果不重复 sink事务或幂等键 协调与提交延迟
防止buffer无界增长 backpressure与容量 上游吞吐下降、排队延迟

两个推演练习

双输入watermark。 输入A的watermark是100,输入B是40。窗口 [40,50) 能否因为A已到100而关闭?

不能。算子进度受最小输入watermark 40限制;除非B被正确标记idle或其watermark继续推进,否则提前关闭会误判B后续合法记录。

checkpoint后外部写成功。 source offset 10与状态S进入checkpoint,随后处理offset 10并成功写HTTP,下一checkpoint前崩溃。恢复后会怎样?

source从10重放,managed state从S继续,因此内部状态可保持一致;HTTP效果可能再次发生。只有sink幂等或参与checkpoint提交,才能扩大保证范围。

工程结论

watermark决定“等到什么时候”,state保存“已经算到什么”,checkpoint决定“失败后从哪里恢复”,backpressure决定“下游跟不上时怎样减速”。任何一个都不能替另外三个工作,更不能单独承诺外部副作用恰好一次。

下一篇进入PBFT。同一份业务需求若把故障从崩溃扩大到任意作恶,复制数量、消息验证和安全证明都会发生结构性变化。

参考资料