Home › Data Engineering › Spark Small Files Problem Kills Reads — Compact It
Intermediate 5 min · September 23, 2026
Spark Small Files Problem Destroys Read Performance

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.

N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Written from production experience, not tutorials.

Follow
✓ Production
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
Before you start⏱ 13 min
  • ✓Spark reads/writes on partitioned Parquet or Delta tables
  • ✓Basic partitioning concepts (partition columns and directories)
  • ✓Familiarity with Spark UI task counts and durations
 ● Production Incident 🔎 Debug Guide
⚡Quick Answer
  • 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
✦ Definition~90s read
What is Spark Small Files Problem Destroys Read Performance?

The small-files problem is fixed-cost economics applied to storage: every file charges listing (driver enumerates paths), planning (one split object per file in driver memory), scheduling (minimum one task per file, each with ~1s launch tax), and requests (per-GET billing and throttling on object stores). These tolls don't scale with size — a 200KB file pays the full toll for 1/1000th of a 200MB file's freight.

★
Picture a warehouse where every paperclip ships in its own truck.

Past roughly a million files per table (or under ~64MB average), tolls exceed work: tasks start longer than they read, planners OOM before executing, and request bills eclipse storage bills.

Writers mint runts whenever task count exceeds data shape: 200 default partitions across 96 hourly sub-partitions, 2-minute streaming triggers across 200 tasks, dynamic overwrites touching thousands of directories. Each combo of task × partition × trigger becomes a file, and small combos become KB files.

Time compounds it linearly — 19,200 runts a day is 7M a year — while reads degrade super-linearly as listing, planning, and throttling interact.

Compaction is the economy of scale that reverses it: merging 87 files of ~2MB into one 128MB+ file pays one toll instead of 87 for identical bytes. Write-time sizing prevents the mint; OPTIMIZE-class tools consolidate history; Auto Compaction sweeps continuously. All three enforce the same unit economics — full trucks or don't dispatch.

Plain-English First

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.

📊 Production Insight
The $3,100-vs-$94 bill ratio ended every 'data growth' argument in one slide — tolls, not terabytes, owned the budget, and compaction paid back in under a month.
🎯 Key Takeaway
Every file pays listing, planning, scheduling, and request tolls regardless of size — land files at 128MB-1GB so the toll stays under ~1% of read cost.

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.

file_census.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
from pyspark.sql import functions as F

# Per-partition file census: the metric that predicts this failure a year early.
# Run against the table path; flag partitions above 3x their 128MB-ideal count.
files = spark.read.format("binaryFile").load("s3://lake/events/")
# path like .../dt=2026-09-22/hr=03/part-00000-....parquet
census = (
    files.withColumn("partition", F.regexp_extract("path", r"((dt=[^/]+/)(hr=[^/]+/)?)", 1))
    .groupBy("partition")
    .agg(F.count("*").alias("files"),
         (F.sum("length") / 1024 / 1024).alias("mb"))
    .withColumn("ideal_files", F.ceil(F.col("mb") / 128))
    .withColumn("runt_ratio", F.col("files") / F.greatest(F.col("ideal_files"), F.lit(1)))
    .orderBy(F.desc("runt_ratio"))
)
census.filter("runt_ratio > 3").show(20, truncate=False)
print("Worst offender ratio:", census.first()["runt_ratio"])
📊 Production Insight
The census showed quiet-hour partitions at 400x ideal (3 files needed, 1,200 minted) while busy hours sat near 2x — resizing only the quiet-hour writer fixed 80% of the mint with one config change.
🎯 Key Takeaway
Math per partition (bytes/128MB = files needed), coarsen narrow grains, and let the runt-ratio histogram name the exact writers to resize.

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.

sized_write.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
from pyspark.sql import functions as F
import math

src = spark.read.parquet("s3://lake/staging/dt=2026-09-22/")

# Size n from BYTES per output partition, not from global defaults.
# Quiet hour: 900MB -> 8 files; busy hour logic identical, n differs.
PARTITION_MB = 900  # measure per partition (see file_census.py), or estimate
n = max(1, math.ceil(PARTITION_MB / 128))
print(f"Writing {n} files for ~{PARTITION_MB}MB")

# Narrowing only (200 tasks -> n files): coalesce, no shuffle, cheap.
(src.coalesce(n)
    .write.mode("overwrite")
    .parquet("s3://lake/events/dt=2026-09-22/hr=03/"))

# At pipeline boundaries needing exact balance: repartition by key + count.
# (src.repartition(n, "device_id")
#     .write.partitionBy("dt")
#     .mode("overwrite").parquet("s3://lake/events_balanced/"))

# Backstop against pathological giants (record counts != bytes — generous cap).
# spark.conf.set("spark.sql.files.maxRecordsPerFile", "5000000")
📊 Production Insight
Switching the hourly writer from default-200 to byte-sized coalesce cut daily minted files 19,200 → 640 with zero shuffle cost — the mint stopped the same day the config shipped.
🎯 Key Takeaway
Coalesce (cheap, narrowing) for daily writes, repartition (exact, shuffled) at boundaries, n from bytes/128MB — and maxRecordsPerFile only as a backstop.

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.

optimize_history.sqlSQL
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
-- Compact history in bounded, idempotent passes — recent partitions first.
OPTIMIZE lake.events WHERE dt >= '2026-09-01'
  ZORDER BY (event_time, device_id);

-- Then march backward through history one month per job.
-- OPTIMIZE lake.events WHERE dt >= '2026-08-01' AND dt < '2026-09-01';

-- Verify: file counts and p50 sizes from the operation + table history.
DESCRIBE HISTORY lake.events LIMIT 5;

-- Continuous sweeping so runts never re-accumulate (Databricks):
ALTER TABLE lake.events SET TBLPROPERTIES (
  'delta.autoCompact.enabled' = 'true',
  'delta.autoCompact.maxFileSize' = '134217728'
);
💡Compact recent partitions first, then march backward
New runts land in recent partitions while old ones sit cold — OPTIMIZE the hot tail first for immediate relief, then sweep history in bounded monthly jobs. Each pass is idempotent and logged, so interrupted compactions just resume; and ZORDER on the same pass buys skipping gains that pure compaction doesn't.
📊 Production Insight
One weekend OPTIMIZE plus ZORDER took COUNT(*) 40 min → 3 min and filtered reads down 6x — same bytes, 87x fewer files, with time-travel intact and zero pipeline logic touched.
🎯 Key Takeaway
OPTIMIZE history in bounded partition passes (ZORDER for skipping bonus), enable Auto Compaction for the future, and verify counts after every pass.

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.

streaming_sized_sink.pyPYTHON
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
from pyspark.sql import functions as F
import math

events = (spark.readStream.format("kafka")
            .option("subscribe", "events").load()
            .select(F.from_json(F.col("value").cast("string"),
                                "device_id STRING, amount DOUBLE").alias("e"))
            .select("e.*"))

def sized_batch(batch_df, batch_id):
    # Size THIS batch's output from its own bytes: no runt factory.
    n_files = batch_df.count()  # small batches: exact ticketing below
    if n_files == 0:
        return
    approx_mb = batch_df.rdd.map(
        lambda r: len(str(r))).sum() / 1024 / 1024
    n = max(1, math.ceil(approx_mb / 128))
    (batch_df.coalesce(n).write.format("delta")
     .mode("append").save("s3://lake/events/"))

(events.writeStream
 .trigger(availableNow=True)  # fuller commits, fewer files
 .foreachBatch(sized_batch)
 .option("checkpointLocation", "s3://lake/_chk/events/")
 .start().awaitTermination())
📊 Production Insight
Moving one Kafka sink from 2-minute triggers to AvailableNow hourly batches plus foreachBatch coalescing cut its daily mint 144,000 → 96 files — a 1,500x reduction from two config lines and a ten-line batch function.
🎯 Key Takeaway
Size triggers for full commits, coalesce each micro-batch from its own bytes, and schedule compaction at table creation — streaming tables are countdowns without sweeps.

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.

📊 Production Insight
Per-table request-spend attribution moved 11 teams to self-serve compaction in one quarter — visible tolls beat every wiki page about file sizing ever written.
🎯 Key Takeaway
Gate counts in CI, metric them daily, review grain-plus-sizing-plus-sweeps per table, and bill request tolls to owners — visible tolls self-enforce.
● Production incidentPOST-MORTEMseverity: high

2.7M Files, 40-Minute COUNT(*): The Streaming Table That Ate Itself

Symptom
The events table (4.1 TB total, 2.7M files, 190KB average) degraded gradually then suddenly: partition-pruned reads that took 3 minutes in year one took 40+ minutes, full-table COUNT(*) timed out the BI tool at 60 minutes, and the driver twice OOM'd with 8G heap during query planning — before reading a single byte. Meanwhile the storage bill showed 480M GET/LIST requests per month ($3,100) against 4.1 TB stored ($94): access overhead cost 33x the bytes.
Assumption
The team blamed data growth and bought bigger drivers (8G→16G→32G) plus more executors, reasoning that 4.1 TB 'should' read in minutes on their fleet. When that barely moved the needle, they suspected Parquet encoding and recompressed a sample — identical timings. Both theories measured bytes; the bottleneck was file count, which nobody had graphed. Average file size (190KB) was never computed until week 3 because every dashboard showed terabytes, never file tallies.
Root cause
The streaming writer used the default 200 shuffle partitions across 96 narrow hourly sub-partitions per day — 200 tasks × 96 partitions meant most task-partition combos held a few hundred KB, each written as its own file. Eighteen months × 19,200 files/day compounded to 2.7M files. Every read then paid: driver listing 2.7M paths, planning millions of splits (the 8G planning OOM), task scheduling where startup exceeded read time 50-to-1, and per-file GET requests that throttled the object store into retry storms.
Fix
History: one Delta OPTIMIZE (ZORDERED by event time) compacted 2.7M files into 31,000 files averaging 132MB — a weekend run, after which COUNT(*) dropped 40 min → 3 min. Write path: coalesce sized from bytes-per-hour (hourly partitions now target 4-8 files of ~150MB via repartition on quiet hours, coalesce on busy), maxRecordsPerFile as a backstop, and Auto Compaction enabled for the table. Bill: LIST/GET charges fell $3,100 → $140/month. Guardrail: a daily file-count metric with a page at 2x the 128MB-ideal count.
Key lesson
  • 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.
Production debug guideMeasure files before bytes — these five checks separate file-count pain from real data growth.5 entries
Symptom · 01
Reads slowing over months while data volume looks stable, drivers straining in planning
→
Fix
Count files and compute average size per partition: list the table path, tally files per partition directory, and divide partition bytes by file count. Averages under ~64MB with thousands of files per partition confirm the diagnosis. Compare against the ideal (partition bytes / 128MB = files you should have) — a 100x gap means the tax exceeds the work and compaction is the fix, not bigger drivers.
Symptom · 02
Confirmed runts: partitions hold hundreds of KB-sized files each
→
Fix
Find the mint: check the writer's partition count (default 200 shuffle partitions is the classic culprit) multiplied by the table's sub-partition fan-out (hours × categories). If tasks-per-partition exceed bytes-per-partition/128MB, the writer is slicing air. Record today's file count per partition as the baseline you'll compact and gate against.
Symptom · 03
Historical table already holds millions of small files
→
Fix
Compact history without rewriting logic: Delta OPTIMIZE (with ZORDER on frequent filter columns to also speed reads) collapses runts into 128MB+ files in place; enable Auto Compaction (databricks) or schedule periodic OPTIMIZE so new runts merge automatically. Run it partition-by-partition on huge tables to bound each job, and verify file counts and p50 sizes after — then re-run your worst query to bank the before/after numbers.
Symptom · 04
Writer still minting runts every run (streaming micro-batches, hourly jobs, wide fan-out)
→
Fix
Size the write from bytes: repartition(n) with n = partition bytes / 128MB when you must widen or narrow (full shuffle, use at boundaries), coalesce(n) to only narrow without a shuffle (cheap, runs inside existing tasks). Add maxRecordsPerFile as a backstop against giant rows. For streaming, use trigger.AvailableNow with foreachBatch coalescing, or Delta's optimized writes + Auto Compaction so the engine merges runts for you.
Symptom · 05
Fixed today but ingestion patterns drift every quarter
→
Fix
Gate it: daily file-count and average-size metrics per table with pages at 2x ideal counts, a CI check that fails writes producing files under ~32MB in tests, and partition-design review for new tables (narrow hourly + high-cardinality sub-partitions need explicit file budgets). Re-graph quarterly — the mint restarts the moment a new job copies the old defaults.
Small-file responses compared
Root CauseHow to ConfirmFixPrevention
Historical runt accumulationAverages under ~64MB; millions of files; planning OOMs before readsOPTIMIZE in bounded partition passes; ZORDER for skipping bonusAuto Compaction on; scheduled OPTIMIZE; file-count metrics daily
Writer slicing air (default partitions × fan-out)Tasks-per-partition far exceed bytes/128MB; census runt ratio above 3xCoalesce (narrow) or repartition (exact) with n from bytes/128MBByte-sized writes in templates; CI assertion on test output sizes
Streaming micro-batch mintFiles-per-day near triggers × tasks; tiny per-trigger outputsAvailableNow triggers; foreachBatch coalescing; optimized writesCompaction cadence set at table creation; grain review per stream
Over-narrow partition grainRunt ratio worst on quiet hours and deep sub-partitionsCoarsen grain (daily over hourly); ZORDER instead of deep directoriesFile-budget math in design docs; quarterly census of top tables
⚙ Quick Reference
4 commands from this guide
FileCommand / CodePurpose
file_census.pyfrom pyspark.sql import functions as FPartition Sizing Math
sized_write.pyfrom pyspark.sql import functions as FWrite-Time Control
optimize_history.sqlOPTIMIZE lake.events WHERE dt >= '2026-09-01'Fixing History
streaming_sized_sink.pyfrom pyspark.sql import functions as FStreaming and Micro-Batches

