【问题标题】:Why mapPartitionsWithIndex cause a shuffle in Spark?为什么 mapPartitionsWithIndex 会导致 Spark 中的洗牌?
【发布时间】:2015-07-25 14:25:48
【问题描述】:

我是 Spark 的新手。我正在检查测试应用程序中的洗牌问题,但我不知道为什么在我的程序中 mapPartitionsWithIndex 方法会导致洗牌!正如你在图片中看到的,我的初始 RDD 有两个 16MB 分区,Shuffle 写入大约 49.8 MB。 我知道mapmapPartitionmapPartitionsWithIndex 不像groupByKey 那样进行改组转换,但我发现它们也会导致Spark 中的改组。为什么?

【问题讨论】:

    标签: apache-spark shuffle rdd


    【解决方案1】:

    我认为您在 mapPartitionsWithIndex 之后执行了一些加入/组操作,这导致了随机播放。

    你可以通过修改你的代码来建立它。

    当前代码

    val rdd = inputRDD1.mapPartitionsWithIndex(....)
    val outRDD = rdd.join(inputRDD2)
    

    修改后的代码

    val rdd = inputRDD1.mapPartitionsWithIndex(....)
    println(rdd.count)
    

    【讨论】:

    • 但是,如果是这样,为什么 spark 将此方法视为 Stage,而不将下一个 join/group 操作视为阶段中的最终操作?
    • 一个阶段对应于所有执行相同代码的任务集合,每个任务都在不同的数据子集上。每个阶段都包含一系列转换,无需打乱完整数据即可完成。 - blog.cloudera.com/blog/2015/03/…
    • 我知道,但我记得每个Stage的最终操作通常是一个shuffle操作,如果我们假设在这个map方法之后有一个join/group操作,那么它应该是最终操作在这个阶段并显示而不是 mapPartitionWithIndex。
    • 不,那将是下一阶段。对于简单的连接工作流程,需要 3 个阶段。第一个用于加载 rdd1+transformation + shuffle 输出,第二个用于加载 rdd2 +transformation + shuffle 输出,第三个用于读取 shuffle 输出并加入这些提要,然后保存到磁盘或执行其他任何操作。
    • 我修改了你告诉我的代码,即count rdd 紧跟在mapPartitionWithIndex 之后,但没有迹象表明为什么会发生洗牌。 i.imgur.com/aPAYFpx.png
    猜你喜欢
    • 2016-07-24
    • 2018-04-21
    • 1970-01-01
    • 2015-04-08
    • 1970-01-01
    • 2018-11-30
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多