Flink调整并行度时是否会丢失旧的Keyed State?
Flink调整作业并行度是否会丢失Keyed State
正常符合操作规范的前提下,调整并行度不会丢失已有的Keyed State,不管是并行度调大还是调小的场景都适用。
背后的核心逻辑
Flink引入了Key Group作为Keyed State分配的最小单元:
- 每个Key的哈希值会映射到对应的Key Group上,Key Group的总数量等于作业首次启动时设置的
maxParallelism(最大并行度,作业启动后不可修改) - 运行时每个并行实例负责处理固定数量的Key Group,对应的状态也会绑定到对应实例的Key Group下
不同调整场景的状态表现
- 并行度调大(如从2调整为5):原有实例持有的Key Group会按规则拆分,分配给新增的实例,所有历史状态都会跟着对应的Key Group迁移,不会丢失
- 并行度调小(如从5调整为3):原有多个实例持有的Key Group会合并到更少的实例中,所有状态都会完整保留
可能出现状态丢失的异常场景
只有触发以下错误操作时才会出现状态丢失:
- 调整后的并行度超过了作业设置的
maxParallelism,导致部分Key Group无法找到对应的实例映射 - 调整并行度时没有从作业原有的checkpoint/savepoint恢复,直接冷启动作业
- 调整并行度的同时修改了键的序列化逻辑、状态名称、状态数据结构,导致原有状态无法被正常读取
标准调整流程:先对运行中的作业触发
savepoint,停止作业后指定新的并行度从该savepoint启动即可,只要满足调整后并行度 ≤maxParallelism的前提,就能保证状态完全恢复。
内容的提问来源于stack exchange,提问作者yanghaogn
相关产品推荐
相关产品推荐

