【发布时间】:2021-08-27 14:01:29
【问题描述】:
由于 spark 提供了作业失败检测和重试机制的功能,我想在我的代码中使用这个功能,这意味着在某些情况下我想将 spark Job 标记为失败,然后 spark 会再次重试该作业。 我试图找出 Job/Stage 状态,我找到了this。并实施:
JavaSparkContext jsc = JavaSparkContext.fromSparkContext(this.spark.sparkContext());
JavaSparkStatusTracker statusTracker = jsc.statusTracker();
for(int jobId: statusTracker.getActiveJobIds()) {
SparkJobInfo jobInfo = statusTracker.getJobInfo(jobId);
for(int stageId: jobInfo.stageIds()) {
SparkStageInfo stageInfo = statusTracker.getStageInfo(stageId);
LOGGER.warn("Stage id=" + stageId + "; name = " + stageInfo.name()
+ "; completed tasks:" + stageInfo.numCompletedTasks()
+ "; active tasks: " + stageInfo.numActiveTasks()
+ "; all tasks: " + stageInfo.numTasks()
+ "; submission time: " + stageInfo.submissionTime());
}
}
在这里,我可以找到作业/任务/阶段状态,但是有什么方法(API)可以将 Spark 作业标记为失败?另外,实现重试机制而不是编写自定义代码进行重试是否是一种好方法?
【问题讨论】: