skip to content

A Spark job writes 40,000 tiny Parquet files a day — what caused it and how do you fix it?

level: seniorimportance: should knowfreq 50%

answer

  1. multiply two numbers to get the count
  2. every task visits every directory
  3. make each directory belong to one task
  4. cap the rows if a file gets huge
  5. the obvious one-file fix is a trap

basics

~20 s

File count is partitions multiplied by the distinct partitionBy values each task touches: 200 shuffle partitions across 200 directories can emit 40,000 files. Fix it by repartitioning on the same columns you partition by before writing, so one task owns each directory.

solid answer

~40 s

Spark writes **one file per in-memory partition per output directory**, so the count is roughly `numPartitions × distinct partitionBy values touched per task`. With `spark.sql.shuffle.partitions` at its 200 default and `partitionBy("dt","country")` spanning 200 directories, every task contributes a sliver to every directory and you get 40,000 files. The direct fix is `df.repartition("dt","country")` immediately before the write: rows for a directory then live in one task, which writes one file per directory. Add `maxRecordsPerFile` to cap any file that would grow too large, or `repartition(n, "dt", "country")` when a directory needs several files. Do **not** reach for `coalesce(1)` — it is narrow, so it throttles the entire upstream computation to a single task. If the partition columns are high-cardinality, the real fix is to stop partitioning by them.

code

python · 13 lines
python
# Before: 200 shuffle partitions x 200 directories = up to 40,000 files
(df.write
   .partitionBy("dt", "country")
   .mode("overwrite")
   .parquet("s3://lake/events/"))

# After: one task owns each directory -> ~200 files
(df.repartition("dt", "country")
   .write
   .option("maxRecordsPerFile", 5000000)
   .partitionBy("dt", "country")
   .mode("overwrite")
   .parquet("s3://lake/events/"))

go deeper

for a junior

Know that Spark writes one file per partition, so the number of output files follows the number of partitions rather than being chosen by the writer.

for a middle

Explain the partitions-times-directories multiplication, and why repartitioning on the same columns used in partitionBy collapses it to one file per directory.

for a senior

Diagnose from real evidence — write-stage task count, directory count, mean file size — and weigh the fix's cost, including the write-stage skew that concentrating a directory into one task creates.

for a principal

Own file-layout policy for a shared lake: target file sizes, which columns a dataset is partitioned by, whether compaction is a scheduled service, and the read-cost these choices impose on every consumer.

