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

如何通过Java代码上传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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.30 21:45:38