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

Apache Spark CSV分区规则及无Shuffle调整小文件分区方法咨询

How to Match Partition Count for Small CSV Files in Spark Without Shuffle

Let’s break down your problem step by step—starting with why Spark partitions your files the way it does, then how to fix the partition count without shuffling, and finally explaining that repartition behavior you observed.

Why the Partition Counts Differ

When Spark reads CSV (or any non-columnar text-based files), it relies on Hadoop’s TextInputFormat under the hood to split data into partitions. The default logic for partition count depends on two core factors:

  • spark.sql.files.maxPartitionBytes: The maximum size of each partition (default is 128MB). Spark splits files into chunks no larger than this value.
  • Total size of your input file(s)

Your smaller CSV (~487k rows) totals roughly 3 chunks of the default max partition size, hence 3 partitions. The larger CSV (~1.39M rows) fills 8 such chunks—matching your cluster’s total 8 slots—so it gets 8 partitions.

Getting 8 Partitions for the Small CSV Without Shuffle

If you want to avoid shuffling entirely, adjust the partitioning logic during the read operation instead of modifying the DataFrame afterward. Here are two reliable methods:

1. Set minPartitions During Read

Explicitly hint to Spark to create at least 8 partitions when loading the CSV by adding the minPartitions option. This works well if your input is split across multiple files:

csv_file = "/FileStore/tables/mnt/training/departuredelays02.csv"
schema = "`date` STRING, `delay` INT, `distance` INT, `origin` STRING, `destination` STRING"
df = (spark
     .read
     .format("csv")
     .option("header","true")
     .option("minPartitions", 8)  # Force minimum 8 partitions
     .schema(schema)
     .load(csv_file)
)
print(df.rdd.getNumPartitions())  # Should return 8

2. Reduce maxPartitionBytes

If your small CSV is a single file, lower the maxPartitionBytes value to force Spark to split the file into more partitions. Calculate a value roughly equal to total_file_size / 8 (e.g., if your file is 300MB, set it to ~37MB):

# Adjust the partition size to split the file into 8 chunks
spark.conf.set("spark.sql.files.maxPartitionBytes", "37000000")  # 37MB

csv_file = "/FileStore/tables/mnt/training/departuredelays02.csv"
schema = "`date` STRING, `delay` INT, `distance` INT, `origin` STRING, `destination` STRING"
df = (spark
     .read
     .format("csv")
     .option("header","true")
     .schema(schema)
     .load(csv_file)
)
print(df.rdd.getNumPartitions())  # Should return 8

This splits the file into smaller chunks during the read phase, no shuffle required.

Explaining the Repartition + Count Behavior

When you ran df = df.repartition(8) followed by df.count(), here’s why you saw 3 tasks in the first stage:

  • The first stage (Stage 123: 3/3) is Spark reading the original 3 partitions of the small CSV.
  • The second stage (Stage 125: 8/8) is the shuffle triggered by repartition(8)—Spark redistributes data from 3 partitions into 8 new ones.
  • The final stage (Stage 128: 1/1) is the count action, which aggregates results from the 8 shuffled partitions.

Repartition always triggers a shuffle because it needs to redistribute data across the cluster to reach the desired partition count. That’s why read-time parameter adjustments are better if you want to avoid shuffling.

内容的提问来源于stack exchange,提问作者Sunil Thapa

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 07:50:13