如何通过Java代码上传Flink任务Jar包至集群提交作业
Java代码实现无
createRemoteEnvironment的Flink Jar包作业提交 要实现和Flink Web UI完全一致的作业提交效果,不需要调用StreamExecutionEnvironment.createRemoteEnvironment()方法,本质是直接模拟UI调用Flink JobManager暴露的官方REST接口完成全流程,和手动在UI点提交的行为完全对齐。
核心流程
Flink Web UI提交作业分为两个固定步骤,代码实现完全复刻该流程即可:
- 第一步:将本地作业Jar包以multipart表单形式上传至Flink集群,获取集群侧生成的Jar唯一标识(jarid)
- 第二步:携带jarid、作业入口类、并行度、启动参数等配置,调用作业运行接口触发作业启动
前置依赖
引入HTTP请求、JSON处理的基础依赖即可,不需要额外引入Flink官方的客户端依赖包:
<dependencies> <!-- HTTP客户端,用于调用Flink REST接口 --> <dependency> <groupId>org.apache.httpcomponents.client5</groupId> <artifactId>httpclient5</artifactId> <version>5.2.1</version> </dependency> <!-- JSON序列化与结果解析 --> <dependency> <groupId>com.fasterxml.jackson.core</groupId> <artifactId>jackson-databind</artifactId> <version>2.15.2</version> </dependency> </dependencies>
完整实现代码
import com.fasterxml.jackson.databind.JsonNode; import com.fasterxml.jackson.databind.ObjectMapper; import org.apache.hc.client5.http.classic.methods.HttpPost; import org.apache.hc.client5.http.entity.mime.FileBody; import org.apache.hc.client5.http.entity.mime.MultipartEntityBuilder; import org.apache.hc.client5.http.impl.classic.CloseableHttpClient; import org.apache.hc.client5.http.impl.classic.HttpClients; import org.apache.hc.core5.http.ContentType; import org.apache.hc.core5.http.io.entity.EntityUtils; import org.apache.hc.core5.http.io.entity.StringEntity; import java.io.File; import java.util.HashMap; import java.util.Map; public class FlinkJarSubmitter { private static final ObjectMapper OBJECT_MAPPER = new ObjectMapper(); // Flink JobManager REST地址,和Web UI访问地址一致 private static final String FLINK_REST_BASE_URL = "http://your-flink-jobmanager:8081"; public static void main(String[] args) throws Exception { // 本地待上传的Flink作业Jar包路径 String localJarPath = "/path/to/your/flink-job.jar"; // 作业入口类全限定名 String entryClass = "com.your.package.FlinkJobMain"; // 作业并行度 int parallelism = 3; // 作业启动参数 String programArgs = "--source-topic test --sink-path hdfs:///output"; // 1. 上传Jar包获取jarid String jarId = uploadJar(localJarPath); System.out.println("Jar上传成功,jarid为:" + jarId); // 2. 提交作业 String jobId = submitJob(jarId, entryClass, parallelism, programArgs); System.out.println("作业提交成功,jobid为:" + jobId); } /** * 上传Jar包到Flink集群 */ private static String uploadJar(String localJarPath) throws Exception { String uploadUrl = FLINK_REST_BASE_URL + "/v1/jars/upload"; try (CloseableHttpClient httpClient = HttpClients.createDefault()) { HttpPost uploadRequest = new HttpPost(uploadUrl); // 构造multipart表单,和UI上传的表单格式完全一致 MultipartEntityBuilder entityBuilder = MultipartEntityBuilder.create(); entityBuilder.addPart("jarfile", new FileBody(new File(localJarPath))); uploadRequest.setEntity(entityBuilder.build()); return httpClient.execute(uploadRequest, response -> { String respStr = EntityUtils.toString(response.getEntity()); if (response.getCode() != 200) { throw new RuntimeException("Jar上传失败,响应信息:" + respStr); } JsonNode respNode = OBJECT_MAPPER.readTree(respStr); if (!"success".equals(respNode.get("status").asText())) { throw new RuntimeException("Jar上传失败,响应信息:" + respStr); } // 从返回的文件路径中截取jarid String filePath = respNode.get("filename").asText(); return filePath.substring(filePath.lastIndexOf("/") + 1); }); } } /** * 触发作业运行 */ private static String submitJob(String jarId, String entryClass, int parallelism, String programArgs) throws Exception { String submitUrl = FLINK_REST_BASE_URL + "/v1/jars/" + jarId + "/run"; try (CloseableHttpClient httpClient = HttpClients.createDefault()) { HttpPost submitRequest = new HttpPost(submitUrl); // 构造提交参数,和UI提交传参完全一致 Map<String, Object> submitParams = new HashMap<>(); submitParams.put("entryClass", entryClass); submitParams.put("parallelism", parallelism); submitParams.put("programArgs", programArgs); submitParams.put("allowNonRestoredState", false); // 如需从savepoint启动可追加savepointPath参数 // submitParams.put("savepointPath", "hdfs:///savepoints/savepoint-xxx"); StringEntity jsonEntity = new StringEntity( OBJECT_MAPPER.writeValueAsString(submitParams), ContentType.APPLICATION_JSON ); submitRequest.setEntity(jsonEntity); return httpClient.execute(submitRequest, response -> { String respStr = EntityUtils.toString(response.getEntity()); if (response.getCode() != 200) { throw new RuntimeException("作业提交失败,响应信息:" + respStr); } JsonNode respNode = OBJECT_MAPPER.readTree(respStr); return respNode.get("jobid").asText(); }); } } }
注意事项
- 运行代码的环境需要能正常访问Flink JobManager的REST端口(默认8081),如果集群开启了Kerberos、基础认证等安全校验,需要在HTTP请求头中追加对应的认证信息
- 上传Jar包大小不能超过Flink配置中
web.upload.max-size设定的阈值,否则会被集群拦截 - 如果作业有额外第三方依赖,建议直接打在作业Fat Jar中,也可以在提交参数的
classpaths字段中指定集群节点上已存在的依赖路径 - 提交成功返回的jobid和Flink Web UI上展示的作业ID完全一致,可以直接用这个ID在UI上查询作业状态、管理作业生命周期
内容的提问来源于stack exchange,提问作者Maria Isla
相关产品推荐
相关产品推荐

