skip to content

Why does cutting a 2 GB input into 50,000 pieces usually run slower than cutting it into 200?

level: juniorimportance: should knowfreq 60%

answer

  1. cost per piece, not per byte
  2. 40 KB of work, full bookkeeping
  3. one coordinator tracks every piece
  4. the serial section, not the reading

basics

~20 s

Each piece carries a fixed cost that does not shrink with its bytes: it is created, dispatched, tracked and collected, and its input is opened and closed. At 40 KB per piece that bookkeeping outweighs the reading, and the coordinating process becomes the bottleneck.

solid answer

~40 s

A piece of the input is one slice a single worker thread reads start to finish, and it costs something before it touches a byte: the coordinating process - the one process that turns the program into a graph and hands out work - has to create it, dispatch it, track it and collect its result, and the lane has to open and close the input. That cost is roughly constant per piece. At 200 pieces a piece is 10 MB and the constant is noise; at 50,000 pieces a piece is about 40 KB and the constant is the job. Worse, the coordination is largely serial in one process, so the cluster finishes pieces faster than it can be handed new ones. The 2 GB is read either way.

go deeper

for a junior

Recall that every piece carries a fixed cost whatever its size, so cutting an input very finely can spend more on managing the work than on doing it.

for a middle

Explain the mechanics: creation, dispatch, input open and close, result collection, and the multiplication of bookkeeping when the pieces feed a step that needs records from every other worker.

for a senior

Recognise the signature in production - high unit churn, low per-lane utilisation, a wall clock that worsens when machines are added - and know that merging neighbouring pieces in place is the cheap correction.

for a principal

The angle is where the platform's floor should sit: a minimum bytes-per-piece guard applied across jobs costs a little parallelism on small inputs and protects the coordinator on every large one.

## What a piece costs before it reads anything A **piece of the input** is one slice of the stored input that a single **worker thread** - one lane inside a worker process - reads and processes from start to finish. Cutting the input finer does not make the bytes cheaper to read; it multiplies everything that is charged *per piece* rather than per byte. In a staged runtime the per-piece charges are roughly: - **Creation and bookkeeping.** The **coordinating process** - the single process that turns the program into a graph, decides the pieces and hands them out - holds a record for every piece: its inputs, its assignment, its attempt number, its status. - **Dispatch and start-up.** Each piece is sent to a lane, set up there, and the lane reports back when it is done. - **Opening the input.** Each piece opens its source, seeks to its offset, reads headers or footers it needs, and closes again. On remote storage that is one or more network round trips whose cost is independent of how many bytes follow. - **Result collection.** Each completion returns a status and some metadata to the coordinating process, which must merge it into the run's state. - **Feeding a wide step.** A **wide step** is one a worker cannot finish from the records it already holds, because it needs records currently sitting on every other worker. Each upstream piece has to organise its records into one group per downstream piece, so the bookkeeping there grows as the product of the two counts. The byte-level cost of the exchange itself belongs to Moving Data Between Workers; only the multiplication is this leaf's concern. ## When the overhead outruns the work With 2 GB cut into 50,000 pieces, each piece holds about 40 KB. The fixed cost per piece is small but not zero - on the order of milliseconds in a staged runtime, and engines differ by an order of magnitude here - while reading and decoding 40 KB is far quicker than that. The run becomes mostly bookkeeping. | piece count | bytes per piece | what dominates | |---|---|---| | 20 | 100 MB | reading, but only a few lanes are busy | | 200 | 10 MB | reading; the fixed cost is noise | | 50,000 | 40 KB | the fixed cost, paid 50,000 times | ## The coordinating process is a serial section The second effect is worse than the arithmetic suggests. One process tracks every piece, so its bookkeeping, its heartbeats and its result handling are a serial section in an otherwise parallel run. Past some count - which depends entirely on the engine, and is the kind of number worth measuring rather than quoting - it cannot hand out work as fast as the lanes finish it, and the cluster idles while the coordinator catches up. A run in that state looks strange on a dashboard: enormous unit churn, low per-lane utilisation, and a wall clock that gets *worse* when you add machines. ## Long-lived lanes change the shape, not the moral All of the above assumes pieces are launched and torn down, which is the finite-job picture. In a **continuous job** - a run over an input with no end - the lanes are long-lived and there is no per-piece launch at all; the author sets a **declared operator width** instead, and it usually stands until the job is restarted. The analogous cost of setting that width too high is still real, but it is a different bill: - per-lane buffers and per-lane state overhead, paid for the life of the job; - network connections between producing lanes and consuming lanes, which grow as the product of the two widths; - fixed per-lane work at every recovery point, so a wider job takes longer to record one. So *do not over-divide* survives the move to a continuous runtime, but the reason changes from launch overhead to standing per-lane cost. ## What this is not - It is **not** the small-files problem. The file count this write leaves behind, and what those files cost the next job and the storage bill, belong to The Output Shape, and repairing them to lakehouse table compaction. - It is **not** skew - records concentrating on a few pieces so that they take far longer than the rest. Here every piece is equally tiny. - It is **not** an argument for the opposite extreme. Very few, very large pieces leave lanes idle, make each retry expensive and push a lane into **spilling**, which is writing part of its working set to local disk because it no longer fits in memory. ## The cheap correction When the count is already too high and the data is already divided, **merging pieces without moving records** - gluing neighbouring pieces together in place - cuts the count without a wide step. It is much cheaper than a full redistribution, and the catch is that it only merges neighbours, so each survivor is correspondingly larger and the merge cannot even out sizes.

  • The job now runs worse after adding machines. How does that fit?
    It fits the coordinator-bound case exactly. More lanes finish tiny pieces faster, so they demand new work faster, and the single process handing work out is already saturated. Extra machines add demand on the serial section without relieving it, so utilisation falls and wall clock can get worse. The fix is fewer, larger pieces, not more machines.
  • If the count is already too high and the data is divided, what is the cheap way down?
    Merge pieces without moving records: glue neighbouring pieces together in place so several become one. That avoids a wide step entirely, which is why it is cheap. The limitation is that it only combines neighbours, so if the pieces were uneven to begin with the survivors stay uneven.

Shipping. Every parcel carries the same paperwork, the same label and the same scan at each hop, whatever it weighs. Splitting one truckload into ten thousand envelopes does not move fewer kilograms - it just means you now spend more on paperwork than on freight, and the dispatch desk, not the road, is what everything is queued behind.

saying these in an interview costs you the question

  • A piece costs only in proportion to the bytes it holds
  • More pieces is always more parallelism, so more is better
  • Adding machines will fix a job drowning in tiny pieces
  • Confusing this with the small-files problem the write leaves behind
  • Assuming the coordinating process scales with the cluster