skip to content

The Quiet Source

One partition of the input goes quiet, stops advancing time, and every other partition's results stop with it. Forcing it forward unblocks them and makes its own records late on return.

on this pageshow

questions

3

A continuous job's hourly group results stopped at 03:00, yet most of its parallel input partitions are still delivering records — why?

level: middleimportance: should knowfreq 50%

answer

  1. output frozen, workers not idle
  2. time advances at the slowest input
  3. a claim needs records to move
  4. the operator takes the minimum
  5. silence pins the minimum in place

basics

~20 s

Time in the job advances at its slowest input. A partition that sends nothing never moves its own claim that nothing older is still coming, and an operator may assert only the minimum across inputs, so busy partitions' groups stay open.

solid answer

~50 s

The job closes a group on a **completeness claim** — a timestamp carried along with the records asserting that nothing older than it is still expected to arrive; the plain word is a *watermark*. That claim is derived per input from the newest moment stamped in the payloads of the records that input has delivered, held back by whatever wait the pipeline chose. An operator fed by several parallel inputs may only assert what is true of all of them, so it takes the minimum. A partition that has delivered nothing since 01:02 has nothing to derive a newer claim from, so its claim sits still, the minimum sits still, and every other partition's groups stay open although its workers are perfectly busy. The symptom is frozen *output* with healthy consumption: no backlog, no failing worker, just nothing being emitted.

code

json · 12 lines
json
{
  "as_of": "03:13:00",
  "wait_behind_newest": "00:02:00",
  "inputs": [
    { "partition": "p0", "newest_occurrence_time": "03:12:40", "claim": "03:10:40" },
    { "partition": "p1", "newest_occurrence_time": "03:11:05", "claim": "03:09:05" },
    { "partition": "p2", "newest_occurrence_time": "03:12:58", "claim": "03:10:58" },
    { "partition": "p3", "newest_occurrence_time": "01:02:00", "claim": "01:00:00" }
  ],
  "operator_claim": "01:00:00",
  "groups_still_open": ["01:00-02:00", "02:00-03:00"]
}

go deeper

for a junior

Remember that grouping by when things happened means the job must decide a moment is finished, and it decides that from the records that have arrived — so an input that sends nothing supplies no evidence that time has moved on.

for a middle

Explain the mechanics: each parallel input carries its own claim, the operator asserts the minimum so it is never wrong, and a silent input pins that minimum. Be able to say why consumption continues while emission stops.

for a senior

Diagnose it from the evidence rather than guessing: backlog near zero, workers busy, asserted time falling behind the clock at one second per second, and one input's claim unmoved. Note the retained memory growing while groups stay open.

for a principal

The judgment call is whether a shared job should ever mix always-on inputs with inputs that go quiet, since one tenant's silence becomes everyone's latency — an isolation decision, not a tuning one.

