Skip to content
Good engineers know 6 min read · maxRecordsPerFile 1 practice problem ↓

The small file problem

Why thousands of tiny files slow every read, how Spark jobs create them, and how to prevent and fix them.

You will learn

  • Why many small files slow down every read
  • The five ways Spark jobs usually create them
  • How to prevent them when writing
  • How to fix them after the fact, and how to detect them

Read first

Comfortable with these? Read on.

TL;DR Each file costs a listing entry, an open, a footer read and often a task, whatever its size. A table of 500,000 files of 100 KB is far slower to read than 400 files of 128 MB holding the same data. Prevent them at write time with sensible partitioning, and compact what already exists.

Why small files hurt

Reading a file has a fixed cost that does not depend on its size:

Cost per fileWhy it adds up
ListingObject stores list about 1,000 keys per request. Listing 1 million files takes 1,000 sequential requests before any planning.
Open and footer readParquet keeps its schema and statistics at the end of each file; every file needs at least one extra request.
MetadataFor Hive tables, every partition and file is tracked by the metastore and the driver. For Delta and Iceberg, every file is an entry in the log or manifests, which grows and slows planning.
TasksSpark packs small files together (counting each as at least openCostInBytes, 4 MB) but scheduling still grows with the file count.
Compression and statisticsTiny row groups compress worse, and min/max statistics over a few rows rarely let a reader skip anything useful.

On object stores like S3, ADLS and GCS, request latency, not bandwidth, dominates. Reading 10 GB as 100,000 files can take ten times longer than the same 10 GB as 80 files.

How Spark jobs create them

  1. Too many shuffle partitions at write time. A write after a join or aggregation produces one file per non-empty partition. 200 partitions writing 50 MB of data gives 200 files of 250 KB.
  2. partitionBy on a high-cardinality column. Each task writes one file per partition value it holds. 200 tasks × 365 dates = up to 73,000 files per run.
  3. Streaming micro-batches. A stream that commits every 10 seconds writes files every 10 seconds: 8,640 batches a day, each with several files.
  4. Frequent small appends. Hourly or per-event loads of a few MB each.
  5. Over-partitioned sources. Reading thousands of files and writing straight back preserves the count.

Walk through one bad write

Step 1 · the job

A daily job aggregates 5 GB of events and writes them partitioned by event_date and country (60 countries). The upstream shuffle has 200 partitions.

The command

daily.write.mode("append").partitionBy("event_date", "country").parquet(path)

Step 2 · the math

Each of the 200 tasks holds rows for most of the 60 countries, and writes one file per country it holds. That is up to 200 × 60 = 12,000 files for one day, averaging about 400 KB.

Step 3 · a year later

365 days × 12,000 = 4.4 million files. Queries spend minutes listing and opening files before reading meaningful data, and the driver needs gigabytes just to plan.

Step 4 · the fix

Repartition by the partition columns before writing, so each (date, country) combination is written by one task: one file per directory. 60 files per day instead of 12,000.

The command

daily.repartition("event_date", "country") \
    .write.mode("append").partitionBy("event_date", "country").parquet(path)
INSERT INTO events_by_country
SELECT /*+ REPARTITION(event_date, country) */ *
FROM daily

Preventing them when writing

TechniqueHowWatch out
Let AQE coalesceWith AQE on, a write after a shuffle uses merged partitions near 64 MBOnly after a shuffle; raise advisoryPartitionSizeInBytes to 128-256 MB for write-heavy jobs
repartition(n)Pick n from the output size: 50 GB / 256 MB ≈ 200A full shuffle
repartition(partition cols)One task per table partition, so one file per directoryA huge partition value becomes one huge task (skew)
coalesce(n)Merge without a shuffleReduces the parallelism of the whole stage
maxRecordsPerFile.option("maxRecordsPerFile", 5_000_000) caps file size from aboveDoes not merge small files
Coarser partitioningPartition by month instead of day, or not at all for small tablesSee Partitioning done right

Delta Lake can do this for you: optimized writes add an adaptive shuffle before writing so each partition gets fewer, larger files, and auto compaction runs a small compaction after a write. They are table properties (delta.autoOptimize.optimizeWrite, delta.autoOptimize.autoCompact) on Databricks and in recent open-source Delta releases; check your version.

Fixing them afterwards

PySparkSpark SQL · Compaction
# Plain Parquet: rewrite one partition into fewer files, then swap it in
(spark.read.parquet(f"{path}/event_date=2025-03-01")
    .repartition(8)
    .write.mode("overwrite").parquet(f"{tmp}/event_date=2025-03-01"))

# Delta Lake
from delta.tables import DeltaTable
DeltaTable.forPath(spark, path).optimize().executeCompaction()
-- Delta Lake
OPTIMIZE events WHERE event_date >= '2025-03-01';

-- Apache Iceberg
CALL catalog.system.rewrite_data_files(table => 'db.events');

With plain Parquet, rewriting in place is unsafe for concurrent readers: they can see the old files deleted before the new ones appear. Table formats make compaction a single atomic commit, which is one of their biggest practical benefits. Old files remain until VACUUM (Delta) or expire_snapshots (Iceberg) removes them.

Detecting the problem

  • DESCRIBE DETAIL table on Delta returns numFiles and sizeInBytes: divide to get the average file size. Under about 32 MB is worth fixing.
  • For Iceberg, query the metadata table db.events.files for file_size_in_bytes.
  • In the Spark UI, a scan with tens of thousands of tasks reading a few hundred KB each, or a long gap before the first job (planning and listing), points to small files.
Why this matters: the healthy target for analytic tables is roughly 128 MB to 1 GB per file. Larger files reduce overhead further but make each task longer and reduce parallelism for small queries.

Common mistakes

coalesce(1) to "fix" small files on big data

One task does all the work, and one giant file reduces read parallelism later.

partitionBy on user_id or timestamp

Millions of directories, each with tiny files.

Compacting plain Parquet in place while others read it

Readers can fail or see duplicates. Use a table format, or write to a new location and swap.

Key takeaways

  • Each file has a fixed cost: listing, opening, metadata and often a task.
  • Small files come from too many write partitions, high-cardinality partitionBy, streaming and small appends.
  • Repartition by the partition columns, let AQE coalesce, and use maxRecordsPerFile to cap size.
  • Table formats make compaction atomic; aim for 128 MB to 1 GB files.

Check yourself

3 questions

1. 200 tasks each hold rows for 30 dates and write with partitionBy("date"). Up to how many files are written?

Show the answer

6,000. Each task writes one file per partition value it holds: 200 × 30.

2. Which change gives one file per date directory?

Show the answer

repartition("date") before the write. Repartitioning by the partition column sends all rows of a date to one task.

3. Why is compaction safer on Delta or Iceberg than on plain Parquet?

Show the answer

The rewrite is committed atomically, so readers see either the old files or the new ones. The new file list is published in one commit; old files stay readable until vacuumed.

Practice it

Interview problems that use this: write the PySpark, run it, and get graded on hidden tests.

Solve: Find Partitions That Need Compaction →

Go deeper

Primary sources: Performance tuning · Delta Lake: optimizations · Iceberg: maintenance