skip to content

Operators

A pipeline is built from small composable steps, each making its own promise about ordering and concurrency. Interviewers ask you to choose between two for a case, since that decides correctness.

on this pageshow

explore

questions

23

A checkout display joins a scan feed and a price feed: why does lockstep pairing emit fewer results than latest-value combination?

level: juniorimportance: must knowfreq 58%

answer

  1. two ways to wait for a partner
  2. one consumes values, the other reuses them
  3. queue per source versus slot per source
  4. pairing runs at the slowest source's rate
  5. latest-value emits on every arrival

basics

~20 s

Lockstep pairing consumes one value from each source per result, so its rate is the slowest source's rate. Latest-value combination keeps each source's most recent value and emits on every arrival, so its rate is the sum of all rates.

solid answer

~40 s

The two rules differ in what they do with a value once they have seen it. **Lockstep pairing** keeps a queue per source: every arrival is appended, and a result is emitted only when every queue is non-empty, taking the front value of each. Each input value is consumed by exactly one output, so the output rate is bounded by the slowest source and the faster source's surplus sits buffered. **Latest-value combination** keeps a slot per source: every arrival overwrites its slot and, once every slot has been filled at least once, emits a result built from that arrival plus the other sources' current values. One value can appear in many results, or in none if it is overwritten first. Neither emits anything until every source has produced at least one value.

code

pseudocode · 12 lines
pseudocode
// lockstep pairing: one queue per source
on value(v) from source i:
    queue[i].append(v)
    if every queue is non-empty:
        emit combine( take_front(queue[0]), take_front(queue[1]) )

// latest-value combination: one slot per source
on value(v) from source i:
    latest[i] = v
    filled[i] = true
    if every filled[k] is true:
        emit combine( latest[0], latest[1] )

go deeper

for a junior

Recall the two rules by their bookkeeping: a queue per source that is consumed, versus a slot per source that is overwritten. That one image gives you the output counts without memorising anything else.

for a middle

Explain the rates: pairing is bounded by the slowest source, latest-value runs at the sum of all rates. Then say what happens to the surplus and why neither emits during the warm-up.

for a senior

Show the consequence downstream. Latest-value output multiplies every write, call or render; pairing accumulates an invisible backlog and pairs against stale partners. Pick the rule from what the consumer of the output must be true.

for a principal

The choice is a contract on the output: is a result an event that happened once, or the current state of the world? Setting that contract per surface, and requiring an explicit key when correctness depends on which value belongs to which, is what keeps teams from arguing about it per feature.

## Two rules, two different promises A join takes several sources that each push values over time and produces one output source. The interesting part is not that values arrive from more than one place — it is the rule the join uses to decide **when it has enough to emit**, and **what it does with the values it did not use**. - **Lockstep pairing** keeps one **queue** per source. Every arriving value is appended to its own queue. The moment every queue is non-empty, the join removes the front value of each queue, combines them, and emits one result. Every input value is consumed by exactly one output. - **Latest-value combination** keeps one **slot** per source. Every arriving value overwrites its slot. Once every slot has been filled at least once, *every* arrival emits a result built from that arrival plus the other sources' current slot values. One input value can appear in many results — or in none, if a newer value overwrites it before any other source moves. That single difference — consume versus reuse — produces every other difference between them. ## Why the output counts differ 1. **Pairing consumes, combination reuses.** A paired result costs one value from each source; a latest-value result costs one arrival from any source. 2. **Pairing runs at the slowest source's rate.** For two sources it emits `min(count of A, count of B)` results, and the surplus from the faster source waits in its queue for a partner that may never arrive. 3. **Combination runs at the combined rate.** After the warm-up, it emits once per arrival across all sources — roughly the sum of the input rates. The one thing both rules share is the warm-up: **neither emits anything until every source has produced at least one value**. A source that stays silent keeps the whole output silent. (With exactly one value per source the two rules produce the same single result — pairing is never *more* talkative than latest-value combination, and it is quieter as soon as any source outpaces another.) ## A worked arrival trace The terminal receives, in this order: a price update `P1`, a scan `S1`, prices `P2` and `P3`, a scan `S2`, price `P4`, scan `S3`, price `P5`. | arrival | lockstep pairing | latest-value combination | |---|---|---| | `P1` | queued — no scan yet | nothing — the scan slot is empty | | `S1` | emits `(S1, P1)` | emits `(S1, P1)` | | `P2` | queued | emits `(S1, P2)` | | `P3` | queued | emits `(S1, P3)` | | `S2` | emits `(S2, P2)` | emits `(S2, P3)` | | `P4` | queued | emits `(S2, P4)` | | `S3` | emits `(S3, P3)` | emits `(S3, P4)` | | `P5` | queued | emits `(S3, P5)` | | **total** | **3 results**, 2 prices left buffered | **7 results** | Notice what the trace exposes beyond the counts. Pairing gave `S3` the price `P3` — three ticks stale — because it pairs by **position in arrival order**, not by time or identity. Latest-value combination gave the freshest price every time, but emitted the same scan three times over. ## Choosing the rule for the surface - **A display that must refresh whenever any input changes** wants latest-value combination: a new price should redraw the total even though nothing was scanned. - **A per-event record where each input value must appear exactly once** wants pairing: one receipt line per scan, never three. - **A downstream step that writes, charges or sends per result** wants pairing, or an explicit rate-reducing step first — latest-value output count is not the event count, and treating it as one duplicates work. - **Neither rule matches an item to its own price** unless the two feeds are guaranteed to move in lockstep by construction. When correctness depends on which price belongs to which scan, carry a shared key in both values instead of relying on order. ## What each rule hides - **Pairing hides the unused tail.** The two buffered prices in the trace never reach the output and nothing reports them. If one source is permanently faster, that tail grows for as long as the pipeline runs. - **Pairing hides staleness.** The partner it pairs with is the oldest unused one, not the current one. - **Latest-value hides write amplification.** Five price ticks turned three scans into seven results; a downstream persist or network call now runs seven times. - **Latest-value hides non-co-occurrence.** A result may combine values that were never simultaneously true of the world — the newest scan with a price that arrived a moment later. - **Both hide the warm-up.** Until every source has spoken once, the output is indistinguishable from a hang.

  • With five sources instead of two, how does each rule's output rate change?
    Lockstep pairing still emits at the slowest source's rate — one result per complete set of five — so adding sources can only slow it. Latest-value combination emits once per arrival across all five, so its rate is roughly the sum of the five rates and grows with every source added.
  • Can either rule emit a result before every source has produced a value?
    No. Pairing needs one unused value in every queue; latest-value combination needs every slot filled at least once. Until then the output is silent, which looks exactly like a hang. If one input is genuinely optional, give it a starting value before the join so its slot is never empty.
  • Under latest-value combination, can a value never appear in any result?
    Yes. If the same source produces a newer value before any other source moves, the older value is overwritten in its slot and is never combined with anything. Latest-value combination makes no promise that every input value is represented in the output; only pairing does that.

