Flink作业如何为特定状态禁用Savepoint快照
Flink单作业区分状态Savepoint持久化的实现方式
完全可以实现,根据你使用的Flink版本选择对应方案即可:
推荐方案:使用Flink原生状态持久化级别(Flink 1.17及以上版本支持)
Flink从1.17版本开始为状态新增了持久化级别配置,可以直接给单个状态指定是否要写入Savepoint,不需要拆分作业或者修改状态后端逻辑:
- 对需要随Savepoint保留的核心状态,不需要额外配置,默认的持久化级别就是
CHECKPOINT_AND_SAVEPOINT,会同时写入Checkpoint和Savepoint - 对不需要持久化到Savepoint的非必要状态,在初始化状态描述符时,将其持久化级别设置为
StatePersistenceLevel.ONLY_CHECKPOINT即可,这类状态只会在Checkpoint流程中生成快照,触发Savepoint时会被自动跳过,不会占用Savepoint存储空间,也不会增加Savepoint生成耗时。
代码示例:
// 1. 定义不需要写入Savepoint的临时状态 ValueStateDescriptor<CacheData> tempStateDesc = new ValueStateDescriptor<>( "temp-compute-cache", TypeInformation.of(CacheData.class) ); // 配置该状态仅在Checkpoint时持久化,跳过Savepoint tempStateDesc.setPersistenceLevel(StatePersistenceLevel.ONLY_CHECKPOINT); ValueState<CacheData> tempState = getRuntimeContext().getState(tempStateDesc); // 2. 定义需要随Savepoint保留的核心业务状态,默认级别即可,无需额外配置 ValueStateDescriptor<BizData> coreStateDesc = new ValueStateDescriptor<>( "core-biz-state", TypeInformation.of(BizData.class) ); ValueState<BizData> coreState = getRuntimeContext().getState(coreStateDesc);
该配置对键控状态(Keyed State)、算子状态(Operator State)均生效。后续从Savepoint启动作业时,标记为仅Checkpoint持久化的状态会自动初始化为空值,不会出现状态不匹配的报错,不需要额外编写兼容逻辑。
注意:不要将需要跨版本迁移、长期保留的核心业务状态设置为
ONLY_CHECKPOINT级别。Checkpoint默认绑定作业生命周期,作业取消时如果开启了Checkpoint自动清理,这部分状态会被直接删除,只有写入Savepoint的状态可以用于作业迁移、版本升级等场景。
低版本Flink折中方案(1.17以下版本)
1.17之前的Flink版本没有提供原生的单状态持久化级别配置,不推荐为了这个需求自定义修改StateBackend快照逻辑(改造成本高、容易引入状态一致性问题),如果暂时无法升级Flink版本,可以将非必要状态拆分到独立算子中,通过算子粒度的快照排除规则实现类似效果,但是灵活性远不如原生的状态级别配置。
内容的提问来源于stack exchange,提问作者sparkless
相关产品推荐
相关产品推荐

