【问题标题】:When are accumulators truly reliable?蓄电池何时真正可靠?
【发布时间】:2015-04-07 14:46:59
【问题描述】:

我想使用累加器来收集有关我在 Spark 作业中处理的数据的一些统计信息。理想情况下,我会在作业计算所需的转换时这样做,但由于 Spark 会在不同情况下重新计算任务,因此累加器不会反映真实的指标。以下是文档对此的描述:

对于仅在操作内部执行的累加器更新,Spark 保证每个任务对累加器的更新只会 应用一次,即重新启动的任务不会更新该值。在 转换,用户应该知道每个任务的更新可能 如果重新执行任务或作业阶段,则应用不止一次。

这很令人困惑,因为大多数 动作 不允许运行自定义代码(可以使用累加器),它们主要从以前的转换中获取结果(懒惰地)。该文档还显示了这一点:

val acc = sc.accumulator(0)
data.map(x => acc += x; f(x))
// Here, acc is still 0 because no actions have cause the `map` to be computed.

但是如果我们在末尾添加data.count(),这是否可以保证正确(没有重复)?显然acc 不是“仅用于内部操作”,因为 map 是一种转换。所以不能保证。

另一方面,有关 Jira 票证的讨论谈论的是“结果任务”而不是“行动”。例如herehere。这似乎表明结果确实可以保证是正确的,因为我们在和 action 之前立即使用acc,因此应该作为单个阶段计算。

我猜这个“结果任务”的概念与所涉及的操作类型有关,它是最后一个包含操作的操作,就像在这个例子中一样,它显示了几个操作是如何分成阶段的(洋红色,图片取自here):

所以假设,该链末端的count() 动作将是同一最终阶段的一部分,并且我可以保证最后一张地图上使用的累加器不会包含任何重复项?

澄清这个问题会很棒!谢谢。

【问题讨论】:

  • 好吧,赏金期结束了,我仍然不知道真正的答案,所以将它授予迄今为止评论最高的答案:-S
  • data.count 不会运行 data.map(...) 但这会做 >>>val data2 = data.map(x => acc += x; f(x)) > >>data2.count()

标签: apache-spark


【解决方案1】:

回答“蓄电池什么时候真正可靠?”

答案:当它们出现在 Action 操作中时。

根据 Action Task 中的文档,即使存在任何重新启动的任务,它也只会更新 Accumulator 一次。

对于仅在操作内部执行的累加器更新,Spark 保证每个任务对累加器的更新只会应用一次,即重新启动的任务不会更新值。在转换中,用户应注意,如果重新执行任务或作业阶段,每个任务的更新可能会应用多次。

Action 确实允许运行自定义代码。

例如

val accNotEmpty = sc.accumulator(0)
ip.foreach(x=>{
  if(x!=""){
    accNotEmpty += 1
  }
})

但是,为什么地图+动作,即。结果任务操作对于累加器操作不可靠

  1. 由于代码中的一些异常,任务失败。 Spark 将尝试 4 次(默认尝试次数)。如果任务每次失败都会给出异常。如果偶然成功,则 Spark 将继续并更新成功状态的累加器值,而失败状态的累加器值将被忽略。判决:妥善处理
  2. Stage Failure :如果执行器节点崩溃,则不是用户故障,而是硬件故障 - 如果节点在 shuffle 阶段出现故障。由于 shuffle 输出存储在本地,如果节点出现故障,则该 shuffle 输出消失。所以 Spark 回到生成 shuffle 输出的阶段,查看哪些任务需要重新运行,并在其中一个还活着的节点上执行它们。在我们重新生成丢失的 shuffle 输出之后,生成 map 输出的阶段已多次执行其中的一些任务。Spark 计算所有这些任务的累加器更新。
    结论:未在结果任务中处理。累加器将给出错误的输出。
  3. 如果某个任务运行缓慢,Spark 可以在另一个节点上启动该任务的推测副本。
    结论:未处理。累加器将给出错误的输出。
  4. 缓存的 RDD 很大,不能驻留在内存中。因此,无论何时使用 RDD,它都会重新运行 Map 操作以获取 RDD,并再次更新累加器。
    结论:未处理.Accumulator 会给出错误的输出。

因此可能会发生相同的函数可能会在相同的数据上运行多次。因此 Spark 不提供任何保证累加器因为 Map 操作而得到更新的保证。

所以在 Spark 中最好使用 Accumulator in Action 操作。

要了解有关 Accumulator 及其问题的更多信息,请参阅 Blog Post - 作者:Imran Rashid。

