Spark本地模式groupBy后分区数始终为1?原因排查
问题原因与解决办法
你的问题核心是Spark自适应查询优化(AQE)和极小数据集的shuffle特性共同导致的,并非设置错误或理解偏差。
具体原因
- 自适应查询优化(AQE)默认开启:Spark 3.x及以上版本默认启用
spark.sql.adaptive.enabled=true,该优化会在shuffle阶段后自动合并数据量极小的分区。你的groupBy结果仅2条数据(NY和CA各一组),所有shuffle分区的数据量远低于合并阈值,最终被合并成1个分区。 - 数据集key数量过少:按
state分组后仅生成2个分组,即使设置了spark.sql.shuffle.partitions=500,大部分shuffle分区都是空的,AQE会直接将这些空分区和小数据分区合并为一个。
验证与解决
方案1:关闭自适应查询优化
在创建SparkSession时添加参数,强制保留设置的shuffle分区数:
val spark = SparkSession.builder() .appName("Test App") .master("local[5]") .config("spark.sql.adaptive.enabled", "false") // 关闭AQE .config("spark.sql.shuffle.partitions", 500) // 设置shuffle分区数 .getOrCreate()
此时执行println(df2.rdd.getNumPartitions),输出会变为你设置的500。
方案2:增大数据集规模
当分组结果数据量足够大,每个shuffle分区数据量达到spark.sql.adaptive.coalescePartitions.minPartitionSize(默认64MB)阈值时,AQE就不会合并所有分区,分区数会接近你设置的spark.sql.shuffle.partitions值。
额外说明
- 即使关闭AQE,若分组key数量远小于shuffle分区数,大部分分区会是空的,但分区总数仍等于设置的
spark.sql.shuffle.partitions值。 - 生产环境不建议随意关闭AQE,它能有效优化小数据量查询性能,仅在需要精确控制分区数的场景下调整即可。
内容的提问来源于stack exchange,提问作者Pepria
相关产品推荐
相关产品推荐

