Spark DataFrame仅1个分区原因及并行度配置咨询
问题分析与解答
为什么仅生成1个分区?
你的代码中Spark只创建1个分区的核心原因在于CSV文件的大小和Spark的文件读取策略:
- Spark读取文件时,依赖Hadoop的
FileInputFormat来划分数据分片(split)。如果你的CSV文件体积远小于默认的HDFS块大小(通常128MB),Hadoop只会生成1个数据分片,对应Spark的1个分区。 local[4]只是设置了本地模式下允许并行运行的最大任务数,它不会主动拆分数据来生成更多分区——分区数量由输入数据的分片逻辑决定,而非集群的核心数。
另外,你的代码中没有指定CSV文件的表头(header=true),导致列名被自动命名为_c0、_c5等,但这不是分区数少的原因。
如何让Spark利用4核并行执行?
要让count操作使用4个Task,你需要主动调整DataFrame的分区数,常见方式有两种:
1. 读取文件时指定分区相关参数
通过设置spark.sql.files.maxPartitionBytes(默认128MB)为更小的值,让Spark将单个文件拆分为多个分片:
from pyspark.sql import SparkSession # 调整单个分区的最大字节数为32MB spark = SparkSession.builder.master("local[4]") \ .config("spark.sql.files.maxPartitionBytes", "32MB") \ .getOrCreate() df = spark.read.csv("annual-enterprise-survey-2021-financial-year-provisional-size-bands-csv.csv") df.createOrReplaceTempView("table") sqldf = spark.sql('SELECT _c5 FROM table WHERE _c5 > "1000"') print(sqldf.count()) print(df.rdd.getNumPartitions()) # 此时分区数会根据文件大小和设置的参数调整 print(sqldf.rdd.getNumPartitions())
2. 显式重分区
使用repartition()方法强制将DataFrame划分为指定数量的分区:
df = spark.read.csv("annual-enterprise-survey-2021-financial-year-provisional-size-bands-csv.csv").repartition(4) # 后续操作基于这个4分区的DataFrame
注意:如果文件本身很小,强制重分区可能会带来额外的shuffle开销,反而降低性能。只有当数据量足够大时,并行处理才有意义。
Action执行时DataFrame分区数的确定逻辑
DataFrame的分区数由以下几个环节共同决定:
- 输入源阶段:
- 对于文件类数据源(CSV、Parquet等):分区数 = 总文件大小 /
spark.sql.files.maxPartitionBytes(向上取整),同时每个文件至少对应1个分区。如果有多个小文件,Spark可能会合并部分分片(由spark.sql.files.openCostInBytes控制)。 - 对于数据库等结构化数据源:分区数由JDBC连接的
numPartitions参数决定。
- 对于文件类数据源(CSV、Parquet等):分区数 = 总文件大小 /
- 转换阶段:
- 过滤(
filter)、投影(select)等窄依赖转换会保留原分区数。 - 宽依赖转换(如
groupBy、join)会触发shuffle,此时分区数由spark.sql.shuffle.partitions(默认200)控制,也可通过repartition()/coalesce()手动调整。
- 过滤(
- Action阶段:Action操作的Task数等于当前DataFrame的分区数,Spark会尽可能并行执行这些Task(受限于集群的核心数或
local[N]中的N值)。
内容的提问来源于stack exchange,提问作者taci
相关产品推荐
相关产品推荐

