凌晨三点,一个连跑两天的频控作业因为 TaskManager 被 OOMKill 而重启,恢复后频次计数从零开始,当天上午同一个用户被同一支广告反复触达了几十次,客诉和超投一起找上门。开篇讲清了 Flink 的运行时,但真正让它区别于无状态流处理的是状态:去重、聚合、开窗、双流 join、频控计数,中间结果全都要存下来。而状态一旦存在,就必须回答那个要命的问题:作业挂了、机器宕了,这些状态怎么不丢、又怎么恢复到一致的那一刻? 这一篇只钻状态管理与 Checkpoint 容错。
一句话定位:Checkpoint 是 Flink 用异步分布式快照(Chandy-Lamport)把无界流切出一致断点的机制——状态后端决定状态存得下多大,barrier 对齐方式决定反压下快照还能不能按时打出来,TTL 与 key group 则决定这份状态能不能长期跑下去。
TL;DR
- 有状态是 Flink 的核心能力:去重、聚合、窗口、双流 join、频控计数都要状态;无状态处理任何框架都能做,有状态且容错才是 Flink 的护城河。
- 两类状态:Keyed State(作用域 = 每个 key 独立,只能在
keyBy之后用)与 Operator State(作用域 = 每个并行实例,与 key 无关,典型如 Kafka source 的分区 offset)。 - Keyed State 原语:
ValueState/ListState/MapState/ReducingState/AggregatingState;Flink 按当前 key 自动切换状态视图,你写的永远是「当前 key 那一份」。 - 状态后端二选一:
HashMapStateBackend(堆内,快但受内存与 GC 限制)vsEmbeddedRocksDBStateBackend(堆外 + 磁盘,撑得住超大状态、支持增量 checkpoint,代价是序列化与读写开销)。 - Checkpoint 原理:基于 Chandy-Lamport 异步分布式快照——JobManager 触发、source 注入 barrier、barrier 随数据流向下游、每个算子收到 barrier 时异步快照本地状态到 DFS,全程不停流。
- barrier 对齐:aligned(多输入等齐所有 barrier 才快照,反压下会被慢通道拖住)vs unaligned(见首个 barrier 即快照、把 in-flight 数据一并纳入,反压下更快但快照更大)。
- Checkpoint vs Savepoint:Checkpoint 由 Flink 周期性自动做、用于故障自动恢复;Savepoint 手动触发、用于升级/迁移/回滚,格式更稳定。
- 状态 TTL 是保命开关:按 key 累积的状态若不设 TTL 就会无限增长,最终拖垮 checkpoint。
- rescale 靠 key group:改并行度时 keyed state 按 key group 重分配、operator list state 轮转重分、broadcast state 全量复制。
- AdTech 落地:实时频控、曝光去重、按 campaign 聚合都靠 keyed state;用户级大状态用 RocksDB + 增量 checkpoint,并且务必设 TTL。
Table of contents
Open Table of contents
1. 为什么有状态计算是核心
先把算子分成两类:
- 无状态(stateless):
map/filter这类算子,每条记录的输出只取决于它自己,不依赖任何历史; - 有状态(stateful):像「统计每个 campaign 过去 1 分钟的曝光数」「判断这个 requestId 是否已经处理过」,输出依赖之前见过的数据,中间结果必须存下来。
广告实时链路几乎清一色落在后者:实时 CTR 要累计曝光与点击,频控要记住每个用户对每支广告的曝光次数,双流归因则要把曝光和点击各自缓在状态里等对方到齐。状态一旦存在,容错就从加分项变成硬需求——作业连跑几天,中途某台 TaskManager 宕机、被抢占、被 OOMKill 都是常态,不能因为一次宕机就把攒了几小时的状态清零,更不能在恢复之后算重或算漏。Checkpoint 要解决的正是这件事。
不过在讨论状态怎么存、怎么不丢之前,得先分清 Flink 到底有几种状态:它们的作用域不一样,改并行度时的重分配方式也不一样。
2. Keyed State vs Operator State
Flink 的托管状态按作用域分成两大类,下图先给出全貌:
Keyed State 随 key 切换视图、随 key 基数增长;Operator State 绑定并行实例、与 key 无关;rescale 时二者的重分配方式截然不同。
Keyed State(作用域 = 每个 key)只能在 KeyedStream(也就是 keyBy 之后)访问:Flink 会按当前处理的那条记录的 key 自动切换状态视图,所以你在算子里读写的永远是「当前 key 那一份」,不需要自己维护一张按 key 索引的 map。它提供五种原语:
ValueState<T>:单值,如去重标记、频次计数器;ListState<T>:列表,如窗口内的元素缓冲;MapState<K,V>:映射,如按维度分桶计数;ReducingState<T>/AggregatingState<I,O>:增量归约与聚合,来一条并一条,省内存。
落到频控场景,一个 KeyedProcessFunction 就够用了。留意 open 里那几行 TTL 配置——它不是可选项,而是 §7 会反复强调的保命开关:
public class FreqCapFunction extends KeyedProcessFunction<String, Impression, Alert> {
private transient ValueState<Long> countState;
@Override
public void open(Configuration cfg) {
ValueStateDescriptor<Long> desc =
new ValueStateDescriptor<>("imp-count", Long.class);
// 关键:设置 State TTL,避免状态无限增长
desc.enableTimeToLive(StateTtlConfig
.newBuilder(Time.hours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.cleanupInRocksdbCompactFilter(1000)
.build());
countState = getRuntimeContext().getState(desc);
}
@Override
public void processElement(Impression imp, Context ctx, Collector<Alert> out) throws Exception {
long c = Optional.ofNullable(countState.value()).orElse(0L) + 1;
countState.update(c); // 操作的就是「当前 key」那一份
if (c > 5) out.collect(new Alert(ctx.getCurrentKey(), c));
}
}
Operator State(作用域 = 每个并行实例)则与 key 无关,直接绑在算子的某个并行子任务上。最常见的是 ListState——Kafka source 就是用它记住各分区的消费 offset,故障后能从断点重放,靠的正是这份状态;其次是 BroadcastState,广播流在每个实例上持有同一份只读副本,适合把风控规则、出价配置这类小而多变的数据下发到全部并行实例。
作用域说清楚了,下一个问题随之而来:这些状态究竟躺在哪里,堆内、堆外,还是本地磁盘?
3. 状态后端:HashMap vs RocksDB
状态存在哪里、以什么形式存,由**状态后端(State Backend)**决定,它实际上划定了一个作业的状态上限:
HashMapStateBackend | EmbeddedRocksDBStateBackend | |
|---|---|---|
| 存储 | JVM 堆内(对象形式) | 堆外 + 本地磁盘(RocksDB,序列化字节) |
| 读写速度 | 快(直接操作对象) | 较慢(要序列化 + 可能读磁盘) |
| 状态上限 | 受堆内存/GC 限制 | 可达 TB 级(磁盘) |
| 增量 checkpoint | 不支持(全量) | 支持(只传变更的 SST) |
| GC 压力 | 大(状态越大越痛) | 小(在堆外) |
| 适用 | 状态较小、追求低延迟 | 状态很大(用户级/长窗口) |
选型其实没多少纠结空间。状态小、追求低延迟的短窗口聚合交给 HashMapStateBackend,对象直接躺在堆里,一次读写就是一次引用寻址;而按 userId 维度的频控、跨小时的长窗口这类大状态,只能上 EmbeddedRocksDBStateBackend——它把状态序列化到堆外与本地磁盘,既避开了大堆的 GC 停顿,又能开启增量 checkpoint,让每轮快照的代价随「这一轮变更了多少」而不是「一共攒了多少」增长。代价同样实在:每次读写都要走一遍序列化,热 key 上的 ValueState 更新会明显慢于堆内。
开启只需要两行,快照落到哪套 DFS 也在这里指定:
env.setStateBackend(new EmbeddedRocksDBStateBackend(true /* 增量 */));
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");
状态存得下,只解决了一半问题;另一半是它得丢不了。
4. Checkpoint 原理:Chandy-Lamport 分布式快照
难点在于:数据在无界地流动、算子分散在几十台机器上,怎么给这个「运动中的系统」拍出一张一致的全局快照,还不能把流停下来?Flink 的答案是 Chandy-Lamport 异步分布式快照算法,而落地这套算法的关键载体只有一个:checkpoint barrier。
barrier 从 source 注入、随流逐级传播;算子处理到 barrier 时快照的状态,恰好是「消费了 barrier 之前所有数据」的一致断点——异步进行,不停流。
一轮完整的 checkpoint 分四步:
- 触发:JobManager 的 Checkpoint Coordinator 周期性触发 checkpoint n。
- 注入 barrier:source 在数据流里插入一个编号为 n 的 barrier(一条特殊的标记记录),它随普通数据一起向下游流动。
- 算子快照:每个算子收到 barrier n 时异步把本地状态快照到 DFS(HDFS/S3/OSS),随即把 barrier n 转发给下游,不必等快照落完。
- ACK 与完成:快照落盘后算子向 Coordinator 上报 ACK(携带状态句柄),只有所有算子都 ACK,这一轮 checkpoint 才被标记为完成、可用于恢复。
关键洞察在于 barrier 把无界流切成了「checkpoint n 之前」与「之后」两段:算子处理到 barrier n 时快照下来的状态,恰好包含了 n 之前所有数据的影响、且不含任何 n 之后的数据,一致性正来源于此。而这一切都在后台异步完成,数据处理并不停止。
但这套机制有个前提平时不显眼:barrier 得走得动。一旦算子有多个输入通道,而其中一条被反压拖慢,麻烦就来了。
5. Barrier 对齐:Aligned vs Unaligned
一个算子有多个输入通道是常态——keyBy 之后每个下游 subtask 都要接收上游所有 subtask 发来的数据,而编号相同的 barrier n 不会同时到达。怎么处理这段时间差,直接决定了反压之下 checkpoint 还能不能按时打出来:
Aligned 等齐所有 barrier 再快照,语义简单但反压下会被慢通道拖住;Unaligned 见首个 barrier 即快照、把 in-flight 数据一并存下,反压下更快但快照更大。
- Aligned(对齐,默认):先到 barrier n 的通道被阻塞并缓冲其后续记录,等所有输入通道的 barrier n 都到齐后才快照、再放行。语义清晰,但反压下慢通道的 barrier 迟迟不到,对齐等待越拖越长,checkpoint 随之变慢甚至超时。
- Unaligned(非对齐):算子见到第一个 barrier n 就立即快照本地状态,并把此刻在途(in-flight)的输入/输出缓冲数据一并纳入快照,barrier 则直接越过缓冲跳到输出端。反压下照样能按时完成,代价是快照体积更大、恢复过程更复杂。
经验法则很简单:反压不明显时用默认的 aligned;一旦出现持续反压导致 checkpoint 频繁超时,开启 unaligned checkpoint 往往能先把作业救回来。但它只是止血,反压根因还得另行排查(见反压定位与调优)。
6. Checkpoint vs Savepoint
快照按时打出来之后,接下来的分歧在于这份快照给谁用。Checkpoint 与 Savepoint 底层机制相同,用途却完全不同:
| Checkpoint | Savepoint | |
|---|---|---|
| 触发 | Flink 自动、周期性 | 人工手动触发 |
| 用途 | 故障自动恢复 | 升级 / 迁移 / 回滚 / A-B |
| 生命周期 | 作业管理,通常滚动删除 | 用户管理,长期保留 |
| 格式 | 偏向快速恢复 | 更稳定、可跨版本 |
| 增量 | 支持(RocksDB) | 全量(标准 savepoint) |
日常运维的分工可以记成一句话:计划内的动作走 savepoint,计划外的意外交给 checkpoint。改代码、调并行度、升级 Flink 版本之前先 flink stop --savepoint 停稳再从 savepoint 拉起;线上半夜宕机则由周期性 checkpoint 自动兜底,不需要人工介入。
7. 状态 TTL 与 rescale
上面几节解决的是状态存得下、丢不了、打得出来,但还有两笔账要在上线几周后才结:状态会不会越涨越大,以及并行度还改不改得动。
状态 TTL:StateTtlConfig 让状态项在存活一段时间后自动过期清理(配置见 §2 的代码)。凡是「按 key 累积」的状态必须设 TTL,否则 key 基数只增不减,状态持续膨胀,checkpoint 的耗时与体积跟着一路上扬,最终把作业自己拖垮。用 RocksDB 时还可以打开 compaction filter,让后台合并顺带回收过期项,省去额外的扫描开销。
rescale(改并行度):从 checkpoint 或 savepoint 恢复时,Flink 允许换一个并行度重新展开作业,但三类状态的重分配方式各不相同:
- Keyed State:按 key group(key 的逻辑分片,数量 =
maxParallelism)整片迁移给新的并行实例——这也解释了为什么maxParallelism一旦设定就不宜再改。 - Operator List State:把列表元素在新实例之间轮转重分(round-robin)。
- Broadcast State:全量复制到每一个新实例。
参考
- Apache Flink. Working with State(Keyed / Operator State、状态原语、StateTtlConfig):状态类型、原语与 TTL 配置的官方指南。
- Apache Flink. Checkpointing(开启与配置一致性快照):checkpoint 间隔、语义模式、超时与并发的配置说明。
- Apache Flink. State Backends(HashMap vs RocksDB、增量 checkpoint):状态后端选型与 RocksDB 调优。
- Apache Flink. Checkpointing under Backpressure(Unaligned Checkpoints):aligned / unaligned 快照的取舍与配置。
- Chandy, K. M., Lamport, L. Distributed Snapshots: Determining Global States of Distributed Systems:Flink 一致性快照算法的理论源头。