A job sorts each of its 200 pieces locally and writes them out; why is the combined result not ordered end to end?
answer
- ordering is a property of placement
- a sort never moves records between pieces
- hashing spreads evenly, scatters the order
- contiguous spans, then read in index order
basics
~20 sSorting inside a piece orders only that piece. Unless the routing that filled the pieces gave each one a contiguous span of the key range, piece 3 can hold keys that belong before piece 2's, so reading the pieces in order interleaves ranges.
solid answer
~40 sA local sort is an ordinary in-process sort over the records one worker happens to hold, so 200 of them give 200 independently ordered runs, not one ordered result. Which records a worker holds was decided earlier, by the routing function that sent each record to a destination. If that function hashes the key — the usual choice, because it spreads work evenly — the piece index tells you nothing about where those keys sit in the order. Equal keys still meet, since equal keys hash alike, but that is the only ordering you get. To read the pieces back in order you need the routing itself to be order-aware: boundaries chosen in key space so each piece owns one contiguous span, sorted locally, and then read in piece order.
go deeper
Recall that a sort inside one piece can only order the records that worker holds, and that which records it holds was decided earlier by the routing function.
Explain why hashing the key spreads work evenly but destroys any link between key value and piece index, and what has to change in the routing for pieces to concatenate into an ordered result.
Show judgement about when an end-to-end order is worth its redistribution at all, and say plainly that it is a claim about a complete result, so over an endless input it applies to a bounded chunk or not at all.
Frame ordering as a contract the platform publishes to consumers: ordered-within-piece is cheap and often sufficient, end-to-end ordering is a recurring cost on every run, and picking the wrong one leaks into every downstream reader.
## What a local sort actually guarantees A job's input is cut into **pieces** — one contiguous share of the input that a single worker processes on its own — and each piece is handled by one **unit of work**: one worker computing one piece once, so a retry replaces exactly that unit. When a step sorts, what runs inside that unit is an ordinary sort over the records that worker is holding at that moment. It has no view of any other worker's records and no way to acquire one. So after 200 local sorts you hold 200 independently ordered runs. That is genuinely useful — each run can be streamed, merged or searched as an ordered sequence — but it is not an **ordered result**, because an ordered result is a claim about the output read as a whole: every record in run 0 sorts before every record in run 1, and so on. ## The routing function decides this, not the sort Records are in a given piece because of the **routing function**: a rule applied to each record's key that names the one destination for it, so that equal keys always meet. Most routing is built to spread work evenly, which in practice means deriving the destination from a hash of the key. A hash deliberately destroys any relationship between a key's value and where it lands: - keys that are adjacent in the ordering land in unrelated pieces; - the piece index carries no information about which span of keys the piece holds; - equal keys do meet, because equal keys hash alike — that is the only ordering a hash-based routing function gives you; - sorting afterwards cannot repair this, because the sort never moves a record to a different piece. That last point is the whole answer. Sorting is a step that happens *after* placement, and placement is what an end-to-end order is a property of. ## The two ways to get an ordered result | approach | what moves | who sorts | where it breaks down | |---|---|---|---| | Collect the result onto one worker and sort it there | the entire result, to one machine | one worker, on its own | one machine's memory and local disk; no parallelism; the whole result crosses the network into one place | | Cut the key range at boundaries, then sort each piece | each record once, to the piece owning its span | every worker, in parallel | the boundaries are only as good as the sample they were drawn from | The first is fine when the result is small — a few thousand rows for a report, or a final aggregate. It stops being an option as soon as the ordered result is a large fraction of the input, which is the case an interviewer is asking about. ## The pipeline the second approach runs 1. Sample the **ordering key** — the field the result must be ordered by — across the input. 2. Choose boundary values from that sample so that each piece will receive one contiguous span of the key range, with roughly equal numbers of records per span. 3. Redistribute: a **redistribution** is the step where every worker sends each record it holds to whichever worker will handle that record's key, here using a routing function that compares the key against the boundaries instead of hashing it. 4. Sort within each piece, in parallel, so each piece is an ordered run covering a span nobody else covers. 5. Read the pieces back in index order. This is a read and a concatenation, not a sort: no worker ever has to hold more than its own piece. ## Finite results, endless inputs, and what varies An end-to-end ordered result is a statement about a *complete* result, so it is a finite-input idea. Over an unbounded input there is no complete result to order; runtimes over endless input generally offer ordering only inside a bounded chunk, such as one grouping interval that has closed, or not at all. Runtimes also differ in how much of the pipeline above they do on your behalf: some derive boundaries from a sample whenever an ordered result is requested, some expect the author to supply the boundary values, and the older designs express the whole thing as an explicit extra step in the program. Do not assume the one you learned is the model. ## What it costs, and what it is not - The sample has to be drawn before any record moves, because the boundaries are an input to the routing function. Some designs pay an extra read of the input for it; others fold it into work already happening. Either way the boundaries are fixed before the movement starts. - The movement itself is a full redistribution: every record crosses the network once, priced in bytes rather than rows. - The local sorts have a working set, and if it exceeds what the worker can hold, part of it is written to local disk — a separate subject with its own trade-offs. - Nothing about this makes the *contents* wrong. Locally sorted pieces are correct data with a weaker ordering contract than the reader assumed.
- If the pieces are only ordered internally, is the output wrong, or just differently ordered?The records are all correct and none are lost — only the ordering contract is weaker than a reader might assume. Anyone consuming a single piece gets a valid ordered run; anyone reading the whole output in piece order gets interleaved key ranges. It is a contract problem, not a data problem, which is why it survives review so easily.
- Why not just collect the result onto one worker and sort it there?Because that worker's memory, local disk and single-threaded sort become the ceiling, and the entire result crosses the network into one place. It is the right answer when the result is small — a report of a few thousand rows — and the wrong one as soon as the ordered result is a large fraction of the input.
- Does a hash-based routing function guarantee anything useful about placement?Yes: equal keys hash alike, so all records sharing a key land on the same worker. That is enough for grouping, for folding by key, and for ordering records inside each key. It is not enough for an order across keys, which needs each piece to own a contiguous span of the range.
saying these in an interview costs you the question
- Says sorting every piece in parallel is enough, because concatenation then gives an ordered result.
- Believes a routing function that balances load well also preserves the order of keys.
- Thinks the only way to get an ordered result is to pull everything onto one machine.
- Assumes the sort step can move a record into a different piece to fix the order.
- Treats locally ordered pieces as corrupt data rather than a weaker ordering contract.