## The arithmetic behind the file count Spark's file writer is simple and its output count is predictable. Each task writes to the output path, and when the write is partitioned by columns it opens a separate file for each distinct combination of those column values it encounters. So: ``` files ≈ (number of in-memory partitions) × (distinct partitionBy value combinations seen per task) ``` A job whose last shuffle produced the default 200 partitions, writing with `partitionBy("dt", "country")` where the batch spans 200 date/country combinations, produces up to 200 × 200 = 40,000 files — each a few hundred kilobytes. Nothing has gone wrong internally; the layout is exactly what was asked for. ## Why it hurts - **Read planning.** Every subsequent Spark job that reads this dataset must list and open all those files. Listing dominates on object stores, where each `LIST` and `HEAD` is a network round trip; the metadata phase can take longer than the scan. - **Wasted parallelism on read.** Spark packs small files together toward `spark.sql.files.maxPartitionBytes` and charges each file an opening cost via `spark.sql.files.openCostInBytes`, so tiny files raise per-task overhead relative to real work. - **Lost compression and statistics.** Parquet's dictionary encoding, run-length encoding and row-group min/max statistics all work better over more rows. A 200 KB Parquet file carries proportionally more footer and schema overhead and gives predicate pushdown almost nothing to prune with. - **Storage-layer cost.** On HDFS every file and block consumes NameNode heap; on object storage you pay per request. Either way, small files are a tax paid by every reader forever. ## Diagnosis Count files per directory (`hdfs dfs -count` or an object-store listing) and check the mean size — anything well under ~100 MB for Parquet is a candidate. Then look at the job: what was the partition count going into the write (visible as the write stage's task count in the Spark UI), and how many distinct `partitionBy` values did the batch contain? Multiplying the two should reproduce the observed file count, which confirms the mechanism rather than guessing at it. ## The primary fix: align in-memory partitioning with write partitioning ```python (df.repartition("dt", "country") .write.partitionBy("dt", "country") .mode("overwrite") .parquet("s3://lake/events/")) ``` Hash-partitioning by the same columns puts all rows for a `(dt, country)` pair into one task, so that task writes exactly one file for that directory. File count drops from 40,000 to 200. This is the standard idiom and the answer an interviewer is listening for. Two refinements: - If one directory is much larger than the rest, a single file becomes unwieldy. Use `.option("maxRecordsPerFile", 5000000)` (or the session-level `spark.sql.files.maxRecordsPerFile`) to cap file size — Spark rolls to a new file when the cap is hit. - If every directory holds several hundred MB, use `repartition(400, "dt", "country")` so each directory gets a handful of properly sized files rather than one enormous one. Be aware that `repartition` by the partition columns concentrates each group in one task, so a skewed group produces a slow task. That is a deliberate trade: you are buying good file layout with some write-stage imbalance. ## What not to do `coalesce(1)` before the write is the classic wrong fix. `coalesce` is a narrow transformation with no stage boundary, so the single-partition limit propagates up the entire fused stage — the scan, the filters and the projections all run in one task. You will trade 40,000 small files for a job that takes hours or dies on memory. ## Prevention and the layout question If the `partitionBy` columns are high-cardinality — `user_id`, an hour-plus-minute timestamp, an event UUID — no amount of repartitioning saves you, because the directory count itself is the problem. Partition on the coarsest column your query predicates actually filter on, usually a date, and use clustering (`repartitionByRange` on a secondary column, or a table format's clustering feature) for finer pruning rather than more directories. For a streaming job the same arithmetic applies per micro-batch, and there the answer is usually a periodic **compaction** job that rewrites yesterday's directories into properly sized files, plus a longer trigger interval so each batch writes more. Adaptive Query Execution helps indirectly: with `spark.sql.adaptive.enabled` on (the default since Spark 3.2) and its coalescing feature, an over-provisioned post-shuffle stage may already be merged down before the write, so the multiplier is smaller than the configured 200. It is a mitigation, not a substitute for aligning the write partitioning.

  • Why not just call coalesce(1) before the write?
    `coalesce` is narrow and adds no stage boundary, so the one-task limit applies to everything fused with it — the scan, filters and projections all run single-threaded. You would replace a file-layout problem with a runtime and memory problem. If you genuinely need one file, `repartition(1)` keeps the upstream stage parallel by inserting a shuffle.
  • After repartitioning by the partition columns, one write task runs far longer than the rest. Why, and what do you do?
    Concentrating a directory into one task means the largest directory becomes the slowest task — you converted file skew into task skew. Options: `repartition(n, cols)` so big groups spread over several files, `maxRecordsPerFile` to roll files within the task, or salting the largest group. Some imbalance in the write stage is usually an acceptable price for good file sizes.
  • How would you handle this in a Structured Streaming job that writes every 30 seconds?
    Each micro-batch writes at least one file per partition per directory, so small files are structural. Lengthen the trigger interval so each batch carries more data, keep the write-side partition count low, and run a scheduled compaction job that rewrites completed directories into large files. Do not try to solve it purely inside the streaming query.

saying these in an interview costs you the question

  • Suggests coalesce(1) as the standard fix for small files
  • Thinks the writer already merges files across tasks automatically
  • Blames the file format instead of the partition-times-directory arithmetic
  • Partitions by a high-cardinality column and expects repartition to save it
  • Assumes AQE alone guarantees well-sized output files

context