Spark按列repartition未生成默认200分区,仅返回1个分区原因
问题原因与排查方案
你遇到的这个现象,核心原因大概率是以下几种情况之一:
1. 测试数据的列基数异常
如果你的department列(以及所有测试的列)所有行的值完全相同(比如全为null、全是同一个固定值),哈希分区后所有数据会被分配到同一个分区。但这里要注意:即使所有数据都在一个分区,getNumPartitions()理论上应该返回spark.sql.shuffle.partitions配置的200(仅199个分区为空)。若你得到的结果是1,可先验证数据:
- 执行
df.select("department").distinct().count(),查看该列的唯一值数量。如果结果为1,说明数据本身存在问题,换有足够基数的列测试即可。
2. spark.sql.shuffle.partitions配置未在shuffle时生效
虽然你通过spark.conf.get查到配置值是200,但如果该配置是在创建SparkSession之后动态设置,且测试DataFrame是在配置修改前创建的,那么针对这个DataFrame的shuffle操作可能不会使用新的配置值。
- 验证方式:在创建SparkSession时直接指定该配置,示例代码(Python):
重新创建测试DataFrame并执行from pyspark.sql import SparkSession spark = SparkSession.builder \ .master("local[*]") \ .config("spark.sql.shuffle.partitions", "200") \ .getOrCreate()repartition("department"),再查看分区数。
3. Spark版本或API重载的细节差异
在部分旧版本Spark中,repartition(cols: Column*)的默认分区数逻辑与官方文档描述不符。而显式指定分区数的重载方法repartition(numPartitions: Int, cols: Column*)不受此影响,这也是你指定200分区时能正常工作的原因。
- 若使用Scala API,确认未混淆参数顺序;若使用Python API,确认传入的是正确的列名参数。
额外验证步骤
执行以下代码,直观查看分区情况:
repartitioned_df = df.repartition("department") print(f"分区数: {repartitioned_df.rdd.getNumPartitions()}") # 查看每个分区的行数 print(repartitioned_df.rdd.glom().map(len).collect())
如果输出分区数为1且行数等于原DataFrame总条数,说明所有数据被分配到了一个分区;如果分区数为200但仅一个分区有数据,说明是数据基数问题。
内容的提问来源于stack exchange,提问作者Hawii Hawii
相关产品推荐
相关产品推荐

