Deploying job to Flink cluster

Viewed 170

I want to deploy configured flink job to standalone cluster by java application. I'd like to do it with RestClusterClient.

As far as I understand, I can do it the following way:

  1. Create a jar with my job (for example simple stream from one kafka topic to another with data transformation);
  2. By usage of PackagedProgramUtils create JobGraph from jar file;
  3. Initialize RestClusterClient and submit JobGraph to cluster.
PackagedProgram packagedProgram = PackagedProgram.newBuilder()
                                                 .setJarFile(new File("path/to/my/jar/file"))
                                                 .setArguments(arguments)
                                                 .build();
JobGraph jobGraph = PackagedProgramUtils.createJobGraph(packagedProgram, flinkConfiguration, 1, false);
try (RestClusterClient<StandaloneClusterId> client = new RestClusterClient<>(flinkConfiguration, StandaloneClusterId.getInstance())) {
    client.submitJob(jobGraph).get();
}

However, I want to configure my job and deploy it to cluster in one java application. I found the following way to do that:

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
// add source, any data processing and sink
StreamGraph streamGraph = graphFlow.getConfiguredEnvironment().getStreamGraph("myGraphName", false);
JobGraph graph = streamGraph.getJobGraph();
try (RestClusterClient<StandaloneClusterId> client = new RestClusterClient<>(flinkConfiguration, StandaloneClusterId.getInstance())) {
    client.submitJob(jobGraph).get();
}

This solution does not seem to be acceptable as StreamExecutionEnvironment methods are annotated as @Internal. Which another ways can be used for building JobGraph object that can be uploaded with RestClusterClient?

As I need to upload the job to flink cluster, I cannot use RemoteStreamEnvironment.

Flink version is 1.11.2

0 Answers
Related