【问题标题】:Concurrent operations in spark streaming火花流中的并发操作
【发布时间】:2015-10-29 06:49:02
【问题描述】:

我想了解有关 Spark 流执行的内部机制。

如果我有一个流 X,并且在我的程序中我将流 X 发送到函数 A 和函数 B:

  1. 在函数 A 中,我在 X->Y->Z 上执行了一些变换/过滤操作等以创建流 Z。现在我在 Z 上执行 forEach 操作并将输出打印到文件中。

  2. 然后在函数 B 中,我减少流 X -> X2(例如每个 RDD 的最小值),并将输出打印到文件

每个 RDD 是否都并行执行这两个函数?它是如何工作的?

谢谢

--- 来自 Spark 社区的评论----

我正在添加来自 spark 社区的 cmets -

如果您在驱动程序中的两个线程中执行收集步骤(foreach 在 1 中,可能在 2 中减少),那么它们都将并行执行。无论哪个先提交给 Spark,都会先执行 - 如果您需要确保执行顺序,您可以使用信号量,但我认为顺序无关紧要。

【问题讨论】:

    标签: apache-spark spark-streaming


    【解决方案1】:

    @Eswara 的答案似乎是正确的,但它不适用于您的用例,因为您的单独转换 DAG(X->Y->ZX->X2)在 X 中有一个共同的 DStream 祖先。这意味着当运行操作以触发这些流中的每一个时,转换 X->Y 和转换 X->X2 不能同时发生。将会发生的是 RDD X 的分区将被计算或从内存(如果缓存)以非并行方式分别为这些转换中的每一个加载。

    理想情况下,转换X->Y 将解析,然后转换Y->ZX->X2 将并行完成,因为它们之间没有共享状态。我相信 Spark 的流水线架构会为此进行优化。您可以通过持久化 DStream X 来确保在 X->X2 上更快地计算,以便它可以从内存中加载,而不是重新计算或从磁盘加载。有关持久性的更多信息,请参阅here

    如果您可以提供复制存储级别 *_2(例如 MEMORY_ONLY_2MEMORY_AND_DISK_2)以便能够在同一源上同时运行转换,那将会很有趣。我认为这些存储级别目前只有useful against lost partitions,因为将处理重复的分区来代替丢失的分区。

    【讨论】:

    • -感谢您的解释。这本质上是我的想法。我将尝试执行该应用程序并查看其行为。我不是在集群上运行它,而是现在在本地模式下运行它。我会用结果回复您,如果确实是这种行为,请接受您的回答:)
    • 再跟进。我在每个 forEach 操作结束时调用 rdd.collect(),并遍历返回的列表。我了解收集功能将此信息发送给驱动程序。列表的遍历是顺序的还是并行的(我假设它在驱动程序内部)。列表的遍历也在forEach()的回调函数内
    • 您是正确的,通过collect()传输的数据被发送回驱动程序。您随后对该数据进行的任何遍历都将仅在驱动程序上执行,因此请注意您发回的数据。虽然there are ways 在一台机器上并行处理这些数据(取决于您的驱动程序架构),但通常这确实意味着您将按顺序处理驱动程序数据。
    • 嘿@tsar2512,有帮助吗?
    【解决方案2】:

    是的。 它类似于使用 DAG 和惰性评估的 spark 执行模型,只是流式处理在每批新数据上重复运行 DAG。 在您的情况下,由于完成每个操作(您拥有的 2 个 foreach 中的每一个)所需的 DAG(或较大 DAG 的子 DAG,如果您更愿意这样称呼)没有公共链接一直到源代码,它们完全并行运行。流应用程序作为一个整体在应用程序提交给资源管理器时分配给每个执行器的 X 执行器(JVM)和 Y 核心(线程)。在任何时候,给定的X*Y 任务中的任务(即线程)将执行这些 DAG 中的 一个 的一部分或全部。请注意,应用程序的任何 2 个给定线程,无论是否在同一个执行程序中,都可以执行同一应用程序的不同操作。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2017-04-27
      • 1970-01-01
      • 2019-04-02
      • 2016-02-07
      • 2015-05-15
      • 1970-01-01
      相关资源
      最近更新 更多