Spark - repartition()vs coalesce()

Pra*_*ati 208 distributed-computing apache-spark rdd

根据Learning Spark的说法

请记住,重新分区数据是一项相当昂贵的操作.Spark还有一个优化版本的repartition(),称为coalesce(),它允许避免数据移动,但前提是你减少了RDD分区的数量.

我得到的一个区别是,使用repartition()可以增加/减少分区数量,但是使用coalesce()时,只能减少分区数量.

如果分区分布在多台机器上并运行coalesce(),它如何避免数据移动?

Jus*_*ony 294

它避免了完全洗牌.如果已知该数字正在减少,则执行程序可以安全地将数据保存在最小数量的分区上,仅将数据从额外节点移出到我们保留的节点上.

所以,它会是这样的:

Node 1 = 1,2,3
Node 2 = 4,5,6
Node 3 = 7,8,9
Node 4 = 10,11,12
Run Code Online (Sandbox Code Playgroud)

然后coalesce下至2个分区:

Node 1 = 1,2,3 + (10,11,12)
Node 3 = 7,8,9 + (4,5,6)
Run Code Online (Sandbox Code Playgroud)

请注意,节点1和节点3不需要移动其原始数据.

  • 谢谢你的回复.文档应该更好地说"最小化数据移动"而不是"避免数据移动". (95认同)
  • @Niemand我认为当前的文档非常清楚:https://github.com/apache/spark/blob/128c29035b4e7383cc3a9a6c7a9ab6136205ac6c/core/src/main/scala/org/apache/spark/rdd/RDD.scala#L376 Keep记住所有`repartition`都是``coalesce`,`shuffle`参数设置为true.如果有帮助,请告诉我. (15认同)
  • 是否应该使用`repartition`而不是`coalesce`? (9认同)
  • 是否可以减少现有分区文件的数量?我没有hdfs,但文件很多。 (2认同)
  • @JustinPihony,感谢您的好回答。在答案中给出的示例中有“12 个分区”,那么“repartition(2)”是否等于“coalesce(2)”?或者 `repartition(2)` 会比 `coalesce(2)` 慢? (2认同)
  • 重新分配将在统计上更慢,因为它不知道它正在缩小......尽管可能他们可以优化它.在内部,它只是用`shuffle = true`标志调用coalesce (2认同)
  • @Niemand-是的,我的生产代码在运行repartition时比在coalesce中运行得更快。repartition将数据平均分配,因此输出到文件并运行分析更快。请参阅我的答案以获取更多详细信息。 (2认同)
  • 我们有更优化的合并版本吗? (2认同)

Pow*_*ers 137

贾斯汀的答案很棒,而且这种反应更深入.

该repartition算法执行完全随机播放并创建具有均匀分布的数据的新分区.让我们创建一个数字从1到12的数据框架.

val x = (1 to 12).toList
val numbersDf = x.toDF("number")
Run Code Online (Sandbox Code Playgroud)

numbersDf 我的机器上包含4个分区.

numbersDf.rdd.partitions.size // => 4
Run Code Online (Sandbox Code Playgroud)

以下是分区上数据的划分方式:

Partition 00000: 1, 2, 3
Partition 00001: 4, 5, 6
Partition 00002: 7, 8, 9
Partition 00003: 10, 11, 12
Run Code Online (Sandbox Code Playgroud)

让我们对该repartition方法进行全面改组,并在两个节点上获取此数据.

val numbersDfR = numbersDf.repartition(2)
Run Code Online (Sandbox Code Playgroud)

以下是numbersDfR我的机器上数据的分区方式:

Partition A: 1, 3, 4, 6, 7, 9, 10, 12
Partition B: 2, 5, 8, 11
Run Code Online (Sandbox Code Playgroud)

该repartition方法创建新分区并在新分区中均匀分布数据(对于较大的数据集,数据分布更均匀).

coalesce和之间的区别repartition

coalesce使用现有分区来最小化混洗的数据量. repartition创建新分区并进行完全洗牌. coalesce导致具有不同数据量的分区(有时具有大小不同的分区)并repartition导致大小相等的分区.

