Skip to content
Charles Shao
Go back

Flink 深挖 · 状态管理与 Checkpoint 容错

views

开篇讲清了 Flink 的运行时——但真正让 Flink 区别于「无状态流处理」的,是状态:去重、聚合、窗口、join、频控计数,全都要把中间结果存下来。而一旦有了状态,就必须回答一个问题:作业挂了、机器宕了,这些状态怎么不丢、怎么恢复到一致的那一刻? 这一篇只钻状态管理与 Checkpoint 容错。

一句话定位:Checkpoint 是 Flink 用「异步分布式快照」把无界流切出一致断点的机制;状态后端决定状态能存多大,barrier 对齐决定反压下 checkpoint 快不快。

TL;DR

Table of contents

Open Table of contents

1. 为什么有状态计算是核心

先分清两种算子:

广告实时链路里几乎全是有状态计算:实时 CTR 要累计曝光/点击、频控要记每个用户对每个广告的曝光次数、双流归因要把曝光和点击 join 起来。状态一旦存在,容错就成了硬需求:一个作业跑几天,中途某台 TaskManager 宕机是常态,不能因为一次宕机就把攒了几小时的状态清零、或算重/算漏。这正是 Checkpoint 要解决的。

2. Keyed State vs Operator State

Flink 的状态按作用域分两大类:

Keyed State vs Operator State 两种状态作用域对比:左侧 Keyed State 作用域是每个 key 独立,含 ValueState/ListState/MapState/ReducingState/AggregatingState;右侧 Operator State 作用域是每个并行实例,含 ListState(如 Kafka 各分区 offset)与 BroadcastState(动态规则下发)

Keyed State 随 key 切换视图、随 key 基数增长;Operator State 绑定并行实例、与 key 无关;rescale 时二者重分配方式不同。

Keyed State(作用域 = 每个 key):只能在 KeyedStreamkeyBy 之后)访问。Flink 按当前处理的 key 自动切换状态视图——你在代码里操作的永远是「当前 key 那一份」。原语:

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)**决定:

HashMapStateBackendEmbeddedRocksDBStateBackend
存储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

Checkpoint barrier 快照流程:JobManager 的 Checkpoint Coordinator 周期性触发 checkpoint n,Source 注入 barrier,barrier 随数据流向 Keyed Agg → Window → Sink 逐级传播,每个算子收到 barrier 时异步快照本地状态到 DFS 状态后端,完成后向 Coordinator 上报 ACK

barrier 从 source 注入、随流传播;算子处理到 barrier 时快照的状态,恰好是「消费了 barrier 之前所有数据」的一致断点——异步、不停流。

流程:

  1. 触发:JobManager 的 Checkpoint Coordinator 周期性触发 checkpoint n。
  2. 注入 barrier:source 在数据流里插入一个编号为 n 的 barrier(一个特殊标记记录),barrier 随普通数据一起向下游流动。
  3. 算子快照:每个算子收到 barrier n 时,异步把本地状态快照到 DFS(HDFS/S3/OSS),然后把 barrier n 转发给下游。
  4. ACK 与完成:算子快照完成后向 Coordinator 上报 ACK(带状态句柄);当所有算子都 ACK,该 checkpoint 才标记为「完成」,可用于恢复。

关键洞察:barrier 把无界流切成「checkpoint n 之前 / 之后」两段。算子处理到 barrier n 时快照的状态,恰好包含了「所有 n 之前的数据的影响、且不含任何 n 之后的数据」——这就是一致性的来源。整个过程异步进行、不停止数据处理

5. Barrier 对齐:Aligned vs Unaligned

当一个算子有多个输入通道(如 keyBy 后每个下游要收上游多个 subtask 的数据)时,各通道的 barrier n 不会同时到达。怎么处理决定了 checkpoint 在反压下快不快:

Aligned vs Unaligned Checkpoint 对比:Aligned 对齐模式下先到 barrier 的输入通道被阻塞、缓冲其后续数据,等所有输入的 barrier 都到齐才快照,反压下慢通道 barrier 迟到会导致对齐等待久、checkpoint 变慢甚至超时;Unaligned 非对齐模式下见首个 barrier 即刻快照并把 in-flight 缓冲数据一并纳入快照,反压下也快,代价是快照更大、恢复更复杂

Aligned 等齐所有 barrier 再快照,语义简单但反压下会被慢通道拖住;Unaligned 见首个 barrier 即快照、把 in-flight 数据一并存下,反压下更快但快照更大。

经验:正常低反压场景用 aligned;持续反压导致 checkpoint 频繁超时时,开启 unaligned checkpoint 往往能救急(反压根因的排查见反压定位与调优)。

6. Checkpoint vs Savepoint

两者底层机制相同(都是快照),用途不同:

CheckpointSavepoint
触发Flink 自动、周期性人工手动触发
用途故障自动恢复升级 / 迁移 / 回滚 / A-B
生命周期作业管理,通常滚动删除用户管理,长期保留
格式偏向快速恢复更稳定、可跨版本
增量支持(RocksDB)全量(标准 savepoint)

日常运维记住:改代码/升级 Flink 版本前先 flink stop --savepoint,从 savepoint 恢复;线上故障则由 checkpoint 自动兜底。

7. 状态 TTL 与 rescale

状态 TTLStateTtlConfig 让状态项在一段时间后自动过期清理(见 §2 代码)。对「按 key 累积」的状态必须设 TTL,否则 key 基数无限增长 → 状态无限膨胀 → checkpoint 越来越慢。RocksDB 可用 compaction filter 在后台合并时顺带清过期项。

rescale(改并行度):Flink 从 checkpoint/savepoint 恢复时可以改并行度:

参考


views
Share this post on:

Previous Post
Flink 深挖 · 反压定位与性能调优
Next Post
Flink 深挖 · DataStream API:ProcessFunction、状态、Timer 与 Side Output