Pairing is a two-part form: each half is filed once, and a spare half waits in the tray for its match. Latest-value combination is a scoreboard: any single change redraws the whole board from the current numbers.

saying these in an interview costs you the question

  • Says both joining rules produce the same number of results
  • Thinks latest-value combination uses each value exactly once
  • Believes lockstep pairing discards the faster source's surplus values
  • Expects a first result before every source has produced a value
  • Treats every latest-value result as a matched pair of simultaneous events
  • Assumes the paired partner is always the most recent value
open as a page

In a stream pipeline, how do a mapping step, a predicate filter, and a running-total accumulator differ?

level: juniorimportance: must knowfreq 70%

basics

~20 s

A mapping step returns exactly one output element per input; a predicate filter returns zero or one, dropping the rest; a running accumulator emits a value derived from every element seen so far, so it carries state between elements.

open as a page

A stream pipeline assembled once at start-up writes a log line immediately, before any run — why?

level: middleimportance: must knowfreq 62%

basics

~20 s

Building the chain runs ordinary code. The expressions handed to each step are evaluated as the chain is described, so a log line or a computed value written there fires once at assembly, not on each later run of the pipeline.

open as a page

In a two-source join, what happens to the output when one source completes early or emits nothing at all?

level: middleimportance: must knowfreq 62%

basics

~20 s

It depends on the joining rule. Lockstep pairing ends as soon as a completed source's queue is empty, discarding whatever is buffered elsewhere. Latest-value combination keeps going on the completed source's last value until all sources complete. A source that emits nothing leaves the output empty.

open as a page

A moderation queue classifies each arriving item with an external call: what does concatenating those calls sequentially cost against merging them concurrently?

level: middleimportance: must knowfreq 72%

basics

~20 s

Sequential concatenation keeps one call in flight: results arrive in queue order, throughput is capped at one call per call-latency. Bounded merging runs several at once, trading that ordering for throughput and putting the bound's worth of load on the dependency.

open as a page

Mapping each queued item to its own asynchronous call yields a source of sources: what does a flattening step then do?

level: middleimportance: must knowfreq 64%

basics

~20 s

A flattening step subscribes to each inner source the mapping produced and re-emits their elements onto one outer sequence. Its strategy also fixes the ordering of results, how many inner calls run at once, and whether superseded ones are cancelled.

open as a page

In a stream pipeline, which clock-driven operator suits a dial turned in bursts, and which suits a continuously drifting reading?

level: middleimportance: must knowfreq 62%

basics

~20 s

A quiet-period debounce fits the dial: it emits one settled value once the turning stops. Periodic sampling fits the drifting reading: it emits whatever is current on each tick, because that source never falls quiet.

open as a page

A nightly report pipeline stamps the same date on every run — how do you make each run recompute it?

level: seniorimportance: must knowfreq 56%

basics

~20 s

Replace the baked-in value with a factory the pipeline calls when a run starts: a step that accepts a function returning a source defers construction to execution, so each run computes its own date instead of reusing the one captured during assembly.

open as a page

In a payroll stream pipeline, what goes wrong when the mapping step also writes an audit row for each element?

