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

Apache Flink:作业取消后如何通过Java客户端从Checkpoint恢复

核心问题原因

每次通过RemoteStreamEnvironment提交作业时,Flink默认生成随机作业ID,导致新作业无法关联之前保留的Checkpoint,只能从头运行。要实现恢复,需要固定作业ID并指定Checkpoint恢复路径。

具体实现方案

1. 固定作业ID

通过设置固定的作业ID,让Flink识别出这是同一个作业的恢复请求,而非新作业。

2. 指定Checkpoint恢复路径

需要配置从外部化的Checkpoint路径恢复,路径可以是手动指定的具体Checkpoint目录,或者通过Flink Rest API自动获取最新的Checkpoint路径。


代码示例:手动指定恢复路径

import org.apache.flink.api.common.JobID;
import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.JobManagerOptions;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import org.apache.flink.streaming.api.environment.RemoteStreamEnvironment;

public class FlinkJobRecovery {
    public static void main(String[] args) throws Exception {
        // 固定作业ID(可提前生成并记录,比如用JobID.generateRandomId()生成一次后固定)
        String fixedJobId = "your-fixed-job-id-hex-string";
        
        // 构建配置,设置作业ID和恢复路径
        Configuration config = new Configuration();
        config.setString(JobManagerOptions.JOB_ID, fixedJobId);
        // 指定之前保留的Checkpoint路径(比如hdfs:///flink/checkpoints/chk-1234)
        config.setString(CheckpointingOptions.RECOVERY_PATH, "your-external-checkpoint-path");
        
        // 创建RemoteStreamEnvironment
        StreamExecutionEnvironment env = StreamExecutionEnvironment.createRemoteEnvironment(
                "jobmanager-host",
                8081, // JobManager端口
                config,
                "/path/to/your/job.jar" // 作业Jar包路径
        );
        
        // 配置Checkpoint保留策略(和之前提交时一致)
        env.getCheckpointConfig()
           .setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
        
        // 构建作业拓扑(必须和之前提交的拓扑完全一致)
        // ... 你的作业逻辑代码 ...
        
        // 提交作业
        env.execute("Your Job Name");
    }
}

代码示例:自动获取最新Checkpoint路径

如果不想手动维护Checkpoint路径,可以通过Flink Rest API获取指定作业的最新完成Checkpoint,再动态配置恢复路径:

import org.apache.flink.configuration.CheckpointingOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.configuration.JobManagerOptions;
import org.apache.flink.streaming.api.environment.RemoteStreamEnvironment;
import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment;
import com.fasterxml.jackson.databind.JsonNode;
import com.fasterxml.jackson.databind.ObjectMapper;
import java.net.URI;
import java.net.http.HttpClient;
import java.net.http.HttpRequest;
import java.net.http.HttpResponse;

public class FlinkAutoRecovery {
    public static void main(String[] args) throws Exception {
        String jobManagerHost = "jobmanager-host";
        int jobManagerPort = 8081;
        String fixedJobId = "your-fixed-job-id-hex-string";
        String jobJarPath = "/path/to/your/job.jar";
        
        // 调用Flink Rest API获取最新完成的Checkpoint路径
        HttpClient httpClient = HttpClient.newHttpClient();
        String checkpointApiUrl = String.format("http://%s:%d/jobs/%s/checkpoints", jobManagerHost, jobManagerPort, fixedJobId);
        
        HttpRequest request = HttpRequest.newBuilder()
                .uri(URI.create(checkpointApiUrl))
                .GET()
                .build();
        
        HttpResponse<String> response = httpClient.send(request, HttpResponse.BodyHandlers.ofString());
        
        // 解析JSON响应,提取最新Checkpoint的外部路径
        ObjectMapper mapper = new ObjectMapper();
        JsonNode rootNode = mapper.readTree(response.body());
        JsonNode completedCheckpoints = rootNode.get("completed");
        if (completedCheckpoints == null || completedCheckpoints.size() == 0) {
            throw new RuntimeException("No completed external checkpoints found for job: " + fixedJobId);
        }
        // 取最新的第一个Checkpoint(按时间倒序排列)
        JsonNode latestCheckpoint = completedCheckpoints.get(0);
        String latestCheckpointPath = latestCheckpoint.get("externalPath").asText();
        
        // 配置作业ID和恢复路径
        Configuration config = new Configuration();
        config.setString(JobManagerOptions.JOB_ID, fixedJobId);
        config.setString(CheckpointingOptions.RECOVERY_PATH, latestCheckpointPath);
        
        // 创建Remote环境并提交作业
        StreamExecutionEnvironment env = StreamExecutionEnvironment.createRemoteEnvironment(
                jobManagerHost,
                jobManagerPort,
                config,
                jobJarPath
        );
        
        // 保持Checkpoint配置和之前一致
        env.getCheckpointConfig()
           .setExternalizedCheckpointCleanup(CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);
        
        // 构建作业拓扑(必须与原作业一致)
        // ... 你的作业逻辑 ...
        
        env.execute("Auto-Recovery Job");
    }
}

注意事项

  • 作业拓扑必须与原作业完全一致:算子数量、类型、并行度、状态描述等不能修改,否则恢复会失败。
  • 固定作业ID需提前确定:可以第一次提交时生成并记录(比如JobID.generateRandomId().toHexString()),后续提交复用该ID。
  • Checkpoint路径必须是外部化的:确保之前提交作业时设置了RETAIN_ON_CANCELLATION,且路径未被清理。
  • 依赖兼容性:作业Jar包的依赖版本需与Flink集群版本一致,避免兼容性问题。

内容的提问来源于stack exchange,提问作者oceansize

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 21:45:41