Flink 1.18 Application模式多Job的Savepoint恢复与定时触发问题
在Flink 1.18 Application模式单Pod多Job场景下的Savepoint管理方案
针对你在K8s Docker Pod中以Application模式运行多Job的场景,以下是关于定时触发Savepoint和重启恢复的具体实现方案:
一、每小时为每个Job触发Savepoint
方案1:通过Flink REST API + 定时任务(推荐)
利用Flink的REST接口实现外部定时触发,不受Job进程状态影响,可靠性更高:
- 获取运行中JobID:调用Flink REST的
GET /jobs接口,提取所有Job的ID; - 触发Savepoint:对每个JobID调用
POST /jobs/{jobId}/savepoints接口,指定共享存储路径(如S3、NFS PV等); - 定时执行:在Pod内通过cron或K8s CronJob定时执行脚本。
示例Shell脚本(需确保镜像包含curl和jq):
#!/bin/bash FLINK_REST_URL="http://localhost:8081" SAVEPOINT_ROOT="s3://your-shared-bucket/flink-savepoints" SHARED_META_PATH="/shared-meta" # 获取所有运行中的JobID JOB_IDS=$(curl -s $FLINK_REST_URL/jobs | jq -r '.jobs[] | select(.status=="RUNNING") | .id') for JOB_ID in $JOB_IDS; do # 获取Job名称 JOB_NAME=$(curl -s $FLINK_REST_URL/jobs/$JOB_ID | jq -r '.name') # 触发Savepoint RESPONSE=$(curl -X POST $FLINK_REST_URL/jobs/$JOB_ID/savepoints \ -H "Content-Type: application/json" \ -d "{\"targetDirectory\": \"$SAVEPOINT_ROOT/$JOB_NAME\", \"cancelJob\": false}") # 提取并记录最新Savepoint路径 SAVEPOINT_PATH=$(echo $RESPONSE | jq -r '.operationLocation' | xargs curl -s | jq -r '.savepointPath') echo $SAVEPOINT_PATH > $SHARED_META_PATH/$JOB_NAME-latest.txt done
将脚本加入Pod的cron配置(如0 * * * * /path/to/trigger-savepoints.sh),或用K8s CronJob定期调用Pod的REST接口。
方案2:Java代码内定时触发
在Job进程内部通过调度器定时触发,适合不需要外部依赖的场景:
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.TimeUnit; public class MultiJobApp { private static final String SAVEPOINT_BASE = "s3://your-shared-bucket/flink-savepoints"; private static final ScheduledExecutorService scheduler = Executors.newScheduledThreadPool(2); public static void main(String[] args) throws Exception { // 启动Job1并定时触发Savepoint StreamExecutionEnvironment env1 = StreamExecutionEnvironment.getExecutionEnvironment(); // 构建Job1业务逻辑... var job1 = env1.executeAsync("UserBehaviorJob"); scheduler.scheduleAtFixedRate(() -> { try { var savepointPath = env1.executeSavepointAsync(SAVEPOINT_BASE + "/UserBehaviorJob", job1.getJobID()).get(); // 记录Savepoint路径到共享存储(如ConfigMap、数据库) recordLatestSavepoint("UserBehaviorJob", savepointPath); } catch (Exception e) { e.printStackTrace(); } }, 1, 1, TimeUnit.HOURS); // 启动Job2并定时触发Savepoint StreamExecutionEnvironment env2 = StreamExecutionEnvironment.getExecutionEnvironment(); // 构建Job2业务逻辑... var job2 = env2.executeAsync("OrderProcessingJob"); scheduler.scheduleAtFixedRate(() -> { try { var savepointPath = env2.executeSavepointAsync(SAVEPOINT_BASE + "/OrderProcessingJob", job2.getJobID()).get(); recordLatestSavepoint("OrderProcessingJob", savepointPath); } catch (Exception e) { e.printStackTrace(); } }, 1, 1, TimeUnit.HOURS); Thread.currentThread().join(); } private static void recordLatestSavepoint(String jobName, String path) { // 实现将路径写入共享存储的逻辑 } }
二、Pod重启时从Savepoint恢复Job
核心逻辑是启动Job前读取对应Job的最新Savepoint路径,指定恢复配置:
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.environment.SavepointRestoreSettings; public class MultiJobApp { private static final String SAVEPOINT_BASE = "s3://your-shared-bucket/flink-savepoints"; public static void main(String[] args) throws Exception { // 恢复Job1 StreamExecutionEnvironment env1 = StreamExecutionEnvironment.getExecutionEnvironment(); env1.setRestartStrategy(RestartStrategies.noRestart()); // 依赖Savepoint恢复,关闭内部重启策略 String job1Savepoint = getLatestSavepoint("UserBehaviorJob"); if (job1Savepoint != null) { env1.setSavepointRestoreSettings(SavepointRestoreSettings.forPath(job1Savepoint)); } // 构建Job1业务逻辑... env1.executeAsync("UserBehaviorJob"); // 恢复Job2 StreamExecutionEnvironment env2 = StreamExecutionEnvironment.getExecutionEnvironment(); env2.setRestartStrategy(RestartStrategies.noRestart()); String job2Savepoint = getLatestSavepoint("OrderProcessingJob"); if (job2Savepoint != null) { env2.setSavepointRestoreSettings(SavepointRestoreSettings.forPath(job2Savepoint)); } // 构建Job2业务逻辑... env2.executeAsync("OrderProcessingJob"); Thread.currentThread().join(); } private static String getLatestSavepoint(String jobName) { // 从共享存储读取对应Job的最新Savepoint路径 // 示例:读取共享文件路径 try { return new String(Files.readAllBytes(Paths.get("/shared-meta/" + jobName + "-latest.txt"))).trim(); } catch (Exception e) { return null; } } }
关键注意事项
- Savepoint必须存储在Pod可访问的共享存储中,禁止使用本地路径;
- 每个Job需设置唯一名称,用于关联对应的Savepoint;
- 定时触发Savepoint时需避免并发执行,可通过检查Job状态或分布式锁实现;
- 定期清理旧Savepoint,避免存储成本过高。
内容的提问来源于stack exchange,提问作者zhen ma
相关产品推荐
相关产品推荐

