KEEL · 龙骨 · A CURRICULUM FOR THE AI ERA
04 · 状态与容错 — keel 龙骨
上一章的窗口状态全在内存的 dict 里。进程一崩,sum 就没了。这一章讲状态放哪、多久存一次快照、重启后从哪继续,以及「恰好一次」到底是不是免费的。
上一章的窗口状态全在内存的 dict 里。进程一崩,
sum就没了。这一章讲状态放哪、多久存一次快照、重启后从哪继续,以及「恰好一次」到底是不是免费的。
现场
流处理进程跑了一个小时,正在累计当天的窗口。运维手一抖 kill -9。重启之后会发生两件事之一:要么窗口从零开始重算(报表数字先跌后涨,运营看到断崖),要么从某个快照继续(数字平滑,但可能重复计算了一小段时间)。
要保住「不重复也不丢」,绕不开两个问题:状态存在哪、它和外面的写入怎么对齐。第二个问题才是「恰好一次」真正难的地方——不是引擎自己的状态对不对,是引擎写完状态和往外部系统写完数据这两件事,没法用一个原子操作罩住。
全链路
flowchart TD
subgraph ENG["流处理进程"]
direction TB
IN["事件流<br/>ev05.jsonl"] --> ST["窗口状态<br/>dict{w_end: {sum,cnt,fired}}"]
ST --> EMIT["按触发规则输出<br/>(w_end, sum, key)"]
ST -->|"每 K 个事件"| SNAP["checkpoint<br/>state-*.json<br/>index + max_ts + windows"]
end
EMIT --> SK{"sink 类型"}
SK -->|"非幂等"| APP["直接 append<br/>sink-*.jsonl"]
SK -->|"幂等(按键去重)"| DEDUP["查已写 key<br/>key=(事件序号,窗口结束)"]
DEDUP --> APP
CRASH(["kill -9 / 模拟崩溃<br/>停在两次快照之间"]) -. 崩溃 .-> ENG
SNAP -. "重启:从最后一份快照继续" .-> RESUME["重放 index 之后的事件"]
RESUME -. "重放会重新输出<br/>上次快照之后的那些窗口" .-> APP
APP -. "同一个 key 写两遍" .-> DUP["下游多算 1 条"]
style IN fill:#e8f0fe,stroke:#3b5bdb,color:#1a1a2e
style ST fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style EMIT fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style SNAP fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style SK fill:#fff4e6,stroke:#e8590c,color:#1a1a2e
style APP fill:#f3f0ff,stroke:#7048e8,color:#1a1a2e
style DEDUP fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style CRASH fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
style RESUME fill:#e6fcf5,stroke:#0ca678,color:#1a1a2e
style DUP fill:#ffe3e3,stroke:#c92a2a,color:#1a1a2e
精确定义
**状态(state)**是跨事件累积的东西:窗口的累计值、去重的键集合、会话的区间。它是流处理区别于无状态转发的根本。
**checkpoint(快照)**是状态在某个一致性点上的副本,加上「已经处理到输入流的哪个位置」。恢复时两者一起还原,才能接着算。
状态后端决定状态存在哪:内存后端快但状态量受限于堆;RocksDB 后端把状态放磁盘、支持超过内存的量,代价是访问要序列化。
至少一次(at-least-once):每条记录至少被处理一次。故障重放会导致重复。
恰好一次(exactly-once):对引擎内部状态而言是「每条记录的结果只生效一次」;对外部系统而言,只有在外部写入支持事务或幂等时才能做到。
幂等键:把「同一条逻辑写入」映射到一个稳定标识,重放时用它去重。本章用的是 (事件序号, 窗口结束)。
一次完整运行
原始输出见 lab/evidence/stream-data-pipeline/05-fault-tolerance.txt,代码在 lab/exp/stream-data-pipeline/fault_engine.py。
先拿基线(不中断跑完):
processed_upto_index=30 written=4
final_windows={"10": 55, "20": 155, "30": 255}
sink:
{"key": "13:10", "window_end": 10, "sum": 55, "watermark": 10, "kind": "fire"}
{"key": "24:20", "window_end": 20, "sum": 138, "watermark": 21, "kind": "fire"}
{"key": "25:20", "window_end": 20, "sum": 155, "watermark": 21, "kind": "late-update"}
{"key": "flush:30", "window_end": 30, "sum": 255, "watermark": 24, "kind": "flush"}
窗口 [20,30) 出现了两条记录:先 24:20 触发时是 138,之后一条迟到事件到达(25:20)把它更新成 155。这就是第 03 章说的「一个窗口触发多次」。
再真跑一次崩溃:每 8 个事件存一次快照,跑 2.6 秒后 kill -9:
崩溃时 sink 已写 = 0 条
崩溃时快照进度:index=8 max_ts=9 windows={"10": 40}
进程被杀,快照停在 index=8。重启从快照继续:
resume_from_index=8 watermark=4
processed_upto_index=30 written=4
final_windows={"10": 55, "20": 155, "30": 255}
BASE final_windows={"10": 55, "20": 155, "30": 255}
RUN final_windows={"10": 55, "20": 155, "30": 255}
=> 一致:快照 + 重放没有算错最终结果
BASE 和 RUN 完全相同。崩溃恢复的正确性不是「数字看起来对」,是「和不中断跑出来的结尾一模一样」。 想把这句话变成回归测试,就比对这两行的字符串。
至少一次与恰好一次:差多少
上面那次崩溃恰好落在快照边界上,所以没产生重复。要量化重复,得让崩溃点落在「已经输出、但还没被下一次快照覆盖」的区间。用确定性崩溃(处理到第 14 个事件后直接退出):
非幂等 sink:
simulated_crash_after_index=14 written=1
崩溃时快照停留在 index=8
恢复后:written=4 final_windows={"10": 55, "20": 155, "30": 255}
total_records=5 unique_keys=4 duplicates=1
重复的键:2 次 "key": "13:10"
幂等 sink(key=(事件序号,窗口结束) 去重):
崩溃时 written=1
恢复后:written=3 skipped=1
total_records=4 unique_keys=4 duplicates=0
把两组数字并列:
| 方案 | 崩溃后写入 | 恢复后写入 | 总记录 | 去重后 | 多算 |
|---|---|---|---|---|---|
| 非幂等 sink | 1 | 4 | 5 | 4 | 1 条(键 13:10) |
| 幂等 sink | 1 | 3(跳过 1) | 4 | 4 | 0 |
同一份输入、同一个崩溃点,区别只是 sink 认不认幂等键,下游就多算 1 条。 这个 1 看起来无关紧要,但把它换成「每 8 个事件一次快照、每秒几万条事件」,重复的规模就是「快照间隔内的全部输出」。
还有一个容易被忽略的事实:幂等 sink 在崩溃时写进的那 1 条,恢复后没有重复写(skipped=1)。它靠的不是「不重放」,而是「重放了但认出来是同一个键」。想做成这样,输出记录里必须带一个跨重启稳定的键——本机用的是 (事件序号, 窗口结束),因为事件序号在重放时不会变。如果键里有随机的 UUID 或处理时间戳,重放就会产生新键,去重立刻失效。
官方文档怎么规定(检索日期 2026-10-05)
以下为 Apache Flink 官方文档陈述,与本地切片的实测分开:
- Flink 的 checkpoint 机制要求「可重放的数据源」加「持久化的状态存储」,默认把快照放 JobManager 堆内存,生产建议放持久化文件系统(Flink 1.20 Checkpointing 文档)。本机切片用的是本地 JSON 文件,位置等价于「持久化存储」。
enableCheckpointing(n)可以选EXACTLY_ONCE或AT_LEAST_ONCE;文档明确说 at-least-once 适用于「追求极低延迟(几毫秒级)」的场景(同上)。这与本地「幂等 sink 换来的是不重复、付出的是每次写入都要查重」是同一个取舍。- Flink 的 Kafka sink
DeliveryGuarantee.EXACTLY_ONCE会把消息写进一个 Kafka 事务,在 checkpoint 时提交;消费者只有读已提交的数据(isolation.level)才不会看到重复。文档同时指出代价:记录在 checkpoint 写完之前对外不可见,并要求transactionalIdPrefix全局唯一,还推荐把 Kafka 的transaction.timeout.ms调到大于「最大 checkpoint 时长 + 最大重启时长」(Flink 1.20 Kafka connector 文档)。 - Flink 的 Kafka source 文档明确:source 不依赖已提交位点做容错,位点提交只用于暴露消费进度给监控(同上)。这一条和本地第 02 章看到的现象对照很有意思——Connect 的 sink 位点是真的用来恢复的,而 Flink source 的位点提交只是为了可观测:它的恢复靠的是 checkpoint 里的状态。
生产边界
- 教学替身 vs 真实依赖:本机快照是「同步写文件 + 原子 rename」。真实引擎的 checkpoint 是异步的、分布式的,还会做 barrier 对齐;状态量大时要用增量快照。本机的 sink 是本地文件,真实 sink 要么支持事务(Kafka 事务、两阶段提交的 JDBC),要么必须自带幂等键。
- 要盯的指标:checkpoint 成功率与时长、两次 checkpoint 的间隔、状态大小、sink 的幂等键冲突率(这个数从 0 变正,说明在重复投递)。
- 失败策略:恢复必须从最后一份完整快照开始,不能用「一半写坏」的快照;写快照要原子(临时文件 + rename),否则崩在写的过程中会得到一份不可读的状态。
- 快照间隔是取舍:间隔越短,恢复后要重放的事件越少、重复越少,但正常运行时 checkpoint 的开销越大。本机每 8 个事件一次;真实系统按「可接受的重放时长」反推间隔。
动手
- 跑
exp05_fault.sh,把--checkpoint-every从 8 改成 3,重跑确定性崩溃那组,确认duplicates变小——解释为什么。 - 把幂等键从
(事件序号, 窗口结束)改成(处理时间戳, 窗口结束)(在fault_engine.py里改key的拼法),重跑,观察幂等 sink 的duplicates是否回到非 0,并解释为什么。 - 断言题:让基线跑完后把
final_windows写成文件,崩溃恢复后再写一次,用diff断言两者相同。
自测
- checkpoint 里必须包含哪两类信息,缺了哪一类会出现「恢复后重复」或「恢复后丢数」?
- 为什么「恰好一次」在外部系统上做不到,除非外部写入支持事务或幂等?Flink 的 Kafka sink 是用哪种方式做的?
- 本机实验里,非幂等 sink 多算的是哪一条?它为什么会多算?
- 幂等键里为什么不能放 UUID 或处理时间?放进去会让哪一步失效?
- checkpoint 间隔从 8 调到 3,重复会变多还是变少?代价是什么?
↓ 下一步:05 章 · 数仓分层与建模 —— 算完的结果落到哪几层表里,宽表和预聚合换来什么。