【讨论】:

  • 您好!您基本上引用了我在问题本身和@Daniel Daravos 回答中的相同内容,但我认为这不是全部情况,因为文档似乎与它本身以及来自 cloudera 的其他专家分析相冲突。您是 Spark 代码/设计贡献者,还是像我们其他人一样只是从外部搞清楚?
  • @DanielL。 : 不,我只是 Spark 用户。但是已经为此工作了很长时间。实际上,我想为答案添加 foreach() 部分,因为可以自定义操作。因此,将来如果有任何 OP 出现,他们可以很好地理解累加器。
  • 当然,但我认为您弄错了(文档不清楚)。事实上,我目前的理解是,累加器在“结果任务”而不是动作上确实可靠。例如,如果您只执行sc.readFile(...).map(...<accumulator use here> ...).saveAsFile(...),则每个任务将只计算一次,并且累加器将是可靠的,因为整个操作将在惰性求值下作为一个单元发生,没有可能重新运行的中间结果(推测与否) .到目前为止,我的经验反映了这一点,这就是我寻找权威答案的原因。
  • 是的,但是如果执行器节点在执行操作时失败。那么我认为结果会有所不同。以及不应该使用map + action = result的原因。 val a = sc.readFile(...); val b = a.map(...<accumulator used here>...);b.saveAsTextFile(...);val c = b.map(....);println(c.count());.现在,如果您检查该值,它将使累加器原始值加倍。因为该映射用于执行操作两次。如果你在行动中使用了累加器。那么你会确定它会调用一次。而且我不会在这里处理缓存问题,就好像 RDD 被驱逐一样。它会再次重新运行 map。
  • 但这确实是我的观点,鉴于这个例子,你可以依靠累加器是 2x :-)
【解决方案2】:

累加器更新会在任务成功完成后发送回驱动程序。因此,当您确定每个任务都将执行一次并且每个任务都按照您的预期执行时,您的累加器结果就可以保证是正确的。

我更喜欢依赖 reduceaggregate 而不是累加器,因为很难列举所有可以执行任务的方式。

  • 动作启动任务。
  • 如果某个操作依赖于较早的阶段,并且该阶段的结果未(完全)缓存,则将启动较早阶段的任务。
  • 检测到少量慢速任务时,推测执行会启动重复任务。

也就是说,有很多简单的情况下可以完全信任累加器。

val acc = sc.accumulator(0)
val rdd = sc.parallelize(1 to 10, 2)
val accumulating = rdd.map { x => acc += 1; x }
accumulating.count
assert(acc == 10)

这是否可以保证是正确的(没有重复)?

是的,如果推测执行被禁用。 mapcount 将是一个阶段,所以就像你说的,一个任务不可能多次成功执行。

但累加器会作为副作用更新。因此,在考虑如何执行代码时,您必须非常小心。考虑这个而不是accumulating.count

// Same setup as before.
accumulating.mapPartitions(p => Iterator(p.next)).collect
assert(acc == 2)

这也会为每个分区创建一个任务,并且每个任务将保证只执行一次。但是map 中的代码不会在所有元素上执行,只会在每个分区中的第一个元素上执行。

累加器就像一个全局变量。如果您共享对可以递增累加器的 RDD 的引用,那么其他代码(其他线程)也可以导致累加器递增。

// Same setup as before.
val x = new X(accumulating) // We don't know what X does.
                            // It may trigger the calculation
                            // any number of times.
accumulating.count
assert(acc >= 10)

【讨论】:

  • 我了解这些情况。但是您如何解释官方文档:“对于仅在操作内部执行的累加器更新,Spark 保证每个任务对累加器的更新只会应用一次,即重新启动的任务不会更新值”???鉴于此,在某些情况下它们肯定必须可靠吗?我不认为“从不”是答案,它可能比这更复杂......
  • 你是对的。我刚刚对其进行了测试,并且不计算来自失败任务的累加器更新。然后我想如果你禁用了投机执行,只要你确定没有什么可以多次触发计算,累加器可能是值得信赖的。我会仔细看看代码。我可能不得不修正我对累加器的深深不信任:)。
  • 你我都有! ;-) 我对代码的看法给我留下了比答案更多的问题,但如果你确实得出了一些结论,请告诉我。如果你这样做,你会得到赏金,呵呵。
  • 是的,我知道驱动程序会跟踪失败的任务(重试),并且可以对更新进行更粗粒度的控制。此外,我相信在完成所有工作人员任务之前,不会将报告发送回驱动程序。不过,我会尝试在一天左右的时间内挖掘更多内容
  • 感谢您的赏金!我在阅读代码后重写了答案。我希望它现在更值得。我仍然会避免使用累加器,因为在一般情况下很难对它们进行推理。
【解决方案3】:

我认为 Matei 在参考文档中回答了这个问题:

正如在https://github.com/apache/spark/pull/2524 上讨论的那样,这是 在一般情况下很难提供良好的语义 (非结果阶段内的累加器更新),用于以下 原因:

  • 一个 RDD 可以作为多个阶段的一部分进行计算。为了 例如,如果您更新 MappedRDD 中的累加器,然后 洗牌,那可能是一个阶段。但是如果你再调用 map() 在 MappedRDD 上,然后随机播放结果,你会得到第二个 该地图是管道的阶段。你要数这个吗 累加器是否更新两次?

  • 如果出现以下情况,可以重新提交整个阶段 shuffle 文件被定期清理器删除或由于 节点故障,因此任何跟踪 RDD 的东西都需要这样做 很长一段时间(只要 RDD 在用户中是可引用的) 程序),这将是相当复杂的实现。

所以我要走了 暂时将其标记为“不会修复”,结果部分除外 在 SPARK-3628 中完成的阶段。

【讨论】:

  • 我确实读过。但什么构成“结果”阶段的定义尚不清楚。我使用此文本来撰写问题并推断替代方案,但我无法验证假设、推断是否正确。我希望有人了解该过程来验证它。 (并可能进一步解释)。不过谢谢!
猜你喜欢
  • 1970-01-01
  • 2021-06-19
  • 1970-01-01
  • 2014-02-23
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多