Apache Flink 状态后端与 Checkpoint 深度实战:从 Chandy-Lamport 分布式快照到 RocksDB 增量检查点
一、流计算的真正难点不是算得快,而是「记得住」
批处理的容错哲学是重算:任务挂了,把输入分区再读一遍,输出覆盖一遍,幂等天然成立。流计算没有这个奢侈——输入是无界的,你不能从头再放一遍三个月的交易流水。于是「状态」成了流计算里最昂贵、也最容易出事的东西:一个实时风控作业需要记住每个用户最近两小时的累计金额,一个会话归因作业需要记住每个会话的事件序列,一个双流 Join 需要记住两侧尚未匹配的行。
Flink 之所以能从众多流引擎里跑出来,核心不是吞吐比别人高多少,而是它把「分布式快照」做成了一种持续的、轻量的、与数据流同行的过程:checkpoint 不暂停计算,代价可控,且保证精确一次(exactly-once)。理解 checkpoint 与状态后端,是判断一个人是否真的用过 Flink 的分水岭。
二、State Backend:状态到底存在哪
Flink 1.15 之后有一个重要拆分,很多人至今没跟上:StateBackend 只负责「本地状态怎么存」,CheckpointStorage 才负责「快照存到哪」。以前这两件事混在 FsStateBackend / RocksDBStateBackend 里,现在是正交的两个概念。
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.enableCheckpointing(30_000, CheckpointingMode.EXACTLY_ONCE);
env.setStateBackend(new EmbeddedRocksDBStateBackend(true)); // true = 增量 checkpoint
env.getCheckpointConfig().setCheckpointStorage("s3://bucket/flink/ckpt");
CheckpointConfig ck = env.getCheckpointConfig();
ck.setMinPauseBetweenCheckpoints(10_000); // 两次 checkpoint 之间至少间隔
ck.setCheckpointTimeout(5 * 60_000);
ck.setTolerableCheckpointFailureNumber(3); // 容忍偶发失败,避免作业直接退出
两种本地后端的选择,本质是一场 trade-off:
| 维度 | HashMapStateBackend | EmbeddedRocksDBStateBackend |
|---|---|---|
| 存储位置 | JVM 堆内哈希表 | 堆外(RocksDB LSM-Tree) |
| 读写延迟 | 纳秒~微秒,无序列化 | 读写需序列化,有读放大 |
| 状态规模 | 受堆大小与 GC 约束 | 可远超内存,TB 级可行 |
| Checkpoint | 只能全量 | 支持增量 |
| 典型场景 | 小状态、极致延迟 | 大状态、长周期作业 |
经验分界线大概是 单算子状态 5~10GB:之下用堆内,之上必须 RocksDB。别被「RocksDB 更稳」误导,小状态用它纯属自找延迟。
三、状态原语:别把业务数据塞进成员变量
新手最常犯的错误是在 RichFunction 里用一个 HashMap 成员变量缓存业务数据。这在单并发下能跑,一旦并行度大于 1 或者发生故障恢复,数据就是错的——因为那块内存不归 Flink 管,不会进快照,恢复后是空的。
正确姿势是使用 Keyed State,它由 Flink 托管、自动快照、自动按 key 分区:
public class RiskCounter extends KeyedProcessFunction<String, Trade, Alert> {
private transient ValueState<Long> amount;
private transient ValueState<Long> lastUpdate;
@Override
public void open(Configuration cfg) {
ValueStateDescriptor<Long> desc = new ValueStateDescriptor<>("amount", Types.LONG);
desc.enableTimeToLive(StateTtlConfig.newBuilder(Time.hours(2))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build());
amount = getRuntimeContext().getState(desc);
lastUpdate = getRuntimeContext().getState(
new ValueStateDescriptor<>("lastUpdate", Types.LONG));
}
@Override
public void processElement(Trade t, Context ctx, Collector<Alert> out) throws Exception {
long cur = amount.value() == null ? 0L : amount.value();
cur += t.getAmountCents();
amount.update(cur);
if (cur > 1_000_000L) out.collect(new Alert(t.getUserId(), cur));
}
}
三个必须记住的事实:
- 只有
KeyedStream才有 Keyed State,状态与 key 绑定,key 基数直接决定状态规模。用户维度 1 亿 key,再小的状态也是灾难。 StateDescriptor的名字 + 类型序列化器 = 状态的身份。改了序列化器实现而不做兼容处理,savepoint 会直接恢复失败。- TTL 是惰性的,不是定时器。它靠读路径过滤 + RocksDB compaction filter 增量清理,所以过期数据不会立刻释放磁盘,只是「不再被看见」。
四、Chandy-Lamport:把「全局一致性」变成一次顺流而下的旅行
分布式系统最怕「全局一致快照」这个词——听起来就得暂停全世界。Flink 的做法恰恰相反:
- JobManager 的
CheckpointCoordinator向所有 Source 注入带checkpointId的 checkpoint barrier; - barrier 随数据一起在数据流里向前流动(它是 in-band 的,不需要额外通道);
- 算子收到某个输入通道的 barrier,就阻塞该通道并缓存后续数据,同时继续处理其他通道(这就是 barrier 对齐 / alignment);
- 等所有输入通道的 barrier 都到齐,算子快照本地状态,向下游广播 barrier,然后解除阻塞;
- 状态快照异步上传到远端存储,不阻塞数据处理;
- 所有算子 ack 后,coordinator 标记该 checkpoint 完成。
用伪代码表示核心逻辑:
void onBarrier(InputChannel ch, Barrier b) {
if (!aligned[ch.index]) {
aligned[ch.index] = true;
blockChannel(ch); // 缓存 barrier 之后的数据
if (allChannelsAligned()) {
Snapshot s = localState.deepCopy(); // 堆内需拷贝;RocksDB 只是文件引用
asyncWrite(s, checkpointStorage); // 异步上传
broadcastBarrierDownstream(b);
unblockAllChannels();
}
}
}
为什么这就够了一致性? barrier 把流切成两半:barrier 之前的数据已处理并计入快照,barrier 之后的数据在快照完成后才处理。因此每个算子的本地快照拼接起来,恰好对应某个「逻辑时刻」的完整全局状态——这正是 Chandy-Lamport 分布式快照算法在数据流通路上的实现。关键点在于:不需要暂停所有算子,只需要每个算子对齐自己的输入。
五、反压会拖垮对齐:Unaligned Checkpoint
生产环境最常见的 checkpoint 事故:某条通道反压,barrier 排在该通道几 GB 缓冲区数据之后,alignment 时间超过 checkpoint.timeout,checkpoint 连续失败,作业实际上失去了容错能力。
Flink 1.11 引入 Unaligned Checkpoint 解决这个问题:barrier 越过(overtake)缓冲区中的数据优先被处理,同时把「被越过的数据」和「in-flight 数据」一并写进快照。恢复时这些数据会被重放。代价是快照变大,收益是极端反压下 checkpoint 依然能完成。
env.getCheckpointConfig().enableUnalignedCheckpoints();
env.getCheckpointConfig().setAlignedCheckpointTimeout(Duration.ofSeconds(2)); // 2s 未对齐则降级
工程上推荐「对齐优先、超时降级」的混合模式:常态下走对齐路径保持快照体积,反压时自动退化到 unaligned 保住可用性。这是目前性价比最高的配置。
六、RocksDB 增量 checkpoint:只上传新增的 SST 文件
全量 checkpoint 意味着每次把整个 RocksDB 序列化上传——状态 200GB、每 30 秒一次,这显然不可能。增量 checkpoint 的巧妙之处在于借用了 LSM-Tree 的一个特性:SST 文件是不可变的。
上一轮已经上传过的 SST 文件,这一轮几乎不需要再传。Flink 只上传本轮新增文件,并在 _metadata 里记录文件引用关系;CompletedCheckpointStore 维护共享文件的引用计数,只有当某个文件不再被任何保留的 checkpoint 引用时,才在 subsuming 之后删除。
s3://bucket/flink/ckpt/chk-42/
_metadata # 元数据:算子 → 状态句柄 → 文件引用列表
shared/ # 跨 checkpoint 共享的 SST 文件(引用计数管理)
taskowned/ # 该算子独占、不共享的状态(如 RocksDB 的临时文件)
由此产生一个运维陷阱:不要手动删除 S3 上的旧 checkpoint 目录。 你删掉的 shared/ 文件很可能正被新 checkpoint 引用,会导致恢复时文件缺失。正确做法是用 Flink 自身的保留策略:
state.checkpoints.num-retained: 3
state.backend.incremental: true
state.backend.local-recovery: true # 本地恢复:节点软失败时秒级拉起
RocksDB 侧的常见调优:
EmbeddedRocksDBStateBackend rocks = new EmbeddedRocksDBStateBackend(true);
rocks.setRocksDBOptionsFactory((cfg, opts) -> {
opts.setWriteBufferSize(64 * 1024 * 1024);
opts.setMaxWriteBufferNumber(4);
opts.setLevelCompactionDynamicLevelBytes(true); // 减少空间放大
opts.setMaxBackgroundJobs(4); // 别抢光 CPU
});
七、Savepoint vs Checkpoint:一个给机器,一个给人
| 维度 | Checkpoint | Savepoint |
|---|---|---|
| 触发 | 自动、周期性 | 手动(或作业停止时) |
| 目的 | 故障自动恢复 | 版本升级、迁移、回滚、A/B |
| 生命周期 | 自动清理 | 人工管理,需显式删除 |
| 自包含性 | 依赖 shared 文件 | 自包含、可拷贝 |
状态演化(schema evolution)是升级时的硬骨头。Flink 1.17 之后支持 savepoint 内置的状态迁移,但改了 StateDescriptor 的类型或序列化器仍需显式处理。更稳妥的路线是用 State Processor API 离线重写:
SavepointWriter.newSavepoint(env, 128)
.withOperator(OperatorIdentifier.forUid("risk-counter"), transform)
.write("s3://bucket/sp-new");
八、端到端精确一次:状态对了,Sink 也要对
状态一致不等于端到端精确一次。外部系统必须在 checkpoint 完成的那一刻才真正提交,否则故障回滚时数据已经漏到了下游。Flink 通过 TwoPhaseCommitSinkFunction 扮演协调者:
public class ExactlyOnceSink extends TwoPhaseCommitSinkFunction<Txn, IN, Context> {
protected void invoke(Txn txn, IN value, Context ctx) { txn.producer.send(value); }
protected void preCommit(Txn txn) { txn.producer.flush(); }
protected void commit(Txn txn) { txn.producer.commitTransaction(); }
protected void abort(Txn txn) { txn.producer.abortTransaction(); }
}
coordinator 在 checkpoint completed 回调里 commit,失败时 abort。一个经典坑:外部系统的事务超时必须大于 checkpoint 间隔,否则事务被上游系统提前终止,commit 阶段直接报错。
九、生产调优清单(血泪版)
- 给每个有状态算子显式
uid()——DAG 变化时状态的身份靠 uid 而非算子下标,否则状态全部错位。 maxParallelism定好后不要改——key 到 key-group 的映射绑定它,改了 savepoint 直接作废。- checkpoint 间隔 1 分钟起步;重点盯
alignment duration和checkpointed data size,前者飙高说明反压已影响容错。 - 大状态(>100GB)必须是 RocksDB + 增量 + local recovery,三者缺一不可。
- 状态 TTL 务必配置,key 无限增长是生产事故的头号来源,且它只会在几周后才暴露。
- 版本升级前在影子环境做一次 savepoint 恢复演练,不要赌。
- RocksDB 目录挂 SSD,避免容器 overlayfs 与网络盘。
十、结语
状态管理是流计算从「玩具」走向「基础设施」的分水岭。当你能画出 barrier 在 DAG 上的完整旅行路径、能解释 RocksDB 共享文件的引用计数为何不能手动删、能说清 savepoint 的身份语义由什么决定,你就掌握了在故障面前让作业「像什么都没发生过一样继续」的能力——这才是有状态流计算真正的价值所在。

发表评论 取消回复