【发布时间】: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 票证的讨论谈论的是“结果任务”而不是“行动”。例如here 和here。这似乎表明结果确实可以保证是正确的,因为我们在和 action 之前立即使用acc,因此应该作为单个阶段计算。
我猜这个“结果任务”的概念与所涉及的操作类型有关,它是最后一个包含操作的操作,就像在这个例子中一样,它显示了几个操作是如何分成阶段的(洋红色,图片取自here):
所以假设,该链末端的count() 动作将是同一最终阶段的一部分,并且我可以保证最后一张地图上使用的累加器不会包含任何重复项?
澄清这个问题会很棒!谢谢。
【问题讨论】:
-
好吧,赏金期结束了,我仍然不知道真正的答案,所以将它授予迄今为止评论最高的答案:-S
-
data.count 不会运行 data.map(...) 但这会做 >>>val data2 = data.map(x => acc += x; f(x)) > >>data2.count()
标签: apache-spark