开篇讲清了 Flink 的运行时——但真正让 Flink 区别于「无状态流处理」的,是状态:去重、聚合、窗口、join、频控计数,全都要把中间结果存下来。而一旦有了状态,就必须回答一个问题:作业挂了、机器宕了,这些状态怎么不丢、怎么恢复到一致的那一刻? 这一篇只钻状态管理与 Checkpoint 容错。
一句话定位:Checkpoint 是 Flink 用「异步分布式快照」把无界流切出一致断点的机制;状态后端决定状态能存多大,barrier 对齐决定反压下 checkpoint 快不快。
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 要累计曝光/点击、频控要记每个用户对每个广告的曝光次数、双流归因要把曝光和点击 join 起来。状态一旦存在,容错就成了硬需求:一个作业跑几天,中途某台 TaskManager 宕机是常态,不能因为一次宕机就把攒了几小时的状态清零、或算重/算漏。这正是 Checkpoint 要解决的。
2. Keyed State vs Operator State
Flink 的状态按作用域分两大类:
Keyed State 随 key 切换视图、随 key 基数增长;Operator State 绑定并行实例、与 key 无关;rescale 时二者重分配方式不同。
Keyed State(作用域 = 每个 key):只能在 KeyedStream(keyBy 之后)访问。Flink 按当前处理的 key 自动切换状态视图——你在代码里操作的永远是「当前 key 那一份」。原语:
ValueState<T>:单值,如去重标记、频次计数器;ListState<T>:列表,如窗口内元素缓冲;MapState<K,V>:映射,如按维度分桶计数;ReducingState<T>/AggregatingState<I,O>:增量归约/聚合,来一条并一条,省内存。
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 压力 | 大(状态越大越痛) | 小(在堆外) |
| 适用 | 状态较小、追求低延迟 | 状态很大(用户级/长窗口) |
选型直觉:小状态、低延迟(如短窗口聚合)用 HashMap;大状态(如按 userId 的频控、长时间窗口)用 RocksDB——RocksDB 把状态放堆外/磁盘,避开了大堆的 GC 停顿,还支持增量 checkpoint 让快照代价随「变更量」而非「总状态量」增长。
env.setStateBackend(new EmbeddedRocksDBStateBackend(true /* 增量 */));
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");
4. Checkpoint 原理:Chandy-Lamport 分布式快照
难点在于:数据在无界地流动、算子分布在多台机器上,怎么给这个「运动中的系统」拍一张一致的全局快照,还不能把流停下来?Flink 用的是 Chandy-Lamport 异步分布式快照算法,核心道具是 checkpoint barrier。
barrier 从 source 注入、随流传播;算子处理到 barrier 时快照的状态,恰好是「消费了 barrier 之前所有数据」的一致断点——异步、不停流。
流程:
- 触发:JobManager 的 Checkpoint Coordinator 周期性触发 checkpoint n。
- 注入 barrier:source 在数据流里插入一个编号为 n 的 barrier(一个特殊标记记录),barrier 随普通数据一起向下游流动。
- 算子快照:每个算子收到 barrier n 时,异步把本地状态快照到 DFS(HDFS/S3/OSS),然后把 barrier n 转发给下游。
- ACK 与完成:算子快照完成后向 Coordinator 上报 ACK(带状态句柄);当所有算子都 ACK,该 checkpoint 才标记为「完成」,可用于恢复。
关键洞察:barrier 把无界流切成「checkpoint n 之前 / 之后」两段。算子处理到 barrier n 时快照的状态,恰好包含了「所有 n 之前的数据的影响、且不含任何 n 之后的数据」——这就是一致性的来源。整个过程异步进行、不停止数据处理。
5. Barrier 对齐:Aligned vs Unaligned
当一个算子有多个输入通道(如 keyBy 后每个下游要收上游多个 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 | |
|---|---|---|
| 触发 | Flink 自动、周期性 | 人工手动触发 |
| 用途 | 故障自动恢复 | 升级 / 迁移 / 回滚 / A-B |
| 生命周期 | 作业管理,通常滚动删除 | 用户管理,长期保留 |
| 格式 | 偏向快速恢复 | 更稳定、可跨版本 |
| 增量 | 支持(RocksDB) | 全量(标准 savepoint) |
日常运维记住:改代码/升级 Flink 版本前先 flink stop --savepoint,从 savepoint 恢复;线上故障则由 checkpoint 自动兜底。
7. 状态 TTL 与 rescale
状态 TTL:StateTtlConfig 让状态项在一段时间后自动过期清理(见 §2 代码)。对「按 key 累积」的状态必须设 TTL,否则 key 基数无限增长 → 状态无限膨胀 → checkpoint 越来越慢。RocksDB 可用 compaction filter 在后台合并时顺带清过期项。
rescale(改并行度):Flink 从 checkpoint/savepoint 恢复时可以改并行度:
- 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 一致性快照算法的理论源头。