Apache Flink:作业取消后如何通过Java客户端从Checkpoint恢复
Flink Java Client 从保留的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
相关产品推荐
相关产品推荐

