Two endless inputs are joined on a shared key. Why must the job hold records from both inputs, not just one?
answer
- neither input ever ends
- the partner may not exist yet
- both sides are held, not one
- a bound on the moments' difference
- dropped when no match remains possible
basics
~20 sNeither input ends, so a record on either side may meet its partner later. Both sides are therefore held while the other is awaited, and a bound on how far apart their two moments may be is what lets a held record ever be dropped.
solid answer
~50 sA join over two finite inputs can read one side to its end, build a lookup from it, then pass the other side over that lookup. Neither endless input has an end, so that plan never gets past its first step. At the instant a record arrives on the left, its partner on the right may simply not have been produced yet — and the same is true with the sides swapped, which is why both are held rather than one. What the job keeps is every record that could still take part in a match, grouped by the join key. Nothing in the join itself says when one of those records may be discarded, so the retained set grows until the predicate supplies a limit: a bound on how far apart the two records' moments may be. Once that difference can no longer be satisfied, the record can never match and is dropped.
go deeper
Recall the core fact: two endless inputs never line up, so records from both sides are held while their partners are awaited, and a limit on how far apart the two moments may be is what lets anything be dropped.
Explain the mechanics: why build-then-probe needs an input that ends, what precisely is held and grouped by which key, and why a limit on the moments' difference gives a held record an expiry.
Show the production consequence: the retained set is arrival rate times holding time on each side, an asymmetric bound holds the two sides for different lengths, and a predicate with no limit is an unbounded commitment made silently.
Frame the trade the bound represents: match completeness bought with memory that scales with traffic rather than with how often a distant pair actually occurs, and decide whether that belongs in a low-latency job at all.
## The plan a finite join uses, and why it cannot start here A join over two **finite** inputs has a shape that systems of this class rely on: read one input to its end, build a lookup structure keyed by the join key, then pass the other input over that structure and emit a row for every match. The plan has a precondition buried in its first word — one input must **reach an end**. An endless input never reaches one. The consequence is not that the join is slow; it is that step one never completes, and no amount of memory or parallelism changes that. A join between two endless inputs therefore cannot be build-then-probe at all. It has to be an incremental match: each arriving record is checked against whatever the job is currently holding from the other input, and is then itself held, because its own partner may not exist yet. ## Both sides are held, and not out of symmetry At the instant a record arrives on the left, exactly one of three statements is true of its partner on the right, and the job can tell them apart from none of them: - the partner arrived earlier and is already in the set held from the right input; - the partner exists but has not been produced yet; - the partner will never exist. Only the first can be acted on immediately. For the other two, the job's options are to hold the arriving record or to give up a match it could have made. Swap left and right and the argument is unchanged, which is why **both** inputs are held. "Hold the small one and stream the other past it" is a statement about how the two inputs are physically brought together on the cluster, not about what an endless-to-endless match retains; that placement choice and its cost are owned by the movement subject and are not settled here. What is held deserves precise words, because the same phrase means three different things elsewhere in this category: - **What** — every record from both inputs that could still take part in a match: the whole record, or at least the join key plus the fields the output needs. An aggregate over an interval can collapse its records into a single running value; a join cannot, because when the partner finally arrives the output row must carry the original record's fields. - **Per what** — grouped by the join key, so an arriving record's probe touches only the held records under its own key. - **For how long** — until the predicate can no longer be satisfied, which is exactly the thing the join by itself does not state. One exception is worth carrying: where the match is genuinely one-to-one and a record can pair at most once, a record may be released the moment it matches. In a many-to-many match it cannot, because a later arrival on the other side may legitimately pair with it again. ## The bound in the predicate is what ends the wait A **time-bounded match** is a join whose predicate includes a limit on how far apart the two records' moments may be — "the same order identifier, and the payment's moment within thirty minutes of the order's moment". That limit does two things, and the second is the one candidates miss: 1. It gives a held record an expiry. Once time has moved far enough that no future arrival on the other side could satisfy the limit, the record can never match again and may be dropped. 2. It converts "how much is retained" from a hope into arithmetic: arrival rate multiplied by holding time, once for each side. The moments being compared are whatever the pipeline **assigned** — a field inside the payload, a moment the source recorded, or, where nothing was assigned, the moment the job reached the record. A bound compared against the last of these is really bounding arrival order, which is a different quantity and usually not what was asked for. The job also needs to know when time has moved far enough. That comes from a running claim that no record older than a stated moment will still arrive — a **completeness claim**, the mechanism the industry commonly calls a watermark. This category relies on that claim; how it advances and what stalls it are a separate subject. ## Where real designs differ | The question | How systems of this class differ | |---|---| | How the bound is written | as a limit on the difference between the two moments, or as a shared interval that both records must fall into; the second quietly misses a pair that straddles an interval edge even when the two moments are seconds apart | | Where held records physically sit | in the worker's own memory, or in a per-worker on-disk structure; this changes throughput and restore cost, not what must be held, and the choice belongs to the retained-set subject | | How the job advances | a runtime that runs continuous work as a rapid succession of small finite jobs must carry the held records from one such job to the next; a record-at-a-time runtime keeps them inside long-lived operators; the oldest model in this family is a finite two-phase disk-to-disk pass with no endless input at all, where a match means re-running the whole pass over the next slice | | The fate of a record that never matched | dropped silently under one contract, emitted once with a null counterpart after the bound passes under another | What this reasoning does **not** settle: which input travels to which machine, what happens when one join key value dominates the records, and how far a retained set may grow in general. Each is a neighbouring subject with its own cost model.
- Does holding both sides always mean twice one side's memory?Only when the two arrival rates and the two holding times match. A bound can be asymmetric — say the right record may follow the left by up to thirty minutes but never precede it — and then the left input is held for the full bound while the right needs holding only briefly. Unequal arrival rates skew it further, so size each side separately rather than sizing one and doubling.
- Why can an aggregate over an interval keep one value per key while a join cannot?An aggregate only has to produce a summary, and sums, counts and merged summaries have the property that any two of them combine into one, so the group costs a single value. A join has to emit the record's own fields next to its partner's, which cannot be reconstructed from a merged value. That is why the retained size of a match is counted in records rather than in groups.
- What happens if the predicate carries no limit on how far apart the moments may be?Nothing ever becomes impossible to match, so nothing can be dropped on correctness grounds and the held set grows for the life of the job. The job does not fail at once; it degrades as the retained records outgrow memory, then outgrow whatever is behind it. How a retained set is bounded in general, and what else can bound it, is a neighbouring subject — but at the level of the join itself, the bound in the predicate is the reason expiry is possible.
Two friends agree to meet at a station but neither knows the other's train. Each must wait, because the other may still be arriving — waiting is not politeness, it is the only way a meeting can happen at all. The station fills up with waiting people until somebody states the rule "wait at most twenty minutes, then leave"; that rule is what caps the crowd, and how large the crowd gets is the number of people arriving per minute multiplied by those twenty minutes.
saying these in an interview costs you the question
- Says only the smaller input is held while the other streams past
- Assumes the join can simply wait for the partner to arrive
- Thinks the join completes when one of the inputs ends
- Claims a time bound alone makes the retained set small, whatever the rate
- Assumes every held record is released at its first match, whatever the join's cardinality