是更快coalesce还是repartition更快?

coalesce可能repartition比同等大小的分区运行速度更快,但不等大小的分区通常比较慢.在过滤大型数据集后,您通常需要重新分区数据集.我发现repartition整体速度更快,因为Spark可以使用相同大小的分区.

如果您想了解更多详情,请阅读此博客文章.

  • @Powers的答案很棒,但分区A和B中的数据是不是偏斜了?它是如何均匀分布的? (7认同)
  • @anwartheravian - 分区A和分区B的大小不同,因为`repartition`算法不会为非常小的数据集分配数据.我使用`repartition`将500万条记录组织成13个分区,每个文件介于89.3 MB和89.6 MB之间 - 这非常相同! (7认同)

Har*_* Ck 21

这里需要注意的另一点是,Spark RDD的基本原理是不变性.重新分区或合并将创建新的RDD.基本RDD将继续存在其原始分区数.如果用例要求在缓存中保留RDD,则必须对新创建的RDD执行相同操作.

scala> pairMrkt.repartition(10)
res16: org.apache.spark.rdd.RDD[(String, Array[String])] =MapPartitionsRDD[11] at repartition at <console>:26

scala> res16.partitions.length
res17: Int = 10

scala>  pairMrkt.partitions.length
res20: Int = 2
Run Code Online (Sandbox Code Playgroud)


kas*_*sur 17

什么从如下的代码和代码文档被认为coalesce(n)是一样的coalesce(n, shuffle = false),并repartition(n)是一样的coalesce(n, shuffle = true)

因此,coalesce和repartition都可用于增加分区数

使用shuffle = true,您实际上可以合并到更多的分区。如果您有少量分区(比如 100 个),可能有几个分区异常大,这将很有用。

另一个需要强调的重要说明是,如果您大幅减少分区数量,您应该考虑使用shuffled版本coalesce(与repartition这种情况相同)。这将允许您在父分区上并行执行计算(多任务)。

但是,如果您正在执行剧烈的合并,例如 to numPartitions = 1,这可能会导致您的计算在比您喜欢的节点数少(例如,在 的情况下为一个节点numPartitions = 1)。为避免这种情况,您可以通过shuffle = true. 这将添加一个 shuffle 步骤,但意味着当前的上游分区将并行执行(无论当前分区是什么)。

另请参阅此处的相关答案


Ale*_* S. 11

即使在@Rob 的回答中提到的分区号减少的情况下,重新分区 >> 合并也有一个用例,即将数据写入单个文件。

@Rob 的回答暗示了一个好的方向,但我认为需要一些进一步的解释来了解幕后发生的事情。

如果您需要在写入前过滤数据,那么重新分区比合并更合适,因为合并将在加载操作之前下推。

例如: load().map(…).filter(…).coalesce(1).save()

翻译成: load().coalesce(1).map(…).filter(…).save()