## What the job is actually waiting for A job over an endless input has no natural end to any group. To close an hourly aggregate it needs an assertion, and the assertion it uses is the **completeness claim**: a timestamp carried alongside the records as the job runs, meaning *nothing older than this is still expected to arrive*. The ordinary word for it is a **watermark**. Two things about it matter here: - It is a **running guess**, not a fact read off the data. Nothing in the input states that a moment is finished. - It is not the stored high-watermark column a scheduled extract keeps between runs to remember where it stopped. Same word, different mechanism. The claim is built from **occurrence time** — the moment stamped in the payload by whatever produced the record (**event time**), as opposed to the wall clock of the worker that eventually reached it. A job that never assigned occurrence times is silently grouping by arrival and does not have this problem, because it also does not have correct numbers. ## Why the slowest input governs The job reads several **parallel inputs** — partitions of the source, each an ordered sub-stream with its own records and therefore its own claim. Downstream of them sits an operator that must assert one claim for everything it has seen. It may only assert what is true of **all** its inputs: | input | newest occurrence time seen | its own claim | can the operator claim it? | |---|---|---|---| | p0 | 03:12 | 03:10 | no — p3 may still send 01:30 | | p1 | 03:11 | 03:09 | no — same reason | | p2 | 03:12 | 03:10 | no — same reason | | p3 | 01:02 | 01:00 | yes — this is the minimum | Taking the **minimum** is what makes the claim sound: if the operator claimed 03:10 while p3 could still deliver 01:30, it would close the 01:00 group and then discover it was wrong. The minimum is the rule every design in this class follows; what varies between engines is how often it is recomputed and whether the runtime offers any way to take an input out of the minimum. ## The failure mode: an input that sends nothing A claim derived from arriving records cannot move without arriving records. Ordinary reasons an input goes quiet: - an overnight region whose users are asleep while the rest of the world is awake; - a low-traffic tenant sharing a pipeline with busy ones; - a producer paused for a deployment, or one whose keys happen to land on no partition this hour; - a partition that exists only because the source was widened in advance of traffic that has not arrived yet. None of these is a fault. The partition is healthy, it simply has nothing to say — and that silence pins the whole job's asserted time to the last thing it said. ## Telling it apart from what it resembles The distinguishing evidence is that **output is frozen while consumption is not**: 1. **Is anything being read?** If the job is consuming at its usual rate and the unprocessed backlog is near zero, the job is not behind; it is waiting. 2. **Is one piece of work still running?** A single long-running unit while the rest finished is uneven work, not a silent input — the other pieces would have emitted. 3. **Is the gap between the asserted time and now growing steadily, at the rate of the clock?** That is the signature: it grows one second per second, because nothing advances it at all. 4. **Which input is the minimum coming from?** The per-input claims are the only reading that identifies the culprit, and the answer is the one that has not moved. Retained memory keeps growing while this lasts, because every open group is still holding its accumulators — the job gets more expensive the longer the silence runs. ## What varies between engines - **Where the claim may advance.** A record-at-a-time runtime can move it between any two records. A runtime that runs continuous work as a rapid succession of small finite jobs can only move it at the boundary of each of those, so the stall appears in steps rather than smoothly. An engine that only ever makes a single pass over a finished bounded input has no running claim at all: it reaches the end of the input and everything closes, so this failure cannot occur there. - **Where the minimum is taken.** If one reader multiplexes several source partitions, there is an inner minimum inside that reader before the operator's own minimum. A silent partition can therefore stall a job even when the reader itself is visibly emitting records from its other partitions. - **What you are given to see.** Some runtimes surface a per-input claim you can read directly; on others you infer it from the gap between the newest occurrence time and the point at which groups stop closing.

  • Is anything lost while the claim is stalled, or is it only slow?
    Nothing is dropped yet: records keep being consumed and folded into their groups. What you pay is latency on every result and memory, because each open group still holds its accumulators and nothing is releasing them. The loss appears later, if you force time forward and the quiet input then delivers records behind the claim.
  • Can this happen in a single pass over a finished, bounded input?
    No. A bounded pass has an end: when an input is exhausted its time can be advanced to the end of time and everything closes. The stall is specific to an endless input, where the claim is derived from records that are still arriving, so an input that stops delivering also stops supplying evidence.

A walking group marks a checkpoint passed only when every member has called it in, because until then someone may still be behind. One member quietly went home at the first stop and calls nothing in. Everyone else is walking fast and reaching checkpoints, but the leader's board still shows the group at the first stop — and if the leader eventually marks that member absent to move the board on, that member's own times are retroactively wrong.

saying these in an interview costs you the question

  • Says an empty partition cannot affect other partitions' results
  • Blames a slow worker without checking whether anything is emitting
  • Thinks the job is behind, when its backlog is zero
  • Says the operator should take the newest claim across inputs
  • Assumes the runtime supplies timestamps rather than the pipeline assigning them
  • Confuses it with the stored bookmark column a scheduled extract keeps
open as a page

A job stops counting a partition silent for ten minutes toward the minimum that advances its completeness claim. What does that buy, and what does it cost when the partition wakes?

level: seniorimportance: should knowfreq 42%

basics

~20 s

Excluding the silent input lets the minimum advance from the remaining inputs, so frozen groups close and latency returns to normal. The price is exact: its own records are late by construction the moment it wakes.

open as a page

A regional feed produces only between 22:00 and 06:00 and shares one job with always-on feeds. How would you design time advancement for both?

level: principalimportance: nice to knowfreq 30%

basics

~20 s

Treat it as a placement decision, not a setting: isolate the quiet feed in its own job, have its producer emit periodic filler records, force it forward on a timeout, or accept the stall — each bills a different owner.

open as a page