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:
- Create a jar with my job (for example simple stream from one kafka topic to another with data transformation);
- By usage of PackagedProgramUtils create JobGraph from jar file;
- 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