skip to content

Which two quantities set the piece count for a finite job reading 800 GB on 400 worker threads?

level: middleimportance: must knowfreq 74%

answer

  1. two pressures, never one
  2. bytes per piece, and lanes available
  3. size proposes, lanes floor, overhead ceilings
  4. even waves beat a ragged tail

basics

~20 s

A target size per piece and the total worker-thread count. Size keeps each piece inside one thread's memory and worth its fixed cost; thread count keeps the lanes busy in even waves. Reconcile them rather than using either alone.

solid answer

~50 s

Two pressures pull against each other. Bytes per piece: the records a piece holds in flight have to fit the memory one worker thread is given, and a lost piece is recomputed whole, so teams carry a target size per piece - commonly a few hundred megabytes for file-backed input, though the number varies with engine, format and memory budget. Worker threads: a piece occupies one lane for its whole life, so 401 pieces on 400 lanes is two waves and roughly double the wall clock of one. For 800 GB at a 256 MB target that is about 3,200 pieces, which is eight clean waves over 400 lanes. If the size figure fell below the lane count you would raise it and accept smaller pieces - but stop before pieces get so small that per-piece overhead outweighs the work.

go deeper

for a junior

Recall that the input is cut into pieces and one worker thread handles one piece at a time, and that two different numbers are in play: how many pieces there are and how many lanes exist to run them.

for a middle

Explain both pressures and the arithmetic that reconciles them: bytes divided by a target size, compared against the lane count, with per-piece overhead as the ceiling. Be able to say what breaks at each extreme.

for a senior

Show you size the step that actually matters rather than the whole job, that you know a large piece makes retry expensive as well as memory tight, and that you check the resulting wave count rather than trusting a round number.

for a principal

The angle is policy: a default sizing rule for hundreds of jobs, expressed as a computation from current bytes and current lanes rather than a literal number, with the cost of being wrong in each direction made explicit to the teams that inherit it.

## The unit being counted A **piece of the input** is one slice of the stored input that a single **worker thread** reads and processes from start to finish. The **piece count** is how many such slices the run is divided into. A worker thread is one lane inside a **worker process** - one operating-system process on one machine, holding its own memory and several lanes. Four numbers get called *the count* and only one of them is being chosen here: - **the piece count** - how many slices the input is cut into; this is the decision; - **the worker-thread count** - how many lanes exist to run them; - **the machine count** - how many machines those lanes are spread over; - **the output-file count** - what the run leaves behind, which is downstream of this decision. The machine count only sets how many lanes exist. The piece count sets how many of those lanes can be busy. ## Pressure one: bytes per piece A piece is simultaneously the unit of memory, the unit of retry and the unit of fixed overhead, and all three argue about its size. - **Memory.** Everything a piece holds in flight - its records, the partial aggregate it is building, its sort buffer - has to fit the share of the worker process's memory that one lane gets. When it does not, the engine **spills**: it writes part of the working set to local disk and reads it back, trading throughput for survival. - **Retry.** When a lane or a machine is lost, the piece is redone from its input. A 2 GB piece costs 2 GB of recomputation; a 256 MB piece costs 256 MB. - **Fixed overhead.** Every piece is created, dispatched, tracked and collected whether it holds a gigabyte or forty kilobytes. Teams therefore carry a **target size per piece**. The value varies with the engine, the file format and the memory a lane is given; what does not vary is the reason for having one - it is the size at which the fixed cost is amortised while the working set still fits. ## Pressure two: worker threads and waves A piece occupies one lane for its whole life, so pieces run in **waves**: as many start as there are free lanes, and the next wave starts as lanes free up. | pieces | lanes | waves | last wave | |---|---|---|---| | 400 | 400 | 1 | full | | 401 | 400 | 2 | one piece, 399 lanes idle | | 800 | 400 | 2 | full | | 3,200 | 400 | 8 | full | The 401 row is the one interviewers reach for: one extra piece roughly doubles the wall clock of a single-wave job. Aiming at a small multiple of the lane count - two to four times is a common rule of thumb, not a rule - also lets one unlucky slow piece overlap with other work instead of becoming the whole tail. ## Reconciling the two 1. Divide the bytes the step will actually handle by the target size. For 800 GB at a 256 MB target that is about 3,200 pieces. 2. Compare with the lane count. 3,200 against 400 lanes is eight even waves - take it. 3. If the size figure lands *below* the lane count, raise it toward the lanes and accept smaller pieces. 20 GB at the same target is 80 pieces, which leaves 320 lanes idle for the whole run. 4. Stop raising once pieces become trivially small, because the fixed per-piece cost then starts to dominate the real work. Neither quantity alone is the answer: the size target proposes the count, the lane count puts a floor under it, and per-piece overhead puts a ceiling on it. ## What each extreme actually costs | | too few pieces | too many pieces | |---|---|---| | lanes | most sit idle | all busy, much of it bookkeeping | | memory per piece | large, spilling likely | trivially small | | retry after a loss | whole large piece redone | cheap | | coordination | light | the coordinating process becomes the bottleneck | | what you see | long wall clock on an idle-looking cluster | long wall clock with enormous unit churn | ## Finite and continuous jobs answer this differently In a **finite job** - a run over an input that ends - the bytes can be measured before anything starts, so the count can be derived from them. In a **continuous job** - a run over an input with no end - nothing about the input can be measured in advance, so the author states a **declared operator width**: how many copies of each operator run. The same two pressures reappear in a different currency: instead of bytes per piece, the throughput one lane can sustain against the arrival rate; instead of retry cost, the state one lane has to hold. And in most continuous runtimes that number stands until the job is restarted, which is why an approximately right answer up front matters more there than in a finite job. ## What an interviewer is listening for A number and the thing that sets that number. *I would raise it* is a weak answer. *This step handles about 800 GB, I want pieces in the hundreds of megabytes so a lane does not spill, that is roughly three thousand, and against four hundred lanes that is eight clean waves* is the answer - and it never once names a setting.

  • Why prefer a piece count that is a small multiple of the worker-thread count rather than exactly equal to it?
    Equal means one wave with no slack: any piece that runs long is the whole tail, and one piece more than the lane count doubles the wall clock. A small multiple gives the scheduler spare work to overlap, so a slow piece is absorbed by other pieces still running rather than leaving the cluster idle waiting for it.
  • Does the same reasoning apply to a continuous job, where no input size can be measured?
    The shape holds, the quantities change. You size the declared operator width from the throughput one lane sustains against the expected arrival rate, plus the state a lane has to hold, rather than from stored bytes. The bigger difference is timing: in most continuous runtimes the width stands until a restart, so there is no cheap correction later.
  • What if the pieces come out at wildly different sizes even though the count is right?
    Then the count is not your problem. Uneven pieces caused by an uneven key distribution are skew - records concentrating on a few pieces - and the diagnosis and remedies for that belong to Skew and Stragglers. Choosing the count assumes a roughly even division; it cannot fix one that is not.

saying these in an interview costs you the question

  • More pieces always means more parallelism, so set it as high as possible
  • Setting the piece count equal to the machine count
  • Sizing purely by bytes and never checking against available worker threads
  • Treating a piece as free to create, so overhead never matters
  • Assuming every engine derives the right count for you, continuous jobs included
  • Naming a setting instead of naming a number and what sets it