Key takeaways

1
Small files tax every read with listing, planning, scheduling, and request tolls
size files to 128MB-1GB.
2
Census per partition (files vs bytes/128MB ideal); runt ratios above 3x name the exact writers to fix.
3
Coalesce cheap narrowing daily, repartition exact balance at boundaries
n always from bytes, never defaults.
4
OPTIMIZE history in bounded passes with ZORDER; Auto Compaction sweeps the future continuously.
5
Streaming needs fuller triggers plus per-batch coalescing
micro-batch defaults mint ~144K runts a day.
6
Gate counts in CI, metric daily, bill request tolls to owners
visible tolls enforce packaging discipline.

Common mistakes to avoid

5 patterns
×

Buying bigger drivers to fix slow reads on runt tables

Symptom
32G drivers plan millions of splits slightly faster while reads stay slow and request bills stay huge — tolls don't yield to heap
Fix
Census files first. Compaction (not heap) removes listing, planning, and request overhead at the source.
×

Leaving 200 default write partitions against narrow grains

Symptom
19,200 runts a day from 200 tasks × 96 sub-partitions — each task-partition combo minted as its own KB file
Fix
Size n from bytes-per-partition; coalesce cheap narrowing daily, repartition exact balance at boundaries.
×

Compacting without ZORDER (or ZORDERing without compacting)

Symptom
Half the win left behind: fewer files but no skipping gains, or clustered data still sliced into runts
Fix
Combine them in one pass — OPTIMIZE ... ZORDER BY on frequent filters compacts and clusters together.
×

Treating maxRecordsPerFile as the sizing strategy

Symptom
Wildly uneven files (wide vs narrow rows differ 100x in bytes) while counts look compliant
Fix
Size from bytes (n = MB/128); keep maxRecordsPerFile generous as pathology backstop only.
×

No file-count metric until the postmortem

Symptom
2.7M files discovered at 40-minute reads — dashboards showed terabytes (healthy) while counts compounded silently
Fix
Daily file-count plus average-size metrics per table, paged at 2x ideal; CI gate on test output sizes.
INTERVIEW PREP · PRACTICE MODE

Interview Questions on This Topic

Q01JUNIOR
Reads slow down over months while data volume is flat. What's your first...
Q02JUNIOR
Coalesce or repartition for fixing write-time file sizes?
Q03SENIOR
How do 200 default partitions interact with 96 daily sub-partitions?
Q04SENIOR
What does OPTIMIZE ... ZORDER BY buy beyond plain compaction?
Q05SENIOR
Design streaming writes that never mint runts. What are the three contro...
Q01 of 05JUNIOR

Reads slow down over months while data volume is flat. What's your first metric?

ANSWER
Files per partition and average bytes per file. Counts compounding (not bytes growing) explain gradual-then-sudden slowdowns — tolls exceed freight below ~64MB averages.
FAQ · 6 QUESTIONS

Frequently Asked Questions

01
What file size should Spark tables target?
02
Why did my driver OOM before reading anything?
03
Coalesce vs repartition — which and when?
04
Does ZORDER replace partitioning?
05
How often should OPTIMIZE run?
06
Can small files really cost more than storage?
N
Naren Founder & Principal Engineer

20+ years shipping production backend systems. Written from production experience, not tutorials.

Follow
✓ Verified
production tested
September 27, 2026
last updated
2,085
articles · all by Naren
🔥

That's Spark. Mark it forged?

5 min read · try the examples if you haven't

←
Previous
Spark AnalysisException: Cannot Resolve Column Name
5 / 5 · Spark
Next
Databricks Delta Concurrent Append Exception
→