skip to content

An operator reads eight parallel inputs; seven assert nothing older than 10:00 is still coming, one only 09:00. What may the operator assert?

level: seniorimportance: should knowfreq 58%

answer

  1. each parallel input carries its own value
  2. the assertion must hold for all of them
  3. one summary rule survives; the others lie
  4. healthy throughput, frozen output
  5. read the values side by side

basics

~20 s

09:00 - the minimum. An operator may assert only what its least advanced input supports, because that input can still deliver records older than 10:00. One lagging input therefore holds the whole job's asserted time back.

solid answer

~50 s

Each of the job's parallel inputs - a partition of the source, or one upstream operator instance - carries its own completeness claim, and the receiving operator takes the **minimum** across them. Any other rule is unsound: claiming 10:00 would assert that nothing older than 10:00 is still coming, while an input still standing at 09:00 is free to deliver a 09:30 record a moment later. The rule composes along the whole path, so a job's effective asserted time is the minimum over every route from source to operator, and it can never become more advanced as it travels downstream. The production symptom is distinctive: throughput matches the input rate on every partition, workers are busy, nothing is backing up - and results stop appearing, because one input's claim is the constraint. Diagnose it by reading per-input claim values against each other, not the aggregate.

go deeper

for a junior

Remember the rule and the reason: an operator takes the smallest of its inputs' claims, because any input still behind can deliver an older record at any moment.

for a middle

Explain how the minimum composes through a redistribution, and why the claim can only fall as it travels downstream rather than rise.

for a senior

Recognise the signature in production - every partition keeping pace, workers busy, time-based results frozen - and separate it from a backlog before you touch capacity.

for a principal

Note the coupling this creates: joining a slow source into a fast pipeline gives the slow source control of everyone's freshness, which is an ownership question as much as a technical one.

## Why minimum, and not anything else A completeness claim is an assertion with a precise content: no record with a moment older than this is still expected. An operator fed by several parallel inputs may assert only what all of them support, because any one of them can still deliver. | candidate rule | what it would assert | why it fails | |---|---|---| | the minimum | nothing older than the least advanced input's value is coming | sound: every input has already passed that point | | the maximum | nothing older than the most advanced input's value is coming | the lagging input may deliver a record older than that within a second | | the average | a value no single input supports | unsound for the same reason as the maximum, and harder to reason about | | weighted by record rate | the busy inputs decide | a quiet input is not a finished input; volume is unrelated to which moments remain | The minimum is therefore not a tuning choice; it is the only rule that keeps the assertion true. ## How it composes through a job The combination happens at every operator instance, over every upstream instance that can send to it: 1. A source-facing step assigns each record its moment and produces a claim per parallel input. 2. A step that consumes records without moving them between workers carries each input's claim forward. 3. A step that redistributes records across the network gives every receiving instance many upstream instances, and that instance holds the minimum over all of them. Because each stage takes a minimum, the claim can only fall as it travels, never rise. The asserted time of a result at the end of a job is the minimum over every path that feeds it. In a wide job with several sources joined together, a single laggard anywhere upstream governs everything downstream of the point where it joins - including branches that have nothing to do with it. ## What this looks like in production The symptom is easy to misread because every health signal is green: - Consumption keeps pace with the input rate on every partition; nothing is piling up unprocessed. - Workers are busy; the cluster is not short of capacity. - Results for time-based groups stop appearing, or appear in an old and unmoving range, even though the records are clearly flowing. That combination - caught up on volume, stalled on time - is close to diagnostic. A backlog looks different: records accumulate unprocessed, consumption is below the input rate, and the claim is behind simply because the job has not yet reached the recent records at all. The two feel identical on a dashboard that shows only 'how far behind are we', which is why the useful view is **per-input claim values side by side**, not one aggregate number. What makes an input lag is separate from this mechanism and usually mundane: it reads a source region with a slower producer, it was assigned more work, or its producer buffers before sending. A lagging input is still producing records; an input that has gone entirely silent is a different situation with different remedies. ## What varies between designs - **Granularity of the combination.** Some runtimes track a claim per input channel and recombine per operator instance, which is what makes one laggard's effect precisely traceable. A runtime that executes continuous work as repeated small finite jobs typically computes one value per small job across the whole input, so the same minimum is taken but you cannot see which input caused it in the same way. - **Whether the job pushes back on a runaway input.** Some engines can slow a far-ahead input so it stops accumulating retained data for periods the lagging input has not reached. Others do nothing, and the memory cost of waiting falls where it falls. - **Where the value is observable.** Some expose the combined value only, some expose it per input. This decides whether the diagnosis above takes minutes or a day. ## What an interviewer wants to hear Say 'the minimum, because the assertion has to hold for every input that can still deliver', then give the consequence - one lagging input governs the job - and finish with the diagnosis: compare per-input claims against each other, and separate the caught-up-but-stalled case from a genuine backlog, because the fixes have nothing in common.

  • Where is the minimum recomputed after records are redistributed across the network?
    At each receiving operator instance, over every upstream instance that can send to it. A redistribution gives every downstream instance many upstreams, so it holds the minimum across all of them. Because each stage takes a minimum, the claim never becomes more advanced downstream - a job's asserted time is the minimum over every path feeding the result.
  • How do you tell a claim held back by one input from a job that is simply behind on volume?
    Compare per-input claim values against each other and against consumption. If every input is keeping pace with its input rate and nothing is piling up, yet time-based results are frozen, one input's claim is the constraint. If records are accumulating unprocessed and consumption is below the input rate, that is a backlog and a capacity problem, and the claim is behind only as a consequence.
  • What does waiting on a lagging input cost the workers that are ahead?
    Memory, mostly. Periods that cannot be declared finished stay open, and whatever the job retains for them keeps accumulating on every worker, including the ones far ahead. Some engines can slow a runaway input so the spread stops widening; where they cannot, a long-running laggard turns a time problem into a memory problem.

A convoy is only as far along as its slowest ship. The fleet's position, for any purpose that depends on all of it having passed a point, is the slowest vessel's position - not the average, and not the leader's. Speeding up the leaders changes nothing about that number.

saying these in an interview costs you the question

  • Takes the newest of the inputs' claims so the job keeps moving
  • Averages the inputs' claims to get the operator's value
  • Weights the combination by each input's record rate
  • Assumes frozen output means the cluster is short of capacity
  • Thinks an operator can advance beyond what its inputs assert
  • Expects more parallelism on the operator to advance the claim