分布式系统(32):事件时间、水位线与状态恢复
流处理不能只问“数据到了没有”。一条记录有业务发生时间、进入系统时间和实际执行时间;乱序让三者分离,持续流又没有天然结尾。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 | |
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。同一份业务需求若把故障从崩溃扩大到任意作恶,复制数量、消息验证和安全证明都会发生结构性变化。
参考资料
- Akidau et al., 2015, The Dataflow Model:event time、watermark、trigger与accumulation。
- Carbone et al., 2015, Lightweight Asynchronous Snapshots for Distributed Dataflows:barrier与一致快照。
- Apache Flink 2.3 Time:event/processing time、parallel watermarks、lateness。
- Apache Flink 2.3 Stateful Stream Processing:checkpoint、recovery、aligned/unaligned与保证范围。
- Apache Flink 2.3 Monitoring Back Pressure:输出buffer与反压指标。
