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 官方文档陈述,与本地切片的实测分开:

生产边界

动手

  1. 跑 exp05_fault.sh,把 --checkpoint-every 从 8 改成 3,重跑确定性崩溃那组,确认 duplicates 变小——解释为什么。
  2. 把幂等键从 (事件序号, 窗口结束) 改成 (处理时间戳, 窗口结束)(在 fault_engine.py 里改 key 的拼法),重跑,观察幂等 sink 的 duplicates 是否回到非 0,并解释为什么。
  3. 断言题:让基线跑完后把 final_windows 写成文件,崩溃恢复后再写一次,用 diff 断言两者相同。

自测

  1. checkpoint 里必须包含哪两类信息,缺了哪一类会出现「恢复后重复」或「恢复后丢数」?
  2. 为什么「恰好一次」在外部系统上做不到,除非外部写入支持事务或幂等?Flink 的 Kafka sink 是用哪种方式做的?
  3. 本机实验里,非幂等 sink 多算的是哪一条?它为什么会多算?
  4. 幂等键里为什么不能放 UUID 或处理时间?放进去会让哪一步失效?
  5. checkpoint 间隔从 8 调到 3,重复会变多还是变少?代价是什么?

↓ 下一步:05 章 · 数仓分层与建模 —— 算完的结果落到哪几层表里,宽表和预聚合换来什么。

进入 keel 阅读