【发布时间】:2017-03-17 21:57:09
【问题描述】:
DataFrame repartition() 和 DataFrameWriter partitionBy() 方法有什么区别?
我希望两者都用于“基于数据框列对数据进行分区”?或者有什么不同?
【问题讨论】:
-
对于任何提出这个问题的人,this one 也可能是相关的
标签: apache-spark-sql data-partitioning
DataFrame repartition() 和 DataFrameWriter partitionBy() 方法有什么区别?
我希望两者都用于“基于数据框列对数据进行分区”?或者有什么不同?
【问题讨论】:
标签: apache-spark-sql data-partitioning
repartition()用于对内存中的数据进行分区,partitionBy用于对磁盘上的数据进行分区。它们经常结合使用。
repartition()和partitionBy都可以用来“根据dataframe列对数据进行分区”,但是repartition()对内存中的数据进行分区,partitionBy对磁盘上的数据进行分区。
重新分区()
让我们尝试一些代码来更好地理解分区。假设您有以下 CSV 数据。
first_name,last_name,country
Ernesto,Guevara,Argentina
Vladimir,Putin,Russia
Maria,Sharapova,Russia
Bruce,Lee,China
Jack,Ma,China
df.repartition(col("country")) 将按国家/地区对内存中的数据进行重新分区。
让我们写出数据,这样我们就可以检查每个内存分区的内容。
val outputPath = new java.io.File("./tmp/partitioned_by_country/").getCanonicalPath
df.repartition(col("country"))
.write
.csv(outputPath)
以下是数据在磁盘上的写入方式:
partitioned_by_country/
part-00002-95acd280-42dc-457e-ad4f-c6c73be6226f-c000.csv
part-00044-95acd280-42dc-457e-ad4f-c6c73be6226f-c000.csv
part-00059-95acd280-42dc-457e-ad4f-c6c73be6226f-c000.csv
每个文件都包含一个国家/地区的数据 - part-00059-95acd280-42dc-457e-ad4f-c6c73be6226f-c000.csv 文件包含以下中国数据,例如:
Bruce,Lee,China
Jack,Ma,China
partitionBy()
让我们使用partitionBy 将数据写入磁盘,看看文件系统输出有何不同。
这是将数据写入磁盘分区的代码。
val outputPath = new java.io.File("./tmp/partitionedBy_disk/").getCanonicalPath
df
.write
.partitionBy("country")
.csv(outputPath)
磁盘上的数据如下所示:
partitionedBy_disk/
country=Argentina/
part-00000-906f845c-ecdc-4b37-a13d-099c211527b4.c000.csv
country=China/
part-00000-906f845c-ecdc-4b37-a13d-099c211527b4.c000
country=Russia/
part-00000-906f845c-ecdc-4b37-a13d-099c211527b4.c000
为什么要在磁盘上分区数据?
对磁盘上的数据进行分区可以使某些查询运行得更快。
【讨论】:
如果您运行repartition(COL),您会在计算期间更改分区 - 您将获得spark.sql.shuffle.partitions(默认值:200)分区。如果您随后调用.write,您将获得一个包含许多文件的目录。
如果您运行.write.partitionBy(COL),那么您将获得与 COL 中唯一值一样多的目录。这加快了进一步的数据读取(如果您按分区列进行过滤)并节省一些存储空间(从数据文件中删除了分区列)。
更新:见@conradlee 的回答。他不仅详细解释了应用不同方法后目录结构的外观,还详细解释了两种情况下产生的文件数量。
【讨论】:
注意:我认为公认的答案不太正确!很高兴您提出这个问题,因为这些名称相似的函数的行为在重要和意想不到的方面有所不同,而官方 spark 文档中没有详细记录。
已接受答案的第一部分是正确的:调用df.repartition(COL, numPartitions=k) 将使用基于散列的分区器创建具有k 分区的数据帧。 COL 这里定义了分区键——它可以是单个列或列列表。基于散列的分区器获取每个输入行的分区键,通过partition = hash(partitionKey) % k 将其散列到k 分区的空间中。这保证了具有相同分区键的所有行最终都在同一个分区中。但是,来自多个分区键的行也可能最终位于同一个分区中(当分区键之间发生哈希冲突时)并且某些分区可能是空的。
总之,df.repartition(COL, numPartitions=k) 的不直观方面是
k 分区可能为空,而其他可能包含来自多个分区键的行df.write.partitionBy 的行为完全不同,这是许多用户无法预料的。假设您希望对输出文件进行日期分区,并且您的数据跨越 7 天。我们还假设df 开始时有 10 个分区。当您运行df.write.partitionBy('day') 时,您应该期待多少个输出文件?答案是“视情况而定”。如果df 中起始分区的每个分区都包含每天的数据,那么答案是 70。如果df 中的每个起始分区恰好包含一天的数据,那么答案就是 10。
我们如何解释这种行为?当您运行df.write 时,df 中的每个原始分区都是独立写入的。也就是说,您原来的 10 个分区中的每一个都在 'day' 列上分别进行了子分区,并为每个子分区写入了一个单独的文件。
我觉得这种行为很烦人,希望有办法在编写数据帧时进行全局重新分区。
【讨论】:
df.write().repartition(COL).partitionBy(COL) 这样的东西?我的目标是partitionBy() 行为,但文件大小和文件数量与我最初的大致相同。这很容易实现吗? partitionBy(date) => 70 个文件示例是相关的。我想要大约 10 个文件,每天一个,对于有更多数据的日子可能需要 2 或 3 个。
df.repartition(df.date, 1000)。许多人期望每个分区恰好包含一天的数据。但是,这 1000 个分区中的一些将是空的,而其他分区将包含多天的数据。许多人觉得这不直观(也许你不这么认为,因此会感到困惑)。