When Spark reads files, it packs file chunks into partitions of about spark.sql.files.maxPartitionBytes, which is 128 MB by default. Each partition becomes one task. So the partition count is roughly the total data size divided by 128 MB, adjusted for how files are laid out.
The rules
- Big splittable files are cut into 128 MB pieces. A 1 GB Parquet file gives about 8 partitions, because Parquet can be split at row-group boundaries.
- Small files are packed together into one partition, up to the limit. Spark adds a fixed cost per file (
spark.sql.files.openCostInBytes, 4 MB by default) to account for the time it takes to open each file. So 1,000 files of 100 KB do not all fit in one partition: they pack as if each were 4 MB plus a bit, giving about 30 or so partitions of roughly 128 MB each. - Files that cannot be split are read by one task each. A 5 GB
.csv.gzfile is one gzip stream, so one task reads all 5 GB, and the other cores idle. This is a classic cause of a stage with one huge task.
Check it
df = spark.read.parquet("s3://bucket/orders/")
print(df.rdd.getNumPartitions())Compare with your core count. Far fewer partitions than cores means idle cores. Thousands of tiny partitions means too much scheduling overhead and slow listing.
What you can do
- Lower
maxPartitionBytesif you want more parallelism for heavy per-row work, such as complex parsing. - Raise it if you have many cores but tasks are very short.
- For non-splittable formats, convert to Parquet, or ask the source to write multiple smaller files.
- Compact small files into larger ones after loading.
Partitions at read time and partitions after a shuffle are controlled by different settings, so keep them apart when tuning.