【问题标题】:Can I force a DAG dependency in Spark?我可以在 Spark 中强制使用 DAG 依赖项吗?
【发布时间】:2016-10-25 02:19:37
【问题描述】:

假设我有这样的代码,其中 x 是一个 RDD。

val a = x.map(...)
val b = x.map(...)
val c = x.map(...)
....
// a, b and c are used later on

出于某种原因,我希望 b 仅在 a 执行完成后执行,而 c 仅在执行 for 后执行b 已完成。有没有办法在 Spark 中强制这种依赖关系?

其次,Spark 执行此类代码的默认机制是什么。 abc 是否会并行执行,因为它们之间没有依赖关系?

【问题讨论】:

  • 什么是x?如果他们常规的 scala 仿函数将一个接一个地执行。
  • x 只是一个 RDD。
  • 有文件吗?这类东西通常决定了事情的运作方式。
  • stackoverflow.com/a/31384084/2784721,是的,它们将并行执行,您需要在地图操作结束时添加操作/终端功能。 Spark 很懒惰。顺便说一句,这是什么原因。

标签: scala apache-spark


【解决方案1】:

我所做的通常是重组我的代码,这样当我调用一个强制评估 b 的动作时,它必须在 a 之后调用。我还没有看到直接操作 DAG 的方法。

关于你的第二个问题。 a,b 和 c 在一个动作之前不会被执行。如果动作被并行调用,例如在 Futures 中,那么它们将由 spark 并行运行。如果都按顺序调用,spark会按顺序调用

澄清一下。如果你调用 a.count()。那会阻塞。因此,强制 b 和 c 的动作无法执行,因此它们不会并行运行。但是,如果您调用类似:

Future(a.count())
Future(b.count())
Future(c.count())

然后地图将并行发生。每个动作都将被调用,这将强制对阶段进行评估。然后 spark 将根据您拥有的执行器核心总数来处理每个阶段的任务。

【讨论】:

  • 当然我假设a、b 和c 将在以后使用。在那种情况下,他们当然会被处决。但问题是,它们会并行执行吗?
猜你喜欢
  • 2011-11-18
  • 1970-01-01
  • 2021-02-27
  • 1970-01-01
  • 2019-04-18
  • 1970-01-01
  • 2011-07-02
  • 2014-08-03
  • 1970-01-01
相关资源
最近更新 更多