【问题标题】:How to know the status of a flink job from Java?如何从 Java 中了解 flink 作业的状态?
【发布时间】:2020-06-18 22:29:32
【问题描述】:
我有一个工作正在运行,我有兴趣只使用一次恢复重试,因为在此期间没有触发这个 flink 重启我有一个尝试解决问题的线程,然后当问题解决时flink 将重新启动,但有时线程通常需要更长的时间来解决问题并触发重新启动策略,由于问题仍然失败,然后作业停止但线程可能有另一个迭代,然后应用程序永远不会死,因为我'将其作为 jar 应用程序运行。所以,我的问题:
- 是否可以从 java 代码中知道作业的状态?类似于 (JobStatus.CANCELED == true)。
提前致谢!
亲切的问候
【问题讨论】:
标签:
apache-flink
flink-streaming
flink-cep
【解决方案1】:
非常感谢费利佩。这就是我所需要的,多亏了你,它已经完成了。我在这里分享代码以防其他人需要。
-
准备监听器
final StreamExecutionEnvironment env = StreamExecutionEnvironment.createLocalEnvironment(...);
final AtomicReference<JobID> jobIdReference = new AtomicReference<>();
//Environment configurations
env.registerJobListener(new JobListener() {
@Override
public void onJobSubmitted(@Nullable JobClient jobClient, @Nullable Throwable throwable) {
assert jobClient != null;
jobIdReference.set(jobClient.getJobID());
jobClient = jobClient /*jobClient static public object in the main class*/;
}@Override
public void onJobExecuted(@Nullable JobExecutionResult jobExecutionResult, @Nullable Throwable throwable) {
assert jobExecutionResult != null;
jobExecutionResult.notify();
}
});
-
使用代码:
Preconditions.checkNotNull(jobClient);
final String status = jobClient.getJobStatus().get().name();
if (status.equals(JobStatus.FAILED.name())) System.exit(1);