【问题标题】:Spark: disk I/O on stage boundaries explanationSpark:舞台边界上的磁盘 I/O 解释
【发布时间】:2020-03-01 03:38:54
【问题描述】:

我在官方文档中找不到关于Spark临时数据持久化的信息,只能在this等一些Spark优化文章中找到:

在每个阶段边界,数据由父级中的任务写入磁盘 阶段,然后由子阶段中的任务通过网络获取。 因为它们会产生大量的磁盘和网络 I/O,所以阶段边界可以是 价格昂贵,应尽可能避免使用。

每个阶段边界上的磁盘持久性是否总是适用于:HashJoin 和 SortMergeJoin?为什么 Spark(内存引擎)在 shuffle 之前对 tmp 文件进行持久化?这是为了任务级恢复还是其他什么?

附:问题主要与 Spark SQL API 有关,而我也对 Streaming & Structured Streaming 感兴趣

UPD:在"Stream Processing with Apache Spark book" 找到了关于为什么会发生的提及和更多详细信息。在参考页面上查找“任务故障恢复”和“阶段故障恢复”主题。据我了解,Why = recovery,When = always,因为这是 Spark Core 和 Shuffle Service 的机制,负责数据传输。此外,所有 Spark 的 API(SQL、流式处理和结构化流式处理)都基于相同的故障转移保证(Spark Core/RDD)。所以我想这是 Spark 的普遍行为

【问题讨论】:

  • @thebluephantom 听起来不错。 Spark 的广泛转换是否比 Map-Reduce 行为更快(更优化)?我从未使用过 Map-Reduce,但读到它有时会保留中间结果,然后再保留 map 输出
  • @thebluephantom 与进一步洗牌相比,此(磁盘 I/O)操作的成本是多少?
  • @thebluephantom 另一点是理论上我可以在同一个工人/机器上运行所有 mu 执行器。 Spark 会考虑舞台边界的局部性吗?
  • 与进一步洗牌相比,这个(磁盘 I/O)操作的成本是多少?很难回答,一段字符串有多长
  • @moin1010 1 个问题中的问题太多。改写它。这回答了部分问题。

标签: apache-spark apache-spark-sql


【解决方案1】:

这是一个很好的问题,因为我们听说过内存中 Spark 与 Hadoop,所以有点令人困惑。文档很糟糕,但我运行了一些东西并通过四处寻找最优秀的来源验证了观察结果:http://hydronitrogen.com/apache-spark-shuffles-explained-in-depth.html

假设已经调用了一个动作——如果没有说明,为了避免明显的注释,假设我们不是在谈论 ResultStage 和广播连接,那么我们正在谈论 ShuffleMapStage。我们最初看的是一个 RDD。

然后,借用网址:

  • 涉及随机播放的 DAG 依赖意味着创建单独的阶段。
  • Map 操作之后是 Reduce 操作和 Map 等等。

当前阶段

  • 所有(融合的)Map 操作都在 Stage 内执行。
  • 下一个阶段要求,Reduce 操作 - 例如一个reduceByKey,表示输出是散列或按键排序(K)在地图的末尾 当前阶段的操作。
  • 这些分组数据被写入 Executor 所在的 Worker 上的磁盘 - 或与该云版本相关联的存储。 (我会 如果数据很小,内存中的想法是可能的,但这是一个架构 Spark 文档中所述的方法。)
  • 通知 ShuffleManager 散列的映射数据可供下一阶段使用。 ShuffleManager 跟踪所有 完成所有地图方面的工作后的键/位置。

下一阶段

  • 下一个阶段是一个 reduce,然后通过咨询 Shuffle Manager 和使用 Block Manager 从这些位置获取数据。
  • Executor 可以重复使用,或者是另一个 Worker 上的新 Executor,或者是同一 Worker 上的另一个 Executor。

所以,我的理解是,在架构上,阶段意味着写入磁盘,即使有足够的内存。鉴于 Worker 的资源有限,这种类型的操作写入磁盘是有道理的。当然,更重要的一点是“Map Reduce”实现。我总结了出色的帖子,这是您的规范来源。

当然,这种持久性有助于容错,减少重新计算工作。

类似的方面也适用于 DF。

【讨论】:

    【解决方案2】:

    Spark 不是,也从来不是“内存引擎”。如果您检查内部结构,很明显它既没有针对内存处理进行优化,也没有针对以内存为中心的硬件进行调整。

    相反,几乎所有设计决策都清楚地假设整个数据的大小以及单个任务的输入和输出可以超过集群和单个执行器的可用内存量/executor线程分别。此外,它显然设计用于商用硬件。

    这样的实现可用于恢复或避免重新计算(参见例如What does "Stage Skipped" mean in Apache Spark web UI?),但这是重新调整用途而不是初始目标。

    【讨论】:

    • 您能否更深入地了解为什么 Spark 在数量级上比 Map-Reduce 快?每个教程都说因为 Spark 不会将中间结果保存到磁盘。我明白为什么窄转换更快,并不是所有的转换都可以用单个 Map-Reduce 的映射操作来表示。但是主要的性能影响是由于广泛的转换而发生的。那么为什么 Spark 的广泛转换比 Map-Reduce 更快?
    • @VB_ 每个教程都说因为 Spark 不会将中间结果保存到磁盘 - 他们不是在谈论同一件事。默认情况下,单个 Hadoop 作业(map-only 或 map-reduce)将数据持久化到分布式存储(即 HDFS)——这是作业之间唯一的通信媒介,通常需要磁盘和网络 IO 来写入(这里还有复制) 并在新工作中阅读。相比之下,Spark 阶段默认写入本地存储(就 Spark 而言可能是内存中的 FS)并且不复制数据(尽管对此行为进行了一些讨论)。
    • @user12357420 我错过了关于 DFS 与本地 FS 的观点,你是对的。谢谢你。你能提供任何讨论的链接吗?
    • @VB_Spark 可以在没有 HDFS 的情况下运行,也可以在没有数据本地性的云中运行。迭代处理和 Streaming,见:lightbend.com/blog/…
    • @VB_ 你能提供任何讨论的链接吗? - 我手头没有具体的链接,但你可以搜索 JIRA 和 dev list 以获取有关解耦块的讨论经理,尤其是在动态分配的情况下。 我错过了关于 DFS 与本地 FS 的要点 - 这里还有更多。例如,多个管道 Hadoop 作业将不得不单独获取资源,从而增加更多延迟。在高级别的 RDD 与沿袭 (Spark) 取代 DFS (Hadoop) 作为弹性机制,这是 Hadoop 的核心优势。
    猜你喜欢
    • 2019-12-07
    • 2014-05-20
    • 2016-07-28
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2019-11-12
    相关资源
    最近更新 更多