Spark中shuffle partition与repartition的区别是什么?
先把两个概念的本质说清楚
你首先混淆了两个完全不在同一维度的东西,它们根本不是同类型的分区调整手段:
- shuffle partition:一般指Spark配置项
spark.sql.shuffle.partitions,是个全局生效的默认参数,默认值200,作用是给所有没有手动指定输出分区数的shuffle操作规定shuffle阶段的输出分区数。它本身不是用来主动调整分区的算子,只是个默认规则。 repartition(partitionNum):是Spark提供给用户主动调整数据集分区数的算子,不管你是要增加还是减少分区,调用这个算子一定会触发全量数据shuffle,所有数据会被打乱重新均匀分布到你指定数量的分区里,分区数完全由你传入的参数决定,不受前面说的shuffle partition配置影响。
关于减少分区的效果解答
- 对
repartition来说,当然可以实现减少分区的效果:只要你传入的分区数小于当前数据集的分区数就行,比如当前数据集有500个分区,调用df.repartition(20)就会把全量数据shuffle后重分布到20个分区。但要注意:如果你的需求只是减少分区,优先用coalesce(),这个算子在减分区场景下不会触发全量shuffle,只会合并相邻的现有分区,性能比repartition好很多。 - 对shuffle partition配置来说,它本身没有主动调整分区的能力:它不会主动修改现有数据集的分区数,只有当你执行join、groupBy、distinct这类会触发shuffle的算子,而且没手动给算子指定输出分区数的时候,shuffle后的分区数才会取这个配置的值。如果这个配置的值比shuffle前的分区数小,那shuffle完成后分区数确实会变少,但这只是默认规则生效的结果,不是它主动“减少分区”的作用。如果你全程没有触发任何shuffle操作,哪怕把这个配置改得再小,数据集的分区数也不会有任何变化。
常见误区提醒
别信网上说的“改spark.sql.shuffle.partitions就能全局调整Spark任务分区数”的错误说法,这个参数只管shuffle阶段的输出,初始读入数据的分区数、不触发shuffle的算子转换后的分区数,都和这个参数没关系。
给你一段测试代码跑一遍就懂了:
# 读取数据源,初始分区数和文件切块数相关,这里假设读入后是200分区 df = spark.read.parquet("你的测试数据路径") print(f"初始分区数: {df.rdd.getNumPartitions()}") # 输出200 # 修改shuffle分区默认配置为10 spark.conf.set("spark.sql.shuffle.partitions", 10) # 执行不触发shuffle的过滤操作 df_filter = df.filter("age > 18") print(f"过滤后(无shuffle)分区数: {df_filter.rdd.getNumPartitions()}") # 还是输出200,配置不生效 # 执行触发shuffle的分组操作,没手动指定分区数 df_group = df_filter.groupBy("city").count() print(f"分组后(触发shuffle,未指定分区数)分区数: {df_group.rdd.getNumPartitions()}") # 输出10,配置生效 # 调用repartition手动指定分区数为50 df_rep = df_filter.repartition(50) print(f"repartition后分区数: {df_rep.rdd.getNumPartitions()}") # 输出50,不受shuffle分区配置影响
内容的提问来源于stack exchange,提问作者Raj harini
相关产品推荐
相关产品推荐

