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

Flink 1.18 Application模式多Job的Savepoint恢复与定时触发问题

针对你在K8s Docker Pod中以Application模式运行多Job的场景,以下是关于定时触发Savepoint和重启恢复的具体实现方案:

一、每小时为每个Job触发Savepoint

利用Flink的REST接口实现外部定时触发,不受Job进程状态影响,可靠性更高:

  1. 获取运行中JobID:调用Flink REST的GET /jobs接口,提取所有Job的ID;
  2. 触发Savepoint:对每个JobID调用POST /jobs/{jobId}/savepoints接口,指定共享存储路径(如S3、NFS PV等);
  3. 定时执行:在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 23:47:43