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不需要移动其原始数据.
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可以使用相同大小的分区.
如果您想了解更多详情,请阅读此博客文章.
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减少分区数量要好。
基本上,重新分区允许您增加或减少分区的数量。重新分区会重新分配所有分区中的数据,这会导致完全洗牌,这是非常昂贵的操作。
Coalesce 是 Repartition 的优化版本,您只能减少分区数量。由于我们只能减少分区的数量,因此它的作用是将一些分区合并为单个分区。与重新分区相比,通过合并分区,跨分区的数据移动量更低。因此,Coalesce 是最小数据移动,但说 Coalesce 不进行数据移动是完全错误的说法。
另一件事是通过提供分区数量进行重新分区,它尝试在所有分区上均匀地重新分配数据,而在合并的情况下,在某些情况下我们仍然可能会出现倾斜数据。
我想补充贾斯汀和鲍尔的回答——
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)
合并\xc2\xa0 使用现有分区来最大程度地减少洗牌的数据量。\xc2\xa0Repartition\xc2\xa0 创建新分区\xc2\xa0 并\xc2\xa0 执行完整\n洗牌。
\n合并\xc2\xa0 会产生具有不同数据量的分区\n(有时是具有许多不同大小的分区)\xc2\xa0 和\n重新分区\xc2\xa0 会产生大致相等大小的分区。
\n合并我们可以减少分区,但是我们可以使用修复来增加和减少分区。
\n| 归档时间: |
|
| 查看次数: |
150961 次 |
| 最近记录: |