Flink(三):状态管理与 Checkpoint 容错
Flink(三):状态管理与 Checkpoint 容错
导语:状态是 Flink 的灵魂,Checkpoint 是状态的保险。这一篇是 Flink 面试区分度最高的一章:Keyed/Operator State 的选型、StateBackend 的取舍、Barrier 对齐与异步快照的完整流程、非对齐 Checkpoint 为什么"用空间换时间"、端到端 Exactly-Once 的三段式、重启策略与 State TTL,共 15 题。
一、状态(State)
1. Flink 中的状态分哪几类?怎么选?
答: 按「是否绑定 key」分两大类,这是最基础也最重要的划分:
| 对比项 | Keyed State | Operator State |
|---|---|---|
| 使用前提 | 必须在 KeyedStream 上(keyBy 之后) | 任意算子(含 Source / Sink) |
| 状态粒度 | 每个 key 一份 | 每个并行子任务一份 |
| 获取方式 | RuntimeContext.getState(...) / KeyedProcessFunction | 实现 CheckpointedFunction(或 ListCheckpointed,已废弃) |
| 扩容/缩容(改并行度) | 按 key 自动重分布(hash(key) 重分区),天然支持 | 需自己选 EvenSplit(均分,可能丢语义)或 Union(每子任务拿全量) |
| 典型用途 | 聚合累加、去重集合、会话、join 缓存 | Kafka offset、Sink 事务状态、广播配置 |
选型口诀:
- 「按 key 统计」→ Keyed State(99% 的场景);
- 「与 key 无关、每个实例一份」→ Operator State(Source 的 offset、Sink 的事务状态);
- 「每个实例都要全量」→
BroadcastState(维表/规则全量广播)。
2. ValueState / ListState / MapState / ReducingState / AggregatingState 各自适用什么?
答:
| 类型 | 存什么 | 典型场景 | 注意 |
|---|---|---|---|
ValueState<T> | 单个值 | 计数、上次状态、简单标志 | 未赋值时 value() 返回 null |
ListState<T> | 一个列表 | 需要保留多条记录(如 join 侧的缓存、TopN 候选) | 只能用 add/addAll 追加、update 覆盖,没有按 index 删除,要删除得整体重写 |
MapState<K,V> | 一个 Map | 去重集合、按维度统计(userId → 次数) | 支持 put/remove/contains/iterate,是「可增删的集合」的首选 |
ReducingState<T> | ReduceFunction 增量聚合 | 累加求和、求最大 | 输入输出同类型;每次 add 都会合并 |
AggregatingState<I,O> | AggregateFunction 增量聚合 | 输入输出类型不同(如 Order → 平均价格) | 状态只存累加器,最省空间 |
BroadcastState<K,V> | 广播流的 KV | 规则/维表广播 | 只能通过 BroadcastProcessFunction 访问;广播侧只读、处理侧只写 |
最容易被追问的三个点:
- 「去重该用哪个?」 →
MapState(或ValueState<RoaringBitmap>做基数十亿级的去重)。不要用ListState:它既不能按值删除,又需要遍历判重,时间是 O(n); - 「为什么
AggregatingState比ListState省?」 → 前者状态里只有一个累加器对象,后者是全量元素;前者add是 O(1) 合并,后者处理时是 O(n) 遍历; - 「
MapState的底层实现?」 → 在HashMapStateBackend下就是一个HashMap(堆内存里的 Java 对象);在RocksDBStateBackend下是按 key 存进 RocksDB 的序列化 KV,每个 key 都单独读写磁盘(所以 key 数量多时 RocksDB 会有明显读放大)。
3. StateBackend 有哪两种?怎么选?
答: Flink 1.13 起统一为两种(旧的 MemoryStateBackend / FsStateBackend 已整合):
HashMapStateBackend | EmbeddedRocksDBStateBackend | |
|---|---|---|
| 旧名 | MemoryStateBackend / FsStateBackend | RocksDBStateBackend |
| 状态存哪 | TaskManager 的 JVM 堆内存 | TaskManager 本地磁盘上的 RocksDB |
| 状态大小上限 | 受堆内存限制(大状态必 OOM) | 受本地磁盘限制(可到 TB 级) |
| 读写性能 | 快(内存对象,无序列化开销) | 相对慢(每次读写都要序列化/反序列化 + 可能落盘) |
| Checkpoint | 全量快照(做全量序列化) | 可增量 Checkpoint(只上传新增的 SST 文件) |
| 是否占堆 | 是(大状态会频繁 Full GC) | 否(堆占用小、GC 友好) |
| 适用 | 小状态(< 几 GB)、追求低延迟 | 大状态、需要增量 Checkpoint |
选型逻辑(面试标准答案):
RocksDB 的关键机制(高频追问):
- 状态在 RocksDB 里是「序列化的 KV」:所以同一份状态会同时存在于「内存中的状态对象」和「RocksDB 中的字节」两个层次——这就是它「省内存但慢」的原因;
- 增量 Checkpoint:RocksDB 的 SST 文件不可变,每次 Checkpoint 只上传新增/变更的 SST 文件与 Manifest,因此第一次是全量,之后是增量;
- 本地恢复(Local Recovery):Checkpoint 时把状态文件同时留一份在 TaskManager 本地盘(
state.backend.local-recovery: true),恢复时先从本地读、只从远端补差异,大状态恢复速度提升明显; - TTL 与 Compaction:RocksDB 的 compaction 会带来后台 I/O,和状态读写争抢磁盘,在「大状态 + 高吞吐」时要关注磁盘 IOPS;
state.backend.rocksdb.memory.managed: true(默认):RocksDB 的内存从 Flink Managed Memory 里分配,而不是 JVM 堆,避免影响 GC。
4. State TTL 是什么?怎么用?
答: 状态是「永久累积」的,如果不清理,key 基数大会导致状态无限增长——State TTL 就是给状态加"过期时间":
StateTtlConfig ttl = StateTtlConfig
.newBuilder(Duration.ofHours(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) // 何时刷新"过期时间"
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) // 过期后是否还能读到
.cleanupInRocksdbCompactFilter(1000) // RocksDB 压缩时顺带清理
.build();
ValueStateDescriptor<String> desc = new ValueStateDescriptor<>("s", String.class);
desc.enableTimeToLive(ttl);三个配置点必须能讲清楚:
| 配置 | 可选值 | 含义 |
|---|---|---|
UpdateType | Disabled / OnCreateAndWrite(默认)/ OnReadAndWrite | 「多久没被写就过期」还是「多久没被读写就过期」 |
StateVisibility | NeverReturnExpired(默认)/ ReturnExpiredIfNotCleanedUp | 过期但还没被物理清理的状态,读的时候是否返回 |
cleanup(清理策略) | 全量快照清理 / 增量清理(heap) / RocksDB Compaction Filter 清理 / 后台清理 | 决定"过期"何时真正释放内存 |
TTL 的四个坑(面试加分):
- 过期 ≠ 立即删除:TTL 是惰性/异步的,过期数据可能还占着内存直到被清理——所以不能靠 TTL 精确控制内存;
OnCreateAndWrite下「只读不写」的状态不会被续期:如果业务是「频繁读、偶尔写」,要用OnReadAndWrite,否则会发生「读着读着状态突然没了」;- 清理会带来开销:
cleanupIncrementally和cleanupInRocksdbCompactFilter都有 CPU 成本,配置过激进会拖慢吞吐; - TTL 不能跨 StateBackend 迁移语义:开了 TTL 的状态,其二进制格式与未开不同,从 Savepoint 恢复时改 TTL 配置要谨慎(可能不兼容)。
一句话:TTL 是"防状态爆炸"的第一道防线,但它是最终一致的内存回收,不是精确过期。
5. 状态的「有界性」为什么重要?如何避免状态无限增长?
答: 状态无界 = 迟早 OOM / Checkpoint 超时 / 恢复变慢。这是 Flink 生产事故的头号原因。
状态无限增长的五个典型来源:
| 来源 | 例子 | 对策 |
|---|---|---|
| key 基数爆炸 | 用「请求 ID / TraceID」当 key 做聚合 | 换业务维度 key;或 TTL + 二级聚合 |
MapState/ListState 只加不删 | 去重集合、join 缓存 | TTL、MapState 上限 + LRU 淘汰、定期清理 |
窗口 + allowedLateness 过大 | 大窗口 + 长迟到容忍 | 缩小窗口、缩短 allowedLateness |
| 双流 join 缓存不清理 | 一条流长期不来数据,另一条流缓存一直涨 | 用 intervalJoin(自动清理)或加 TTL/定时器 清理 |
ProcessWindowFunction 缓存全量元素 | 大窗口 + 全量 Iterable | 改用 增量聚合(aggregate) |
排查手段:
- Web UI → Checkpoint 详情:看
State Size与各算子的Checkpointed Data Size,找到最大的那个算子; - 作业指标:
numBytesUsed、numKeys、lastCheckpointSize趋势; - 报错信号:
Checkpoint expired before completing、Metaspace OOM、RocksDB native memory增长。
二、Checkpoint 与容错
6. Checkpoint 与 Savepoint 有什么区别?
答:
| 维度 | Checkpoint | Savepoint |
|---|---|---|
| 谁触发 | Flink 自动(按 execution.checkpointing.interval) | 用户手动(bin/flink savepoint <jobId> 或 REST API) |
| 目的 | 容错(故障后自动恢复) | 运维(升级/改并行度/迁移/暂停) |
| 生命周期 | 通常保留最近若干个(RETAIN_ON_CANCELLATION 可控制取消时是否保留) | 永久保留,直到用户删除 |
| 格式 | 与 Checkpoint 同一套,但 Savepoint 是"标准格式"(跨版本/跨后端可迁移) | 同左 |
| 触发对性能影响 | 需控制间隔,不能太频繁 | 通常主动触发,允许慢一点 |
| 是否可改并行度 | 一般不直接改(需先 Savepoint 或开自适应调度) | 支持(按 key 重分布) |
关键结论:
- Savepoint 是 Checkpoint 的一种"更标准的形态",可以用来跨 Flink 版本升级作业(但要考虑状态 schema 兼容、TTL 配置、UDF 序列化格式);
- 调整并行度:对
Keyed State可以按 key 重分布,对Operator State需要自己实现redistribute(EvenSplit/Union); - 取消作业时:默认 Checkpoint 会被删除(除非配置了
externalized),而 Savepoint 不会被自动删除。
7. Checkpoint 的整体流程是什么?(Barrier 与异步快照)
答: Flink 的 Checkpoint 基于 Chandy-Lamport 异步屏障快照算法,核心是「Barrier 随数据流走,算子收到就快照」。
「异步」体现在哪里(高频追问):
- 快照与处理解耦:算子在处理 Barrier 的同时,把状态复制/持久化到后台(
HashMapStateBackend会把状态内存拷贝一份,RocksDB 直接利用 SST 文件快照),不必停流等快照写完; - 上传与处理解耦:状态快照的上传到远端存储是异步的;
- Barrier 对齐(Aligned Checkpoint)是唯一的同步点:多个输入通道的 Barrier 到达时间不一致时,先到的通道会被阻塞,直到所有通道的 Barrier 都到齐(见第 8 题)。
两阶段:
- 同步阶段:等所有 Barrier 对齐,并读取/冻结当前状态(很快);
- 异步阶段:把状态写入状态后端并上传远端(较慢,但不阻塞数据处理)。
8. 什么是 Barrier 对齐?非对齐 Checkpoint 解决什么问题?
答: 这是 Checkpoint 最核心的机制,也是反压场景下最容易出问题的地方。
对齐(Aligned)Checkpoint 的问题:
非对齐(Unaligned)Checkpoint 的思路:别等了,先把「还在路上的数据」也一起存下来。
| Aligned(默认) | Unaligned | |
|---|---|---|
| 等待 | 等所有通道 Barrier 对齐 | 不等 |
| 快照内容 | 只含算子状态 | 算子状态 + 在途数据(Input Buffer 中的临时数据) |
| Checkpoint 耗时 | 在反压时显著变长 | 基本不受反压影响 |
| 快照体积 | 小 | 大(多存了在途数据) |
| 恢复速度 | 一般 | 可能更慢 / 也可能更简单(恢复时先回放快照中的在途数据) |
| 适用 | 无反压或轻微反压 | 严重反压、大状态、Checkpoint 总是超时 |
开启方式:
execution.checkpointing.unaligned: true
execution.checkpointing.aligned-checkpoint-timeout: 30s # 对齐超时后自动转非对齐(1.14+)
execution.checkpointing.unaligned.max-subtasks-per-channel-state-file: 5
execution.checkpointing.unaligned.forced: false结论:非对齐 Checkpoint = 「用空间换时间」——用更大的快照(多存了在途数据)换取 Checkpoint 的及时完成。当作业长期高负载/反压且 Checkpoint 频繁超时,它是比「调大超时时间」更根本的解法。
9. Checkpoint 的存储位置和保留策略怎么配?
答:
state.backend.type: hashmap | rocksdb # Flink 1.15+ 用 state.backend.type
state.checkpoints.dir: hdfs://nn/flink/checkpoints # 所有 Checkpoint/Savepoint 的目录
state.savepoints.dir: hdfs://nn/flink/savepoints # Savepoint 默认目录
execution.checkpointing.interval: 60s # 触发间隔
execution.checkpointing.timeout: 10min # 单次超时(超时则本次作废)
execution.checkpointing.min-pause: 30s # 两次 Checkpoint 之间的最小间隔(防连击)
execution.checkpointing.max-concurrent-checkpoints: 1 # 并发 Checkpoint 数(默认 1)
execution.checkpointing.tolerable-failed-checkpoints: 3 # 连续失败 N 次才让作业失败
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION要点:
| 配置 | 为什么重要 |
|---|---|
interval vs timeout | interval 是触发频率,timeout 是单次完成时限;timeout 应远大于正常完成时间 |
min-pause | 防止「上一次还没完、下一次就触发」把资源打满;建议 ≥ 状态写入耗时 |
max-concurrent-checkpoints | 默认 1(串行)最安全;调大能提高容错新鲜度,但会与数据处理争抢 I/O |
tolerable-failed-checkpoints | 容忍几次失败,避免偶发抖动直接让作业挂掉 |
externalized | 决定「作业取消后是否保留 Checkpoint」,RETAIN_ON_CANCELLATION 是常见运维选择 |
unaligned-checkpoint | 反压场景的解法(见第 8 题) |
Checkpoint 里除了状态还存了什么(容易被忽略):
- Kafka offset / Source 的消费位点(在 Source 算子的 Operator State 里);
- Sink 的事务信息(两阶段提交的待提交事务);
- 算子/作业的元数据(并行度、Chain 结构、UDF 序列化信息)——所以并行度或拓扑变化时无法直接用 Checkpoint 恢复,要用 Savepoint + 显式的允许变更配置。
10. 端到端 Exactly-Once 如何实现?(三段式)
答: 端到端 EO = Source 可重放 + 状态一致性 + Sink 事务化,缺一不可:
① Source 端:消费位点(offset)必须作为状态被 Checkpoint 一起保存,故障恢复时从 Checkpoint 中的 offset 重新消费(所以必须「可重放」,如 Kafka)。
⚠️ 常见误解:不能依赖 Kafka 的
auto.offset.commit,因为那是异步提交的,和 Checkpoint 不在同一个时间点,会导致「重复消费」或「丢数据」。Flink 的KafkaSource只在 Checkpoint 成功后才提交 offset。
② Flink 内部:如第 7 题,Barrier 快照保证「状态 = 截止到某条数据的精确视图」。
③ Sink 端:两阶段提交(Two-Phase Commit)——这是端到端 EO 的关键:
关键实现类与 API:
- 老 API:
TwoPhaseCommitSinkFunction(1.15 起已废弃); - 新 API(FLIP-147,1.15+):
Sink的SinkWriter/CommittingSinkWriter,如KafkaSink的DeliveryGuarantee.EXACTLY_ONCE。
三个"EO 失效"的陷阱(面试加分):
- Sink 有外部副作用(如调用第三方 HTTP API、发消息到不支持事务的系统)→ 只能做到 At-Least-Once + 下游幂等;
- Checkpoint 成功后才提交事务:如果「提交事务」这一步失败,会回滚重来,因此下游看到的数据是「幂等」的(可能重复提交,靠事务 ID 去重);
- 事务 ID 必须唯一且稳定:Kafka 事务 ID 的生成要包含「作业 ID + 子任务 ID + Checkpoint ID」等,且重启后不能与历史冲突,否则会报
ProducerFencedException(常见的transactional.id冲突问题)。
11. 作业失败了怎么恢复?重启策略有哪几种?
答:
| 重启策略 | 触发条件 | 说明 | 适用 |
|---|---|---|---|
fixed-delay(默认) | 失败后固定间隔重启 | restart-strategy.fixed-delay.attempts(次数)、delay(间隔) | 通用 |
failure-rate | 按时间窗内的失败率判断 | failure-rate-interval 内失败数 ≤ max-failures-per-interval 则重启 | 最常用(避免"抖动导致无限重启") |
exponential-delay | 指数退避重启,间隔逐次翻倍 | 1.19+;有 initial-backoff、max-backoff、jitter | 依赖外部服务恢复较慢时 |
no-restart | 不重启,直接失败 | — | 调试 / 一次性批作业 |
legacy | 无 Checkpoint 时退化为 no-restart | 兼容旧版 | 不推荐 |
cluster-level / job-level | 重启范围 | 前者由 YARN/K8s 重启整个集群(失败会丢状态),后者 JobMaster 原地重启作业(保留 Checkpoint) | 生产用 job-level |
恢复流程:
max-concurrent-checkpoints 与「重启风暴」:
- 重启风暴:某外部依赖(如 MySQL)抖动 → 大量算子失败 → 反复重启 → 重启后又立即超时……缓解手段:
failure-rate策略(限制重启频率)、tolerable-failed-checkpoints、把「外部依赖调用」加熔断/退避、Async I/O 隔离阻塞。 - Region Failover:把 ExecutionGraph 切成 Region(通过 Shuffle 边界),只重启出问题的 Region,其他 Region 继续跑,显著降低恢复成本。
12. Checkpoint 常见失败原因与排查思路?
答: 这是实战必备,面试时能列出来会非常有说服力:
| 现象 / 报错 | 根因 | 对策 |
|---|---|---|
Checkpoint expired before completing | 反压导致 Checkpoint 超时 | ①开启非对齐 Checkpoint;②排查反压源头(见《Flink(四)》);③降低 interval 频率别太密;④max-concurrent-checkpoints 保持 1 |
Checkpoint was declined / Exceeded checkpoint tolerable failure threshold | 连续失败超阈值 | 先解决根因,临时调大 tolerable-failed-checkpoints |
Checkpoint 一直 IN_PROGRESS,State Size 巨大 | 状态膨胀 / ProcessWindowFunction 缓存全量 / key 基数爆炸 | 加 TTL、改增量聚合、缩窗口;查 Web UI 的 State Size |
| Barrier 对齐时间过长 | 数据倾斜 / 某个子任务慢 | 解决倾斜(加盐/两阶段聚合);或开非对齐 Checkpoint |
| 写到远端存储慢 / 超时 | HDFS/S3 带宽瓶颈 | 开启增量 Checkpoint(RocksDB)、调整 HDFS 并发、增大 timeout |
Metaspace OOM | 每个算子/每次 Checkpoint 生成类加载器(动态代码生成) | 调大 Metaspace、限制 classloader.check-leaked-classloader |
| 恢复时反序列化失败 | UDF / POJO / 状态 schema 变化 | 保证兼容(加字段用默认值、不改字段类型);跨版本走 Savepoint + 版本兼容配置 |
ProducerFencedException(Kafka EO Sink) | 事务 ID 冲突 / 有另一个实例在写 | 保证 transactional.id 唯一、不要重复提交同一作业 |
Not all required tasks are currently deployed | 重启期间资源不足 | 检查 Slot/资源池、slotSharingGroup 配置 |
排查顺序(SOP):
13. Flink 的状态一致性语义有哪几级?怎么选?
答:
| 语义 | 含义 | 如何实现 | 代价 |
|---|---|---|---|
| At-Most-Once | 每条数据最多处理一次(可能丢) | 失败不重放;Checkpoint 关掉 | 最低延迟、最低开销 |
| At-Least-Once | 每条数据至少一次(可能重复) | 有 Checkpoint、失败重放 | 需要下游幂等 |
| Exactly-Once | 状态层面精确一次(每条数据对 Flink 内部的影响只有一次) | Checkpoint + Source 位点 + Barrier 快照 | Checkpoint 开销、Sink 依赖事务 |
必须说的两个澄清(面试拉分点):
- Flink 的 Exactly-Once 是「状态层面」的 EO:指「故障恢复后状态被精确重建」,不等于「下游只收到一次数据」。端到端 EO 还要 Sink 配合(见第 10 题);
- 语义是「分级选择」的:
env.enableCheckpointing()的成本并不便宜,对延迟极敏感且可接受重复的场景,At-Least-Once + 幂等下游往往更划算; - Checkpoint 关闭时的一致性:不 Checkpoint 就没有容错,重启即从头开始(真实业务少见),所以生产至少要 At-Least-Once。
14. 如何做到「无停机」的作业升级与扩缩容?
答:
标准流程(有状态作业升级):
① 停止消费 / 或直接停作业(带 Savepoint)
bin/flink stop --savepointPath hdfs://.../savepoints <jobId>
② 修改代码/并行度
③ 从 Savepoint 恢复启动
bin/flink run -s hdfs://.../savepoints/savepoint-xxx -d ...关键注意点:
| 事项 | 说明 |
|---|---|
| 不能直接改并行度用 Checkpoint | Checkpoint 元数据里绑定了并行度与拓扑,必须走 Savepoint(--allowNonRestoredState 可容忍删除的算子) |
| 状态 schema 兼容 | 给 POJO 加字段可以(旧状态该字段为 null/默认值),改字段类型/删除字段要谨慎;Kryo 序列化的类几乎不兼容 |
| Keyed State 并行度变更 | 按 key 自动重分布(hash 一致,恢复正确) |
| Operator State 并行度变更 | 需要实现 CheckpointedFunction 的 redistribute:EvenSplit(均分)或 Union(每子任务全量)——写错会导致 offset 丢失/重复 |
| TTL 配置变更 | 会影响状态序列化格式,尽量别在恢复时改 |
| Source/Sink 位点 | Kafka offset 在 Source 的 Operator State 里,并行度变更时要确保重分布正确 |
| UDF 序列化 | UDF 会随 Savepoint 保存(uid 标识),算子的 uid 必须显式设置且不轻易改,否则恢复时找不到对应状态 |
最佳实践:给所有算子显式设置
uid(singleOutputStreamOperator.uid("order-window")),这样重构代码(改类名/改顺序)后依然能从 Savepoint 恢复;这是很多团队吃了亏才加上的约定。
15. 综合实战:一个 Checkpoint 一直超时的作业,你怎么排查和优化?
答: 完整回答(能力综合题):
第一步:定位——判断「慢在哪」:
第二步:对症优化
| 症状 | 优化 |
|---|---|
| 反压 | ① 定位反压源头(从 Source 往下找第一个 HIGH 的算子);② 提升下游并行度/资源;③ 开启非对齐 Checkpoint;④ 消除数据倾斜 |
| 大状态(Sync Duration 高) | ① 换 RocksDB Backend(堆外 + 增量快照);② 开启增量 Checkpoint;③ 开启本地恢复;④ 减小状态:加 TTL、改增量聚合、缩窗口、控制 key 基数 |
| 上传慢(Async 高) | ① 提高 Checkpoint 存储带宽/并发;② 换更近的存储(同机房 HDFS/OSS);③ 开启 state.backend.incremental: true |
| 对齐慢(Alignment 高) | ① 非对齐 Checkpoint;② 解决倾斜;③ 调整 execution.checkpointing.aligned-checkpoint-timeout 自动降级 |
| 触发太密 | 增大 interval、设置 min-pause,避免 Checkpoint 互相抢资源 |
第三步:配置层面的"稳"
# 推荐起手配置(大状态 + 有反压风险)
state.backend.type: rocksdb
state.backend.incremental: true
state.backend.local-recovery: true
state.checkpoints.dir: hdfs://nn/flink/checkpoints
execution.checkpointing.interval: 3min
execution.checkpointing.timeout: 10min
execution.checkpointing.min-pause: 1min
execution.checkpointing.max-concurrent-checkpoints: 1
execution.checkpointing.unaligned: true
execution.checkpointing.tolerable-failed-checkpoints: 2
execution.checkpointing.externalized-checkpoint-retention: RETAIN_ON_CANCELLATION一句话总结:Checkpoint 超时的根因通常不在 Checkpoint 本身,而在上游的「反压」或状态本身的「膨胀」——先解决数据流的稳态问题(下一篇),再谈 Checkpoint 参数调优。
下一篇:《Flink(四)》聚焦反压、内存模型与性能调优——Credit-based 流控的完整链路、Web UI 反压指标怎么读、反压源头的定位方法、TaskManager 内存模型拆分、数据倾斜与两阶段聚合、大状态调优与常见生产事故。
