Apache spark作业在没有重试的情况下立即失败,设置maxFailures不起作用

tri*_*oid 2 failover self-healing apache-spark

我正在我的计算机上本地测试Apache Spark上的Web爬行/报废程序.

该程序使用一些RDD转换,这些转换采用零星失败的易失性函数.(该函数的目的是将URL链接转换为网页,有时它调用的无头浏览器只是停电或过载 - 我无法避免这种情况)

我听说Apache Spark具有强大的故障转移和重试功能,任何不成功的转换或丢失的数据都可以从头开始从它可以找到的任何资源中重新计算(听起来像魔法吗?)所以我没有在我的任何故障转移或try-catch中码.

这是我的火花配置:

val conf = new SparkConf().setAppName("MoreLinkedIn")
conf.setMaster("local[*]")
conf.setSparkHome(System.getenv("SPARK_HOME"))
conf.setJars(SparkContext.jarOfClass(this.getClass).toList)
conf.set("spark.task.maxFailures","40") //definitely enough
Run Code Online (Sandbox Code Playgroud)

不幸的是,在大多数阶段和个别任务成功之后,工作失败了.最新的登录控制台显示:

Exception in thread "main" org.apache.spark.SparkException: Job aborted due to stage failure: Task 1.0:7 failed 1 times, most recent failure: Exception failure in TID 23 on host localhost: org.openqa.selenium.TimeoutException: Timed out after 50 seconds waiting for...
Run Code Online (Sandbox Code Playgroud)

看起来Spark只是在失败一次之后放弃了懦弱.如何正确配置以使其更加顽强?

(我的程序可以从https://github.com/tribbloid/spookystuff下载,对于稀缺和杂乱无章的代码/文档,我只是启动了几天)

ADD:如果您想自己尝试一下,以下代码可以演示此问题:

def main(args: Array[String]) {
val conf = new SparkConf().setAppName("Spark Pi")
conf.setMaster("local[*]")
conf.setSparkHome(System.getenv("SPARK_HOME"))
conf.setJars(SparkContext.jarOfClass(this.getClass).toList)
conf.set("spark.task.maxFailures","400000")
val sc = new SparkContext(conf)
val slices = if (args.length > 0) args(0).toInt else 8
val n = 100000 * slices
val count = sc.parallelize(1 to n, slices).map { i =>
  val x = java.lang.Math.random()
  if (x > 0.9) throw new IllegalStateException("the map has a chance of 10% to fail")
  x
}.reduce(_ + _)
sc.stop()
println("finished")
}
Run Code Online (Sandbox Code Playgroud)

应该注意的是,在这篇文章中重复了32次相同的IllegalStateException: Apache Spark抛出java.lang.IllegalStateException:未读块数据

bot*_*que 8

我知道这是一个非常古老的问题,但我遇到了完全相同的问题,并在寻找解决方案时遇到了这个问题。

有 3 种主 URL 格式可以在本地模式下提交 Spark 应用程序:

  • local - 一个线程(无并行性),无重试
  • local[K](或local[*]) - 使用K(或核心数)工作线程并设置task.maxFailures为1(见这里)

  • local[K, F](或local[*, F])- 设置task.maxFailures=F,这就是我们所追求的。

有关详细信息,请参阅Spark 文档。


tri*_*oid 7

让我转发最权威的答案:

如果这是本地模式的一个有用功能,我们应该打开一个JIRA来记录设置或改进它(我更喜欢添加spark.local.retries属性而不是特殊的URL格式).我们最初禁用它除了单元测试以外的所有内容,因为90%的时间本地模式中的异常意味着应用程序中存在问题,我们宁愿让用户立即调试,而不是多次重试任务并让他们担心为什么他们会收到这么多错误.

马太