skip to content

A job reads 300,000 files of 50 KB from an object store and idles twenty minutes before its first record — why?

level: seniorimportance: should knowfreq 62%

answer

  1. the lanes were idle, so it is not reading
  2. planning cost scales with file count
  3. one piece per file by default
  4. listing an object store is round trips
  5. overhead per piece exceeds the 50 KB read

basics

~20 s

Planning dominates, not reading. Listing 300,000 objects and turning each into a piece of the input costs far more than the 15 GB of payload, and that work is serial, so idle worker threads cannot absorb any of it.

solid answer

~50 s

The twenty minutes is planning, not reading. Before any record is processed, the reader has to list the tree on the object store — remote blob storage reached over the network, answering listings a page at a time — and measure each file, then turn the survivors into pieces of the input: one slice a worker thread reads end to end, and ordinarily one per file. Three hundred thousand pieces of a few milliseconds' bookkeeping each is minutes of work, while the payload is only 15 GB and would be read in seconds. The giveaway is idle lanes and one busy process. Remedies at read time: pack several files into one piece where the reader supports it, narrow the listing with column-value directories, or supply a prepared file list. The durable repair is fewer, larger files upstream.

go deeper

for a junior

Recall the inheritance: a finite job ordinarily gets one piece of the input per file, so a dataset of hundreds of thousands of tiny files hands the job hundreds of thousands of units of work.

for a middle

Explain the arithmetic: per-piece bookkeeping of a few milliseconds times the file count is minutes, while 50 KB is less than the round trip that opens it, so overhead exceeds the useful read.

for a senior

Show the diagnosis and the fix: idle lanes with one busy process means serial planning, and the remedies are packing files into pieces, narrowing the listing, or a supplied file list — not more machines.

for a principal

Own the upstream policy: a minimum file size for shared datasets, who enforces it and what the rewrite of existing data costs, weighed against every reader that pays the planning bill today.

## The symptom, stated precisely Wall clock is twenty minutes; records flow for the last two. **Worker threads** — the lanes inside each worker process, one of which is occupied by one piece of the input for its whole life — are idle for the first eighteen. One process, the **coordinating process** that turns the program into a graph, decides the pieces and hands them out, is busy the entire time. That shape is diagnostic: when the lanes are idle and one process is not, the job is not slow at reading, it is slow at *deciding what to read*. ## Where the time goes 1. **Listing.** The reader must discover the files. An **object store** answers a listing in pages over the network, and each page is a round trip. Three hundred thousand objects is thousands of sequential round trips, and if the objects are spread across many prefixes the tree is walked rather than scanned. 2. **Measuring.** Deriving a division needs each file's size, which may be another metadata call per file depending on what the listing returned. 3. **Deriving the pieces.** In a **finite job** the pieces are derived from the stored bytes, and the default derivation is **one piece per file**, because a file is the smallest unit whose divisibility can be reasoned about. 4. **Handing them out.** Each piece costs bookkeeping: assignment, launch, status reporting, completion and commit accounting. ## The arithmetic that makes it obvious | quantity | value | consequence | |---|---|---| | files, each 50 KB | 300,000 | 300,000 pieces of the input at one per file | | total payload | ~15 GB | seconds of actual reading spread over a few hundred lanes | | per-piece bookkeeping, order of milliseconds | ~3 ms | ~15 minutes of serial work in one process | | useful read per piece | 50 KB | smaller than the round trip that opens the file | The last row is the whole point of this leaf: the **per-piece scheduling overhead exceeds the read**. When one unit of work costs more to hand out than to perform, more units make the job slower, not faster. ## Why more machines do not help The piece count sets how many lanes can be busy; the machine count only sets how many lanes exist. Here the piece count is enormous and the lanes are already idle, so the constraint is elsewhere: the listing and the piece derivation are essentially serial work in one place, and adding capacity leaves them exactly as long. Doubling the cluster doubles the bill for the idle eighteen minutes. This is also why the problem is invisible in a small test: with three hundred files the same planning work takes under a second. ## What actually helps, at read time - **Pack several files into one piece.** Many readers can group whole files up to a byte target, so the piece count is set by total bytes rather than by file count. Support varies — some readers give one piece per file unconditionally — so confirm it rather than assuming it. - **Narrow the listing.** If the files sit in **column-value directories** — grouped into directories named for one column's value — a filter on that column drops whole directories while the plan is built, and the listing and the piece derivation shrink with them. - **Avoid re-listing.** Some readers accept a prepared list of files, or read a manifest the producer maintained, which converts thousands of round trips into one read. - **Do not reach for the piece count as a knob here.** Merging pieces after they exist does not undo the listing or the derivation that already happened; the expense was incurred before any lane started. The durable repair is to stop inheriting the file population: rewriting many small files into fewer larger ones is a table-maintenance subject in its own right, and the file count a previous job left behind is that job's output-shape decision. This leaf owns only what the next reader pays for it. ## What varies between engines - **Where the listing runs.** Some engines list from the coordinating process, others distribute the listing across workers or accept a supplied file list, which changes whether the eighteen minutes is serial at all. - **Whether files are packed into one piece**, and whether the packing respects locality or simply fills a byte target. - **Continuous jobs** do not derive a piece count from files at all: the author sets a **declared operator width** that stands until the job is restarted, so a continuous reader over a directory of small files pays per-file open cost but not this planning cliff in the same shape. - **How overhead per piece is spent** differs by an order of magnitude between an engine that launches a fresh process per unit of work and one that reuses long-lived lanes, which is why the same file population is merely annoying on one engine and fatal on another.

  • How do you tell this apart from a job that is genuinely slow at reading the data?
    Look at when lanes are busy. A read-bound job has every lane working from early on and a wall clock that tracks bytes; this one has idle lanes, one busy process and a wall clock that tracks the file count. Re-run against a tenth of the files: planning cost falls tenfold, read cost falls with the bytes.
  • Would merging the pieces after reading fix it?
    No. Gluing neighbouring pieces together in place lowers the count for later steps, but the listing and the per-piece derivation were already paid before any lane started. The saving has to happen at plan time — by packing files into pieces, by narrowing the listing, or upstream by leaving fewer files behind.
  • The same 300,000 files live on storage that runs on the compute machines instead. What changes?
    The metadata calls become local rather than network round trips, so listing is much cheaper, and work can be placed near the bytes. The per-piece bookkeeping is unchanged, though, so a large file count still costs; the cliff arrives later rather than never.

saying these in an interview costs you the question

  • Blames network bandwidth while the worker threads sit idle
  • Adds machines to a job bottlenecked on serial planning work
  • Thinks small files are slow only because of decompression
  • Believes the piece count is capped by the machine count
  • Assumes every reader groups small files into one piece
  • Measures the fix on a test set of a few hundred files