A finite job's stored input is divided into eight pieces and the cluster offers 200 worker threads; how many can work at once?
answer
- two counts, not one
- lanes without pieces do nothing
- concurrency is the smaller number
- a piece holds one lane throughout
basics
~10 sEight. The piece count is the ceiling on how many worker threads can be busy; the machine count only decides how many threads exist, so the other 192 lanes stay idle.
solid answer
~40 sEight, and the other 192 lanes sit idle. A **piece of the input** is one slice of the stored data that a single **worker thread** reads and processes from start to finish, so a piece occupies one lane for its whole life and is not subdivided mid-flight. That makes the **piece count** the ceiling on concurrency: the machine count and the worker-thread count only decide how many lanes exist, and a lane with no piece to run does nothing. Granting the job three times the machines gives 600 lanes and still eight busy ones. Two things genuinely raise the ceiling: dividing the stored input into more pieces, or redistributing records after the read so that later steps run at a larger count — and how large that number should be is a separate question.
go deeper
Recall the one-to-one fact: one piece of the input is read start to finish by one worker thread. From there, say out loud which number is smaller — the pieces or the lanes — and that the smaller one is how much is happening at once.
Explain the mechanics: a piece is handed out once and is not subdivided, so lanes idle whenever pieces run out, and the fix is either more pieces from the stored form or a redistribution into more downstream pieces. Name the cost of that redistribution.
Read it off a real run: lanes occupied against lanes available over time tells you immediately whether you are ceiling-bound or capacity-bound, and stops a team buying machines that cannot be used. Say which of the two the evidence supports before proposing anything.
The platform angle is that a cluster sized above the typical piece count is a standing bill for capacity no job can occupy. Weigh a division standard on shared inputs against the migration cost of rewriting how that data is stored.
## What a piece of the input is A **piece of the input** is one slice of the stored data that a single **worker thread** reads and processes from beginning to end. Two surrounding terms matter: a **worker process** is one operating-system process on one machine, holding its own memory and several worker threads; a **worker thread** is one lane inside that process. A piece occupies one lane for its whole life — it is handed out once, it is not cut in half mid-flight so that an idle lane can help, and the lane is free again only when the piece is finished. Nearly everything in this subject falls out of that one sentence. One disambiguation, because the same English word covers both: one piece of the input per worker thread is **not** the same thing as one **column-value directory** on disk — stored files grouped into directories named for one column's value. Both get called a partition in conversation, and they are different objects with different owners. ## Four counts, and only one of them is the ceiling | Count | What sets it | What it limits | |---|---|---| | **the piece count** | in a finite job, how the stored input divides; in a continuous job, the operator width the author declared | how many worker threads can be busy at once | | **the worker-thread count** | lanes per worker process, times the number of processes | how many pieces can be in flight at once | | **the machine count** | how many machines the job was granted | how many worker processes can exist | | **the output-file count** | how many pieces reach the write | what the next job inherits from storage | Concurrency is the **smaller** of the first two. The third sits one step further away: machines let processes exist, processes let lanes exist, and lanes are only useful if there is a piece to put in one. The everyday mistake is to quote the fourth number in the first row's place — to answer "how parallel is this job" with the size of the cluster. ## The worked case Eight pieces, 200 lanes: eight lanes work and 192 are idle from the first second. The run takes roughly as long as its slowest piece, plus the fixed cost of starting and finishing. Triple the machines and there are now 600 lanes, eight of them busy. Halve the machines and there are 100 lanes, eight of them busy, and the run takes the same time — which is the useful version of the observation, because it says the job was paying for capacity it could never use. ## The degenerate case Push the same arithmetic to its end: if the stored form yields exactly **one** piece, the ceiling is one, and a hundred machines finish no sooner than one. This is the case interviewers reach for, because it is the cleanest proof that the piece count and not the machine count is the quantity in play. ## What actually raises the ceiling - **Divide the stored input into more pieces.** More slices, more lanes that can be occupied at once — up to the number of lanes that exist. - **Redistribute records after reading.** A step that a worker cannot finish from the records it already holds, because it needs records currently sitting on every other worker, is a **wide step**; it can hand its output to a larger number of downstream pieces. The read itself stays as wide as the stored form allowed, but everything after the redistribution can be wider. It costs a full pass of records across the network, and that cost is a subject of its own. - **Nothing else.** Not a larger worker process, not more memory, not more lanes per process, not more machines. Those change what a lane can hold, not how many pieces exist to fill lanes with. ## What varies between runtimes The ceiling rule is common to the whole class; how the number gets set is not. 1. Some runtimes derive the pieces from the stored bytes and re-derive them on **every** finite run, so the count tracks the data by itself. 2. Others fix the split count when the job is **submitted**, so the same input divides the same way until the next submission. 3. A **continuous job** — a run over an input with no end — has no stored byte count to measure at all, so the author states a **declared operator width**: how many copies of each operator run. 4. Where continuous work is implemented as a rapid succession of small finite jobs, each small job derives its own read-side division from what arrived, while the stated width still governs the rest. So "more machines, more parallelism" is wrong everywhere in this class, but "the count is recomputed for you each run" is right only in some of it. ## What this question does not decide How many pieces you should ask for, and how to size them, is a separate decision. So is which machine a given piece is scheduled onto, why one piece might be far heavier than the rest, and how large the worker process running a lane should be. This question settles only the relation between the counts.
- The input yields two pieces of very different size on a cluster of 200 lanes. What sets the run's duration?Roughly the larger piece. Two lanes work, the rest are idle, and the finish waits on the slower of the two — so duration tracks the biggest piece rather than the total bytes. Why one piece is far larger than another, and the remedies for it, is a separate subject.
- Does giving each worker process more lanes raise the ceiling in this case?No. More lanes per process raises the worker-thread count, which was already 200 and already far above the eight pieces. It changes how many pieces could be in flight, not how many exist. Only a larger piece count, or a redistribution into more downstream pieces, moves the ceiling.
- After a redistribution, is the downstream count still the number derived from storage?No. The number derived from the stored form governs the read; once records are redistributed, the number of downstream pieces comes from elsewhere — typically a value the run carries rather than one measured from bytes. Picking that value is owned by the sibling subject on choosing how many pieces.
Eight crates arrive at a loading dock and two hundred porters are on shift. Eight porters carry; the rest watch. Hiring more porters is not the fix — repacking the shipment into more crates is.
saying these in an interview costs you the question
- Says adding machines always increases how much of the job runs at once.
- Quotes the machine count when asked for the degree of parallelism.
- Assumes an idle lane will be handed half of another lane's piece.
- Thinks more lanes per worker process helps when there are too few pieces.
- Confuses a piece of the input with a column-value directory on disk.