Apache Spark CSV分区规则及无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 byrepartition(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

