skip to content

Two requests resemble end-to-end ordering — the 100 largest values overall, and each key's records in time order; does either need range boundaries?

level: middleimportance: should knowfreq 45%

answer

  1. ask who reads the output, and how
  2. a bounded answer needs no global order
  3. equal keys together is not ordering
  4. top few: keep k per piece, then merge

basics

~20 s

Neither does. The top few are found by keeping the best 100 inside each piece and merging those candidates, and ordering inside a key only needs equal keys to land together, which any routing on that key already gives. Range boundaries serve a complete result read in order.

solid answer

~40 s

Both look like sorting and neither needs an ordered result. For the top few: each piece keeps only its own best 100 records, so the movement carries pieces multiplied by 100 candidates rather than the input, and one final pass over those candidates gives the answer — no sample, no order-aware routing, no full redistribution of the data. For ordering inside each key: the requirement is a relationship between records that share a key, so all you need is that equal keys reach the same worker, which routing on that key gives you whether it hashes or uses boundaries; each worker then orders each group by the second field locally. Range boundaries earn their cost only when a consumer reads the whole result in order.

code

python · 20 lines
python
import heapq
from itertools import chain


def amount(record):
    return record["amount"]


def local_top(records, k):
    # runs on every piece, in parallel; holds at most k records at a time
    return heapq.nlargest(k, records, key=amount)


def final_top(candidates_per_piece, k):
    # only pieces * k candidates ever cross the network
    return heapq.nlargest(k, chain.from_iterable(candidates_per_piece), key=amount)


# 1000 pieces, k = 100 -> 100_000 candidates compared at the end,
# whatever the size of the input the pieces were read from.

go deeper

for a junior

Recall that the largest few records can be found by keeping the best few inside each piece and then comparing only those, so almost nothing has to move across the network.

for a middle

Explain why ordering within a key needs only that equal keys meet on one worker, and why a routing function derived from a hash of that key already satisfies it.

for a senior

Probe the actual consumer before paying for an order-aware movement, and spot the variants that break the bounded shape: a count that rivals the input, a per-key limit, or a threshold whose result size is data-dependent.

for a principal

Treat ordering as a contract negotiated with consumers rather than a default; a platform that sorts every output end to end is paying a sample and a full movement on every run to serve readers who mostly take a handful of rows.

## The expensive thing you are being asked to avoid Producing an end-to-end ordered result costs a sample of the ordering key, a set of range boundaries derived from it, a full **redistribution** — every worker sending each record it holds to whichever worker will handle that record's key — and a local sort on every piece. That is the right price when the deliverable is a complete result that someone reads in order. Two very common requests look like that request and are not, and recognising them is most of the value of understanding the mechanism. ## Lookalike one: the top few overall "The 100 largest amounts across ten terabytes" sounds like ordering ten terabytes and taking a prefix. It is not, because the answer is bounded and tiny: - each piece keeps only its own best 100 records, using a bounded structure that never holds more than 100 at a time; - what crosses the network is at most `pieces x 100` candidate records — for a thousand pieces, a hundred thousand records, not ten terabytes; - one final pass over those candidates produces the answer, and the correctness argument is simple: any record in the global top 100 is in its own piece's top 100, so discarding the rest can never lose it. The arithmetic only degenerates when the count asked for stops being small relative to the input — a "top few" of the same order as the data is just an ordered result with extra steps. There is also a genuine variation to be aware of: the local step is a pure per-piece reduction and can often be folded into the reading of the piece, whereas the final pass is a single small collection point. Whether a runtime expresses that as one step or two differs; the shape does not. ## Lookalike two: order inside each key "Every event for each account, in time order" is a statement about records that share a key, not about records across different keys. Account A's events never have to be comparable with account B's. So the requirement decomposes into: 1. **all records for one key reach the same worker** — which any routing function on that key already guarantees, including one derived from a hash, because equal keys are routed alike; 2. **each worker orders each group by the second field** — a local sort over one group's records, entirely within one unit of work. No sample, no boundaries, no relationship between piece index and key order. What varies between runtimes is only *how* the ordering inside the group is obtained: some can arrange for a worker's input to arrive already ordered within the group, others sort in the operator itself, and the memory cost differs accordingly, especially for a key with very many records. The conceptual point is unchanged. ## The three requests side by side | request | what it needs | what crosses the network | range boundaries? | |---|---|---|---| | the 100 largest values overall | a bounded local best-of, then one merge | pieces x 100 candidates | no | | each key's records in time order | equal keys together, then a local sort per group | every record once, routed on the key | no | | the complete result, read in order | contiguous spans per piece, then local sorts | every record once, routed by span | yes | ## The question to ask before paying for an order When someone asks for sorted output, the useful probe is: **who reads it, and do they read all of it in sequence?** If the consumer takes a fixed number of extreme records, you want the bounded local best-of. If the consumer works one key at a time, you want equal keys together and nothing more. Only a consumer that streams the whole result in sequence — a file another system ingests in order, a report printed end to end, a merge against another ordered dataset — actually depends on an end-to-end order, and only then is the sample plus the order-aware movement money well spent. One caution on the bounded-best-of shape: it depends on the answer being a fixed count of records. A request for the top few *per key* is again a different thing — it is the second lookalike with a per-group limit added — and a request for values above a threshold has an answer whose size is data-dependent, so the bound the local step relies on disappears.

  • Why is keeping the best 100 inside each piece guaranteed to find the global top 100?
    Because any record in the global top 100 must also be in the top 100 of the single piece that holds it — nothing in that piece can outrank it more than 99 times without outranking it globally too. So discarding everything below each piece's hundredth-best can never discard a member of the answer.
  • When does the keep-the-best-k shape stop paying off?
    When the count stops being small relative to the input. If each of a thousand pieces must contribute a million candidates, the movement is the same order as the data and you have built an awkward ordered result. It also breaks for a threshold rather than a count, because then the answer's size is data-dependent and the local step has no bound to work with.
  • What about the top few records per key rather than overall?
    That is the second lookalike with a limit attached: equal keys still only need to meet on one worker, and each worker then keeps the best few within each group. No comparison across keys is involved, so there is still no need for contiguous spans or a sample of the key range.

saying these in an interview costs you the question

  • Sorts the whole input to answer a request for the largest hundred records.
  • Thinks ordering inside each key requires each piece to hold a contiguous span.
  • Believes routing on a hash cannot support ordering records within a key.
  • Applies the keep-the-best-k trick when the count asked for rivals the input size.
  • Treats any request containing the word sorted as a request for an end-to-end order.