这意味着您的所有数据都将折叠到一个分区中,在那里将对其进行过滤,从而失去所有并行性。即使对于非常简单的过滤器(如column='value'.

重新分区不会发生这种情况: load().map(…).filter(…).repartition(1).save()

在这种情况下,过滤会在原始分区上并行进行。

只是为了给出一个数量级,在我的情况下,当从 Hive 表加载后过滤具有 ~1000 个分区的 109M 行(~105G)时,运行时间从 coalesce(1) 的 ~6h 下降到 repartition(1) 的 ~2m .

具体例子取自AirBnB 的这篇文章,非常不错,涵盖了 Spark 中重新分区技术的更多方面。


Rah*_*hul 10

重新分区:将数据洗牌到新数量的分区中。

例如。初始数据帧分为 200 个分区。

df.repartition(500):数据将从 200 个分区混洗到新的 500 个分区。

Coalesce:将数据混洗到现有数量的分区中。

df.coalesce(5):数据将从剩余的 195 个分区混洗到 5 个现有分区。


小智 9

所有的答案都在这个非常常见的问题中添加了一些丰富的知识。

因此,按照这个问题时间表的传统,这是我的2美分。

在非常特定的情况下,我发现重新分区的速度比合并更快。

在我的应用程序中,当我们估计的文件数量低于特定阈值时,重新分区的工作速度更快。

这就是我的意思

if(numFiles > 20)
    df.coalesce(numFiles).write.mode(SaveMode.Overwrite).parquet(dest)
else
    df.repartition(numFiles).write.mode(SaveMode.Overwrite).parquet(dest)
Run Code Online (Sandbox Code Playgroud)

在上面的代码段中,如果我的文件少于20个,则合并需要花费很多时间,而重新分区要快得多,因此上面的代码也是如此。

当然,这个数字(20)将取决于工作人员的数量和数据量。

希望能有所帮助。


小智 8

repartition -建议在不增加任何分区的情况下使用重新分区,因为它涉及对所有数据进行改组。

coalesce-建议在减少分区数量的同时使用合并。例如,如果您有3个分区,并且要将其减少到2个分区,则Coalesce会将第3个分区数据移至分区1和2。分区1和2将保留在同一Container中。执行器之间的匹配度很高,并且会影响性​​能。

性能明智的coalesce性能要比repartition减少分区数量要好。


Nik*_*iya 8

基本上,重新分区允许您增加或减少分区的数量。重新分区会重新分配所有分区中的数据,这会导致完全洗牌,这是非常昂贵的操作。

Coalesce 是 Repartition 的优化版本,您只能减少分区数量。由于我们只能减少分区的数量,因此它的作用是将一些分区合并为单个分区。与重新分区相比,通过合并分区,跨分区的数据移动量更低。因此,Coalesce 是最小数据移动,但说 Coalesce 不进行数据移动是完全错误的说法。

另一件事是通过提供分区数量进行重新分区,它尝试在所有分区上均匀地重新分配数据,而在合并的情况下,在某些情况下我们仍然可能会出现倾斜数据。


Sal*_*lim 7

我想补充贾斯汀和鲍尔的回答——

repartition将忽略现有分区并创建新分区。所以你可以用它来修复数据倾斜。您可以提及分区键来定义分布。数据倾斜是“大数据”问题空间中最大的问题之一。

coalesce将与现有分区一起工作并洗牌其中的一个子集。它不能像以前那样修复数据倾斜repartition。因此,即使它更便宜,它也可能不是您需要的。


小智 5

对于所有出色的答案,我想补充一点,这repartition是利用数据并行化的最佳选择之一。虽然coalesce提供了一个减少分区的廉价选项,并且在将数据写入 HDFS 或其他一些接收器以利用大写入时非常有用。

我发现这在以镶木地板格式写入数据以充分利用时很有用。


小智 5

对于从 PySpark (AWS EMR) 生成单个 csv 文件作为输出并将其保存在 s3 上时遇到问题的人,使用重新分区会有所帮助。原因是,合并不能进行完全洗牌,但重新分区可以。本质上,您可以使用重新分区来增加或减少分区数量,但使用合并只能减少分区数量(但不能减少 1 个)。以下代码适用于尝试将 csv 从 AWS EMR 写入 s3 的人:

df.repartition(1).write.format('csv')\
.option("path", "s3a://my.bucket.name/location")\
.save(header = 'true')
Run Code Online (Sandbox Code Playgroud)


Shy*_*pta 5

    \n
  • 合并\xc2\xa0 使用现有分区来最大程度地减少洗牌的数据量。\xc2\xa0Repartition\xc2\xa0 创建新分区\xc2\xa0 并\xc2\xa0 执行完整\n洗牌。

    \n
  • \n
  • 合并\xc2\xa0 会产生具有不同数据量的分区\n(有时是具有许多不同大小的分区)\xc2\xa0 和\n重新分区\xc2\xa0 会产生大致相等大小的分区。

    \n
  • \n
  • 合并我们可以减少分区,但是我们可以使用修复来增加和减少分区。

    \n
  • \n
\n


归档时间:

查看次数:

150961 次

最近记录:

7 年 前