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

调用Flink API创建作业时如何添加Java参数?

如果是通过YARN集群模式的Flink API提交作业,核心思路是在Flink配置中添加env.java.opts参数,对应CLI中的-yD参数配置,具体实现如下:

代码示例(Java API)

import org.apache.flink.client.deployment.ClusterClientProvider;
import org.apache.flink.client.program.ClusterClient;
import org.apache.flink.yarn.YarnClusterDescriptor;
import org.apache.flink.yarn.configuration.YarnConfigOptions;
import org.apache.flink.configuration.Configuration;
import org.apache.flink.client.program.PackagedProgram;

import java.io.File;

public class FlinkApiSubmit {
    public static void main(String[] args) throws Exception {
        // 初始化基础Flink配置
        Configuration flinkConfig = new Configuration();
        // 设置YARN队列,对应CLI的--yarnqueue参数
        flinkConfig.setString(YarnConfigOptions.APPLICATION_QUEUE, "test");
        
        // 配置SkyWalking Agent的Java启动参数
        String skywalkingAgentOpts = "-javaagent:/opt/flink-1.13.6/skywalking-agent/skywalking-agent.jar=agent.service_name=540yh::calcu-ontime-task,collector.backend_service=192.168.100.22:30180";
        flinkConfig.setString("env.java.opts", skywalkingAgentOpts);
        
        // 创建YARN集群描述器
        YarnClusterDescriptor yarnDescriptor = new YarnClusterDescriptor(flinkConfig);
        
        // 加载待提交的作业Jar包
        PackagedProgram program = PackagedProgram.newBuilder()
                .setJarFile(new File("yourFlinkApp.jar"))
                .build();
        
        // 部署集群并提交作业
        ClusterClientProvider<?> clientProvider = yarnDescriptor.deployApplicationCluster(flinkConfig, program.getJobGraph());
        ClusterClient<?> clusterClient = clientProvider.getClusterClient();
        
        // 提交作业并输出作业ID
        clusterClient.submitJob(program.getJobGraph())
                .thenAccept(jobId -> System.out.println("作业提交成功,ID: " + jobId));
    }
}

补充说明

  • 如果使用Flink REST API提交作业,可在提交请求的JSON体中,将env.java.opts加入configuration字段中,示例如下:
{
  "entryClass": "com.your.package.YourFlinkJob",
  "programArgs": "",
  "jobName": "calcu-ontime-task",
  "configuration": {
    "env.java.opts": "-javaagent:/opt/flink-1.13.6/skywalking-agent/skywalking-agent.jar=agent.service_name=540yh::calcu-ontime-task,collector.backend_service=192.168.100.22:30180",
    "yarn.application.queue": "test"
  },
  "savepointPath": ""
}
  • 确保YARN所有节点上都能访问到SkyWalking Agent Jar包的路径,否则会导致作业启动失败。

内容的提问来源于stack exchange,提问作者小梨喵

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.25 17:36:12