You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark本地模式groupBy后分区数始终为1?原因排查

问题原因与解决办法

你的问题核心是Spark自适应查询优化(AQE)和极小数据集的shuffle特性共同导致的,并非设置错误或理解偏差。

具体原因

  1. 自适应查询优化(AQE)默认开启:Spark 3.x及以上版本默认启用spark.sql.adaptive.enabled=true,该优化会在shuffle阶段后自动合并数据量极小的分区。你的groupBy结果仅2条数据(NY和CA各一组),所有shuffle分区的数据量远低于合并阈值,最终被合并成1个分区。
  2. 数据集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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 08:43:26