【问题标题】:How to control Spark application per task/stage/job attempt?如何控制每个任务/阶段/作业尝试的 Spark 应用程序?
【发布时间】:2016-12-28 06:10:27
【问题描述】:

我想阻止 Spark 重试 Spark 应用程序,以防引发某些特定异常。我只想在满足某些条件的情况下限制重试次数。否则,我想要默认的重试次数。

请注意,Spark 应用程序只运行一个 Spark 作业。

我尝试在异常情况下设置javaSparkContext.setLocalProperty("spark.yarn.maxAppAttempts", "1");,但它仍然重试整个工作。

我提交 Spark 申请如下:

spark-submit --deploy-mode cluster theSparkApp.jar

我有一个用例,如果输出是由同一个作业的上一次重试创建的,我想删除它,但如果输出文件夹不为空(第一次重试),则作业失败。你能想到任何其他方法来实现这一点吗?

【问题讨论】:

  • 如何提交 Spark 应用程序进行部署?使用的命令行选项和 Spark 属性是什么?顺便说一句,即使你说你的意思是“整个 Spark 应用程序”,你仍然使用“整个工作仍在重试”。一个 Spark 应用程序可以运行/提交一个或多个 Spark 作业。
  • 能否请您改用spark-submit --deploy-mode cluster --conf spark.yarn.maxAppAttempts=1(并在命令行中使用 Spark 设置)。

标签: apache-spark hadoop-yarn


【解决方案1】:

我有一个用例,如果输出是由同一个作业的上一次重试创建的,我想删除它,但如果输出文件夹不为空(第一次重试),则作业失败。你能想到任何其他方法来实现这一点吗?

您可以使用TaskContext 来控制您的 Spark 作业的行为方式,例如重试次数,如下所示:

val rdd = sc.parallelize(0 to 8, numSlices = 1)

import org.apache.spark.TaskContext

def businessCondition(ctx: TaskContext): Boolean = {
  ctx.attemptNumber == 0
}

val mapped = rdd.map { n =>
  val ctx = TaskContext.get
  if (businessCondition(ctx)) {
    println("Failing the task because business condition is met")
    throw new IllegalArgumentException("attemptNumber == 0")
  }
  println(s"It's ok to proceed -- business condition is NOT met")
  n
}
mapped.count

【讨论】:

  • 这里的问题是我不知道我的工作是否由于 businessCondition() 或其他原因而第一次重试失败(除非我在 Spark 之外的某个地方保持这种状态,我想要避免)。所以我能想到的唯一可能的方法是在满足 businessCondition() 的情况下强制 Spark 失败而不重试。
  • addTaskCompletionListener(listener: TaskCompletionListener): TaskContextaddTaskFailureListener(listener: TaskFailureListener): TaskContext 怎么样?以前从未使用过它们,但它们看起来好像可以在这里提供帮助。
  • onApplicationEnd 可以工作,但无法在 SparkListener 中获取 TaskContext。我需要 TaskContext 来确定要删除的确切内容。另外,我不确定是否可以在侦听器界面中找到应用程序是成功还是失败。
猜你喜欢
  • 2017-11-28
  • 2018-09-14
  • 2015-01-09
  • 1970-01-01
  • 2021-02-19
  • 1970-01-01
  • 2015-06-30
  • 2014-03-22
  • 2020-07-27
相关资源
最近更新 更多