You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Flink作业如何为特定状态禁用Savepoint快照

Flink单作业区分状态Savepoint持久化的实现方式

完全可以实现,根据你使用的Flink版本选择对应方案即可:

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.28 14:57:18