level: seniorimportance: must knowfreq 58%

basics

~20 s

The write becomes invisible in the chain and repeats with the run — once per element per subscription, again on resubscription, and a mid-stream failure leaves earlier writes applied. Keep the transform computing only and declare the effect separately.

open as a page

A filtering step is added to an assembled stream pipeline, yet every value still arrives — why?

level: juniorimportance: should knowfreq 48%

basics

~20 s

Each pipeline step returns a new source that wraps the old one; it never modifies the source it was applied to. Discarding that returned value throws the filtering step away and leaves the original, unfiltered chain in the variable.

open as a page

Why must a stream step that carries a running year-to-date total create its state per subscription?

level: middleimportance: should knowfreq 46%

basics

~20 s

Each subscription is an independent run of the chain. State created once when the chain was built is shared by every subscriber, so one run's totals leak into another; state created per subscription starts from the seed and stays private to that run.

open as a page

In a stream pipeline, why does a positional take of the first n elements stop the source while a predicate filter does not?

level: middleimportance: should knowfreq 40%

basics

~20 s

A positional take counts what passes, so after the nth element it knows nothing further is needed: it cancels upstream and completes downstream. A predicate filter judges each element alone, never learns it is finished, and leaves the source producing.

open as a page

Why does a stage that groups stream values into fixed-size batches usually also need an elapsed-time trigger?

level: middleimportance: should knowfreq 50%

basics

~20 s

A count-only group holds values until the count is reached, so when the source slows the partial group waits indefinitely and its oldest value goes stale. An elapsed-time trigger closes the group on age as well, bounding that wait.

open as a page

A price lookup is sent to three redundant sources and the first answer wins; what must happen to the two losers?

level: seniorimportance: should knowfreq 40%

basics

~20 s

All three are subscribed at once; the moment one responds, the selection cancels the other two and delivers only the winner's values. Cancellation is a request the losing sources must honour, so work already in flight, and any side effect it has, may still complete.

open as a page

A lockstep pairing of scan events and weight readings loses one weight reading; what do all the later pairs look like?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Every later pair is shifted by one: each scan is combined with the next item's weight. Pairing matches by position in arrival order, not by identity, so the results stay well formed and no failure is signalled — the terminal validates each item against its neighbour's weight forever.

open as a page

Concurrently merged classification results reach an ordered audit log out of queue order: how would you restore the required ordering?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Four routes, priced differently: drop to sequential concatenation and lose throughput; use an order-preserving concurrent flatten that holds results and releases them in subscription order; attach a sequence number and resequence downstream; or narrow the requirement to per-key ordering.

open as a page

A detail pane loads whichever queue item is highlighted: what breaks if those loads are merged instead of switched to the newest?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Merging leaves every earlier load in flight, so whichever finishes last wins the pane: a stale load can overwrite the highlighted item's data. Switch-to-newest cancels the in-flight load on each highlight change, leaving at most one alive.

open as a page

What makes a quiet-period stage emit nothing for an hour while its dial source is in constant use?

level: seniorimportance: should knowfreq 44%

basics

~20 s

Every arrival restarts the quiet timer, so while the gap between values stays below the window the timer is cancelled before it ever fires. The stage is starved, not slow: it emits only once the source finally falls quiet.

open as a page

Several pipelines fan out to one shared classification service: how would you decide the concurrency bound each flattening step uses?

level: principalimportance: should knowfreq 34%

basics

~20 s

Size it from the dependency's safe concurrency rather than each pipeline's appetite: divide one fleet-wide budget across callers with headroom, check it against arrival rate times call latency, and keep bounds per pipeline so a single backfill cannot consume the service.

open as a page

You must set one clock-driven reduction rule for hundreds of heterogeneous panel feeds - how would you decide and defend it?

level: principalimportance: should knowfreq 36%

basics

~20 s

Decide from measured gap distributions and the consumer's freshness need, not from taste. Classify feeds into bursty and continuous, give each class a default reduction, and publish a staleness bound as the contract so the operator choice can change later.

open as a page

When one source is appended after another, when is the second source subscribed, and what does that delay cost?

level: middleimportance: nice to knowfreq 30%

basics

~20 s

Only after the first source completes. Appending buys a strict ordering guarantee and pays for it with a gap: a second source that is already running loses everything it produced during that window, a first source that never completes starves it, and a first source that fails prevents it from running at all.

open as a page

When a stream pipeline fails at run time, why does the stack trace rarely name the line that declared the failing step?

level: seniorimportance: nice to knowfreq 34%

basics

~20 s

The statements that declared the chain ran during assembly and their frames returned long before anything executed. At failure the stack holds only the machinery currently delivering values and the bodies it invoked, so the declaring line is nowhere on it.

open as a page

How does a stage that fails when the gap between consecutive values grows too long differ from one bounding the whole sequence?

level: seniorimportance: nice to knowfreq 30%

basics

~20 s

A per-gap bound restarts on every value, so a slow but steady source never trips it and it fires only during silence. A whole-sequence bound fires at a fixed age regardless of progress, even while values are still flowing.

open as a page