Spark Small Files Problem Kills Reads — Compact It
Thousands of tiny files turn every read into a listing storm: tasks outnumber bytes 100-to-1.
20+ years shipping production backend systems. Written from production experience, not tutorials.
- ✓Spark reads/writes on partitioned Parquet or Delta tables
- ✓Basic partitioning concepts (partition columns and directories)
- ✓Familiarity with Spark UI task counts and durations
- The small-files problem means thousands of KB-sized files where hundreds of 128MB+ files belong — the driver lists, plans, and schedules more than it reads
- Diagnose fast: count files per partition and divide bytes by files — averages under ~64MB with task counts dwarfing data size confirm it
- Fix on write with coalesce/repartition sized from bytes (plus maxRecordsPerFile), and fix history with Delta OPTIMIZE or Auto Compaction
- Prevent it with partition sizing math (narrow date-partitions breed runts), right-sized shuffle partitions, and a file-count gate in CI
Picture a warehouse where every paperclip ships in its own truck. Trucks (tasks) outnumber cargo 100-to-1, the dock (driver) logs trucks instead of moving goods, and the toll booth (cloud billing per request) charges per truck. That's the small-files problem: thousands of tiny files make every read pay listing and planning overhead that dwarfs the data. Fix it with packaging discipline — fill each truck (128MB+ files) before dispatch, and consolidate small boxes on the shelves (compaction).
Nobody notices the small-files problem on day one. A streaming job writes 5,000 files of 200KB each per day — reads stay fast, dashboards stay green, everyone's happy. Eighteen months later the table holds 2.7 million files, a simple COUNT(*) takes 40 minutes, the driver OOMs during planning, and cloud bills show millions of LIST and GET requests nobody can explain. The data didn't grow that much. The file count did.
Every file costs a fixed tax regardless of size: driver-side listing and planning, one task (minimum) per file, JVM object overhead per split, and a storage API call per read. A 200KB file pays nearly the same tax as a 200MB file for 1/1000th of the payload. Multiply by millions and the tax exceeds the work — tasks spend more time starting than reading, the driver holds millions of block locations in memory, and S3-style stores throttle your LIST requests.
This article makes file sizing a first-class design input. You'll learn the partition-sizing math that prevents runts, the write-time controls (coalesce vs repartition, maxRecordsPerFile) that land files at 128MB+, the compaction tools (OPTIMIZE, Auto Compaction) that fix history, and the CI gates that keep file counts from silently compounding again.
The Fixed Tax: Why 1,000 Tiny Files Lose to One Big One
Reading a file costs a fixed toll plus a per-byte fare — and the toll dominates below ~64MB. The driver lists every path (one RPC per thousand paths on object stores, throttled past burst quotas), builds one split per file (each split a JVM object the planner must hold — millions of splits OOM even generous drivers), and schedules at least one task per file (each task paying ~1s launch tax plus serialization). A 200KB file and a 200MB file pay nearly identical tolls; the small one just carries 1/1000th of the freight.
Object stores sharpen the pain because they bill per request. Our incident's 480M monthly GET/LIST calls cost $3,100 against $94 of stored bytes — the toll exceeded the freight 33-to-1. Throttling compounds it: burst past the store's request quota and reads queue into retry storms, which look exactly like 'slow Spark' in task metrics while the real queue sits in the storage client.
The target falls out of the arithmetic: 128MB-1GB per file keeps tolls under ~1% of total read cost on modern stores, fits Parquet row-group efficiency (128MB row groups mean one file ≈ one row group ≈ one efficient read unit), and keeps driver planning in the thousands of splits, not millions. Below 64MB average you're paying toll; above 1GB you start hurting parallelism and predicate-pushdown granularity. Size for the toll, not the bytes.
Partition Sizing Math: Files per Partition, Not per Table
Averages lie at table level: 4.1 TB across 2.7M files is 190KB average, but the real design unit is the partition — the directory a query actually lists and the task set a write actually fills. Math per partition: files_you_need = partition_bytes / 128MB. A 2GB daily partition needs ~16 files; a 50MB hourly partition needs one file (or a coarser partition grain entirely). Any partition holding 10x its needed count is manufacturing toll.
Narrow partitioning is the runt factory. Hourly plus high-cardinality sub-partitions (hour × country × device = thousands of directories) divide each batch's bytes across thousands of buckets; with 200 default write tasks, most buckets receive kilobytes — each minted as a file. Two fixes compose: coarsen the grain (daily instead of hourly when queries tolerate it — fewer, fuller buckets) and size write tasks to the bucket (n from bytes-per-partition, not the global default).
Verify with a per-partition histogram, not a table average: count files and sum bytes per partition directory, then flag partitions where files exceed 3x the 128MB-ideal. The offenders cluster — usually the narrowest grains and the quietest hours — which tells you exactly which writers to resize and which grains to coarsen, instead of re-partitioning the world.
Write-Time Control: Coalesce, Repartition, maxRecordsPerFile
The cheapest file is the one never minted small: size the write before it lands. Repartition(n) fully reshuffles to exactly n files per partition — precise, auditable, and expensive (a full shuffle), so reserve it for pipeline boundaries where balanced output protects every downstream reader. Coalesce(n) only narrows without a shuffle — cheap, runs inside existing tasks, but can't widen and can leave lopsided files when inputs skew; perfect for the common case of '200 default tasks, 8 files needed.'
Compute n from bytes, not vibes: n = ceil(partition_bytes / 128MB), clamped to at least 1 and bounded above by current parallelism (coalesce can't exceed it). For partitioned writes, apply per-partition logic — repartition(n, partition_cols...) distributes by the partition columns first, so each directory gets its own n files rather than n files globally. Streaming writers add trigger.AvailableNow with foreachBatch coalescing (each micro-batch sized before commit) or Delta optimized writes, which merge small outputs automatically.
maxRecordsPerFile (spark.sql.files.maxRecordsPerFile) is the backstop, not the strategy: it caps runaway giants when rows vary wildly in size, but record counts don't track bytes (wide rows vs narrow rows differ 100x), so it can't replace byte math. Set it generously (millions) to catch pathology while n-from-bytes does the real sizing — and assert output file counts in tests, because writers silently revert to defaults on every copied template.
Fixing History: OPTIMIZE, ZORDER, and Auto Compaction
History doesn't need rewriting — it needs compacting. Delta Lake's OPTIMIZE collapses small files within (optionally: specific) partitions into 128MB+ files in place, with no logic changes and full time-travel preserved. Add ZORDER BY on frequent filter columns (event time, device_id) and the same pass also clusters data for partition-pruning-plus-skipping — our incident's OPTIMIZE cut both file count (2.7M → 31K) and filtered-read I/O, since ZORDER-colocated rows let readers skip whole files.
Scope the run to bound it: OPTIMIZE ... WHERE dt >= '2026-01-01' compacts recent hot partitions first (where new runts also land), then march backward through history in bounded jobs. Each run is idempotent and metrics-visible (files added vs removed in the operation log), so schedule it like any batch: weekly for active tables, monthly for archives. On Databricks, Auto Compaction (delta.autoCompact.enabled... via table property databricks.delta.autoCompact.enabled on older runtimes, now delta.autoCompact) merges runts automatically after writes — the mint gets swept continuously instead of quarterly.
Open-format alternatives compose the same way: Apache Hudi's clustering, Iceberg's rewrite_data_files procedure, and plain repartition-and-overwrite for Parquet directories all convert runts to 128MB+ in place. The principle never changes — compact by partition, verify counts after, bank before/after query times — whatever the engine's spelling.
Streaming and Micro-Batches: Stop Minting Runts
Streaming writers are runt factories by design: a micro-batch every 2 minutes across 200 tasks mints 200 files per trigger — 144,000 files a day before fan-out. Three controls tame it. Trigger sizing: trigger.AvailableNow (or larger micro-batch intervals) processes more data per commit, so each commit's files run fuller. foreachBatch coalescing: size each micro-batch's output with the same byte math before writing — the batch function is just batch code with n from that batch's bytes. Delta optimized writes plus Auto Compaction: the engine buffers small outputs and merges them post-commit, converting the factory into a compactor.
Partition grain matters more in streaming than anywhere else: per-minute or per-hour directories fragment every micro-batch's bytes across thousands of buckets. Prefer daily grains with ZORDER on event time (same pruning power, 24x fuller files) unless hourly SLAs genuinely require hourly directories — and when they do, budget the file count explicitly (triggers-per-day × tasks-per-trigger ÷ compaction ratio) so the mint is a line item, not a surprise.
Checkpoint and commit hygiene completes the picture: idempotent foreachBatch writers (merge on keys, deterministic file sizing) survive restarts without duplicating runts, and exactly-once Delta sinks mean replays don't double-mint. A streaming table without compaction scheduling is a countdown — set the sweep cadence at creation, not at the 2.7M-file postmortem.
Gates and Habits: Keeping Counts Low Forever
File discipline decays because every new job copies yesterday's template with its 200-partition default and hourly grain. Gates make decay visible: a daily metric of file count and average size per table (paged at 2x the 128MB-ideal count), a CI assertion that test writes produce files above ~32MB, and partition-design review for new tables that budgets files-per-day from triggers × tasks × fan-out before the first byte lands.
Review habits are cheap and load-bearing. New tables declare grain (daily default, hourly only with an SLA reason), write sizing (n-from-bytes with the measurement query linked), and sweep cadence (Auto Compaction on, or OPTIMIZE schedule named) in the design doc — three lines that prevent the entire incident class. Quarterly, re-run the census on the top 20 tables by file count: ingestion drifts (new sub-partitions, quieter hours, copied defaults) restart mints silently, and the census catches them at 2x instead of at 2.7M.
Cost attribution seals it: bill LIST/GET request spend per table alongside storage bytes so owners see their toll ratio. The team paying $3,100 in requests for $94 of storage compacted within a week of seeing that ratio; teams shown only terabytes never feel the urgency. Make the toll visible and the packaging discipline follows.
2.7M Files, 40-Minute COUNT(*): The Streaming Table That Ate Itself
- Graph file counts, not just bytes: 4.1 TB looked healthy while 2.7M files killed reads. Average-bytes-per-file per partition is the metric that predicts this failure 12 months early.
- Narrow partitions multiply the default-partition-count mistake: 200 shuffle partitions × 96 daily sub-partitions minted ~19,200 runts/day. Size write tasks from bytes-per-partition, not from global defaults.
- Access overhead can exceed storage 33-to-1 ($3,100 vs $94), so compaction pays for itself in weeks — one OPTIMIZE weekend plus write-time coalescing retired the whole failure class.
| File | Command / Code | Purpose |
|---|---|---|
| file_census.py | from pyspark.sql import functions as F | Partition Sizing Math |
| sized_write.py | from pyspark.sql import functions as F | Write-Time Control |
| optimize_history.sql | OPTIMIZE lake.events WHERE dt >= '2026-09-01' | Fixing History |
| streaming_sized_sink.py | from pyspark.sql import functions as F | Streaming and Micro-Batches |
Key takeaways
Common mistakes to avoid
5 patternsBuying bigger drivers to fix slow reads on runt tables
Leaving 200 default write partitions against narrow grains
Compacting without ZORDER (or ZORDERing without compacting)
Treating maxRecordsPerFile as the sizing strategy
No file-count metric until the postmortem
Interview Questions on This Topic
Reads slow down over months while data volume is flat. What's your first metric?
Frequently Asked Questions
20+ years shipping production backend systems. Written from production experience, not tutorials.
That's Spark. Mark it forged?
5 min read · try the examples if you haven't