Flink Operator部署应用失败:指定Savepoint不存在问题求助
Flink Operator部署升级时Savepoint不存在导致启动失败的问题
环境与配置
使用Flink Operator 1.6.1部署Flink 1.14.5应用,相关配置如下:
Job配置
job: allowNonRestoredState: true args: - <command to run my job> entryClass: org.apache.beam.runners.flink.FlinkPortableClientEntryPoint jarURI: local:///opt/flink/flink-web-upload/beam-runner.jar parallelism: 32 state: running upgradeMode: last-state
Flink配置
flinkConfiguration: fs.s3a.aws.credentials.provider: com.amazonaws.auth.WebIdentityTokenCredentialsProvider high-availability: org.apache.flink.kubernetes.highavailability.KubernetesHaServicesFactory high-availability.storageDir: s3://<path>/high-availability state.backend: rocksdb state.backend.fs.checkpointdir: file:///checkpoints/flink/checkpoints state.backend.incremental: "true" state.checkpoints.dir: s3a://<path>/checkpoints state.savepoints.dir: s3a://<path>/savepoints
问题现象
部署仅包含内部库变更的新版本Docker镜像时,部分FlinkDeployment启动失败,报错信息如下:
Caused by: java.io.FileNotFoundException: Cannot find checkpoint or savepoint file/directory 's3a://<path>/savepoints/savepoint-2254eb-61d7c7a6f6a3' on file at org.apache.flink.runtime.state.filesystem.AbstractFsCheckpointStorageAccess.resolveCheckpointPointer(AbstractFsCheckpointStorageAccess.java:275) at org.apache.flink.runtime.state.filesystem.AbstractFsCheckpointStorageAccess.resolveCheckpoint(AbstractFsCheckpointStorageAccess.java:136) at org.apache.flink.runtime.checkpoint.CheckpointCoordinator.restoreSavepoint(CheckpointCoordinator.java:1639) at org.apache.flink.runtime.scheduler.DefaultExecutionGraphFactory.tryRestoreExecutionGraphFromSavepoint(DefaultExecutionGraphFactory.java:163) at org.apache.flink.runtime.scheduler.DefaultExecutionGraphFactory.createAndRestoreExecutionGraph(DefaultExecutionGraphFactory.java:138) at org.apache.flink.runtime.scheduler.SchedulerBase.createAndRestoreExecutionGraph(SchedulerBase.java:335) at org.apache.flink.runtime.scheduler.SchedulerBase.<init>(SchedulerBase.java:191) at org.apache.flink.runtime.scheduler.DefaultScheduler.<init>(DefaultScheduler.java:140) at org.apache.flink.runtime.scheduler.DefaultSchedulerFactory.createInstance(DefaultSchedulerFactory.java:134) at org.apache.flink.runtime.jobmaster.DefaultSlotPoolServiceSchedulerFactory.createScheduler(DefaultSlotPoolServiceSchedulerFactory.java:110) at org.apache.flink.runtime.jobmaster.JobMaster.createScheduler(JobMaster.java:346) at org.apache.flink.runtime.jobmaster.JobMaster.<init>(JobMaster.java:323) at org.apache.flink.runtime.jobmaster.factories.DefaultJobMasterServiceFactory.internalCreateJobMasterService(DefaultJobMasterServiceFactory.java:106) at org.apache.flink.runtime.jobmaster.factories.DefaultJobMasterServiceFactory.lambda$createJobMasterService$0(DefaultJobMasterServiceFactory.java:94) at org.apache.flink.util.function.FunctionUtils.lambda$uncheckedSupplier$4(FunctionUtils.java:112)
疑问
- 代码中未指定该Savepoint,应用为何会获取到这个不存在的路径?
- 能否在找不到Savepoint时强制使用最新Checkpoint恢复?
解答
问题1:自动获取不存在Savepoint路径的原因
你配置了upgradeMode: last-state,Flink Operator会自动从两个来源读取上次作业的恢复点信息:
- Kubernetes中对应FlinkDeployment资源的注解(
flink.apache.org/last-savepoint-*),这里会记录作业上次终止时的恢复点路径,可能是之前手动触发的Savepoint,或是作业失败时自动生成的Savepoint。 - Flink HA存储目录
high-availability.storageDir中的作业元数据,其中也会保存最近的恢复点记录。
如果这个记录的Savepoint因为存储清理、路径变更、权限问题等原因不存在,就会触发找不到文件的错误。
问题2:找不到Savepoint时使用最新Checkpoint恢复的解决方案
可以通过以下几种方式实现:
1. 显式指定从Checkpoint根目录恢复
修改Job配置,指定restoreFromPath为Checkpoint的根目录,同时更新savepointTriggerNonce触发Operator重新识别恢复路径:
job: # 保留原有其他配置 upgradeMode: last-state restoreFromPath: s3a://<path>/checkpoints savepointTriggerNonce: 12345 # 每次修改作业时更换一个随机数,强制Operator刷新恢复逻辑
Flink 1.14.5支持自动从Checkpoint根目录查找最新的有效Checkpoint进行恢复。
2. 清理旧的恢复点元数据
如果Kubernetes注解或HA存储中的旧Savepoint记录已经失效,可以:
- 删除FlinkDeployment资源的旧Savepoint注解:
kubectl annotate flinkdeployment <your-deployment-name> flink.apache.org/last-savepoint-="" --overwrite - 清理HA存储目录中对应作业的元数据文件夹(
s3://<path>/high-availability/<job-id>),让Operator重新尝试从Checkpoint恢复。
3. 临时切换为无状态启动(紧急场景)
如果不需要保留状态,可以临时将upgradeMode改为stateless,先启动作业,之后再手动从Checkpoint恢复状态,但这种方式会丢失未恢复的状态,仅适合紧急恢复场景:
job: # 保留原有其他配置 upgradeMode: stateless
4. 优化Flink配置增强容错
- 确保
state.checkpoints.dir配置正确,且Flink Pod拥有访问该S3路径的权限; - 调整
state.checkpoints.num-retained参数,保留足够多的历史Checkpoint,避免最新Checkpoint被意外清理。
内容的提问来源于stack exchange,提问作者Lydian
相关产品推荐
相关产品推荐

