我有简单的java代码用于flink job
List<Tuple2> list = new ArrayList<>();
for (int i = 0; i < 10; i++) {
list.add(new Tuple2(Integer.valueOf(i), "test" + i));
}
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.fromCollection(list).print();
env.execute("job1");
我打包了这段代码并创建了一个jar:比方说flink-processor-0.1-snapshot.jar,从提交作业ui上传到jobmanager。上传中没有问题。我看到entryclass有主类(com..xyz.streaming.flinkprocessor)现在,我从“提交作业ui”提交作业,并使用一些参数(--ns.conf1 .file--ns.conf2 xyz.file)指定主类(com..xyz.streaming.flinkprocessor)。作业提交失败。在jobmanager中,我看到以下错误。
org.apache.flink.client.program.OptimizerPlanEnvironment$ProgramAbortException
at org.apache.flink.streaming.api.environment.StreamPlanEnvironment.execute(StreamPlanEnvironment.java:70)
at org.apache.flink.streaming.api.environment.StreamPlanEnvironment.execute(StreamPlanEnvironment.java:53)
at com.abc.xyz.streaming.FlinkProcessor.run(FlinkProcessor.java:114)
at com.abc.xyz.streaming.FlinkProcessor.main(FlinkProcessor.java:53)
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessor
at sun.reflect.NativeMethodAccessorImpl.invoke0(Native Method)
at sun.reflect.NativeMethodAccessorImpl.invoke(NativeMethodAccessorImpl.java:62)
at sun.reflect.DelegatingMethodAccessorImpl.invoke(DelegatingMethodAccessorImpl.java:43)
at java.lang.reflect.Method.invoke(Method.java:498)
at org.apache.flink.client.program.PackagedProgram.callMainMethod(PackagedProgram.java:529)
at org.apache.flink.client.program.PackagedProgram.invokeInteractiveModeForExecution(PackagedProgram.java:421)
at org.apache.flink.client.program.OptimizerPlanEnvironment.getOptimizedPlan(OptimizerPlanEnvironment.java:83)
at org.apache.flink.client.program.PackagedProgramUtils.createJobGraph(PackagedProgramUtils.java:80)
at org.apache.flink.runtime.webmonitor.handlers.utils.JarHandlerUtils$JarHandlerContext.toJobGraph(JarHandlerUtils.java:126)
at org.apache.flink.runtime.webmonitor.handlers.JarRunHandler.lambda$getJobGraphAsync$6(JarRunHandler.java:142)
at java.util.concurrent.CompletableFuture$AsyncSupply.run(CompletableFuture.java:1590)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:748)
2020-06-16 16:14:20,588 ERROR org.apache.flink.runtime.webmonitor.handlers.JarRunHandler - Unhandled exception.
org.apache.flink.client.program.ProgramInvocationException: The program plan could not be fetched - the program aborted pre-maturely.
System.err: Running flink job :2020-06-16T23:14:20.476Z
Error msg: org.apache.flink.client.program.OptimizerPlanEnvironment$ProgramAbortException
at org.apache.flink.streaming.api.environment.StreamPlanEnvironment.execute(StreamPlanEnvironment.java:70)
at org.apache.flink.streaming.api.environment.StreamPlanEnvironment.execute(StreamPlanEnvironment.java:53)
当我使用flink run命令提交同一个作业时,效果很好。没有错误。
flink run-c com..xyz.streaming.flinkprocessor/users//target/flink-processor-0.1-snapshot.jar--ns.conf1 .file--ns.conf2 xyz.file
不知道我在这里错过了什么?非常感谢您的帮助。
1条答案
按热度按时间kknvjkwl1#
发现了问题。我正在用
更改了异常块,请尝试{….}catch(exception exp){ignore exp};在这之后,它开始工作。
谢谢!