skip to content

Time in a Stream

Three clocks disagree: when a thing happened, when the system received it, when a worker reached it. Interviewers probe it: grouping by the convenient one produces numbers nobody notices are wrong.

on this pageshow

explore

questions

20

A streaming job advances a timestamp asserting nothing older than it is still expected. What does that assertion license?

level: juniorimportance: must knowfreq 68%

answer

  1. an endless input supplies no ending
  2. somebody has to declare the hour finished
  3. a timestamp travelling with the records
  4. compared against a group's end
  5. eligible to close, not emitted

basics

~20 s

A completeness claim - a watermark - is a timestamp the job carries with the records, asserting no record older than it is still expected. It makes any group ending at or before it eligible to close.

solid answer

~50 s

An unbounded input never ends, so nothing in the data ever reports that a period is finished. The job asserts it instead: it carries a **completeness claim** - the plain word is a watermark - a timestamp travelling alongside the records saying that no record with an older moment is still expected. Once that claim passes a group's end, the group is *eligible to close*: a time-bounded decision may now be taken on it. Eligibility is not emission - what is actually published, and whether a result is published again later, is a separate contract. The claim is also not the marker used to line up a consistent recovery picture, and not the stored bookmark column a scheduled extract keeps between runs. Where the value comes from, and how often it may advance, differ between engines.

go deeper

for a junior

Recall the one-line definition: a timestamp the job carries with the records saying nothing older is still expected. Know that an endless input supplies no ending of its own, so something has to declare a period finished.

for a middle

Explain that the value is computed from the moments the pipeline assigned, held a chosen duration behind the newest one, and that it advances rather than retreats while the job runs.

for a senior

Show where eligibility to close stops and the publishing contract begins, and keep the claim distinct from the marker used for consistent recovery. Teams conflate the two in incident reviews and reach wrong conclusions.

for a principal

Frame it as a promise to consumers: the claim is where your pipeline decides a period is finished, so it is the point at which your organisation's definition of a final number actually lives.

## The boundary an endless input never supplies A job over a finite input knows when to answer: the input ends, the last record is read, the total is final. A job over an **unbounded input** - records arriving indefinitely, with no end in sight - never receives that signal, and yet it is asked for answers about periods of time: the order count for 09:00-10:00, whether two records fell within five minutes of each other. Group those by **occurrence time** - the moment the thing being recorded actually happened, as written into the record by whatever produced it - and the answers survive a re-run unchanged. But that choice creates a question the data itself cannot answer: *is the 09:00 hour finished?* Records carrying 09:xx moments can still appear long after 10:00, because a producer buffered, a retry path was slow, or a device spent the morning out of coverage. Nothing ever arrives to say 'that was the last one'. So the job asserts it. ## The claim itself **A completeness claim - the market's plain word for it is a watermark - is a timestamp the job carries alongside the records, asserting that no record with a moment older than that timestamp is still expected.** Three properties matter: - It is **carried in band**. It travels with the flow rather than beside it, so every operator downstream sees it in position relative to the records it speaks about. - It is **computed by the job**, not read from a record. A record's moment is data; the claim is an inference drawn from the moments seen so far. - It is **advanced, and generally held non-decreasing** while the job runs, because decisions already taken on the strength of it cannot be untaken. The case where it effectively goes backwards is a restart from a saved recovery point, which resumes an older value and works forward again. ## What the claim licenses A decision bounded by time cannot be taken until something declares that period finished. The claim is that declaration: 1. A group whose end lies at or behind the claim becomes **eligible to close** - the job may treat that period as complete. 2. A timer registered against occurrence time may be considered reached. 3. A time-bounded match between two inputs may stop expecting a partner for a record, because a partner older than the claim is no longer expected. **Eligible to close is not the same as emitted.** Whether the result is published at that instant, published earlier as a speculative figure, or published again later with a new value, is a separate contract the pipeline chooses - and engines differ sharply here, some emitting once when a period is eligible, others emitting an updated running result on every input record. The claim only removes the reason to keep waiting. ## Three things get called a marker; only one is this | what | what it asserts | what it is for | |---|---|---| | the completeness claim | no record older than this moment is still expected | letting a time-bounded decision be taken at all | | the recovery snapshot marker | this point divides 'before' from 'after' in the flow | lining up one consistent restart picture across workers | | a stored bookmark value | everything up to here was already extracted | resuming the next scheduled extract where the last one stopped | The first two both travel with the data and are easy to confuse; they align by different rules and exist for unrelated reasons. The third does not travel at all - it is a value persisted between runs of a scheduled batch extract, wearing the same word. ## Where the value comes from Most pipelines build the claim from the moments they assigned to records: take the largest moment seen so far and hold the claim a chosen duration behind it - a **disorder bound**, the length of time the job agrees to wait for records that run behind. Some sources offer something better: metadata describing how far each parallel input has progressed, which lets the job derive a claim instead of inferring one from payloads. Either way the claim is only as good as the moment assignment beneath it: a pipeline that never assigned a moment is quietly asserting completeness over arrival order. ## What varies, and why your answer should say so - **When it may advance.** A runtime that pushes records through operators one at a time can raise the claim between any two records, usually damped by a short emission interval. A runtime that executes continuous work as a rapid succession of small finite jobs can normally raise it only at the boundary between one of those jobs and the next. - **Whether it exists at all.** A single pass over a finished bounded input needs no running claim: the end of the input is the boundary, and completeness is observed rather than asserted. - **How it is carried.** Some designs emit the claim on a periodic schedule; others attach it to distinguished records in the flow. State the model first - an asserted, advancing, in-band timestamp that licenses time-bounded decisions - and then name which mechanics you are describing.

  • Can a job's completeness claim move backwards while it runs?
    Engines generally hold it non-decreasing per operator: groups already treated as complete cannot be un-completed, so a lower value would make earlier decisions retrospectively wrong. The case where it effectively rewinds is a restart from a saved recovery point, which resumes the older value carried in that point and advances forward again over the replayed records.
  • Does every job need a completeness claim?
    No. A single pass over a finished bounded input has a real ending, so completeness is observed rather than asserted. A stateless per-record transform takes no time-bounded decision at all. The claim is needed only where a running job must decide, before its input ends, that some period is finished.

saying these in an interview costs you the question

  • Says the claim proves every older record has already arrived
  • Treats the claim as the thing that publishes results
  • Confuses it with the marker that lines up a consistent recovery picture
  • Calls it the bookmark value a scheduled batch extract stores between runs
  • Assumes the claim is read off each record rather than computed by the job
open as a page

A job groups records by the moment stamped in the payload — where does one that arrives after its group closed go, and who would notice?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Nowhere visible: a record whose group has already closed is usually discarded without an error, so the totals are simply short. An explicit count of records dropped for lateness is what makes the loss visible at all.

open as a page

A continuous job keeps receiving records stamped earlier than ones it already handled — why is that normal, and what causes it?

level: juniorimportance: must knowfreq 62%

basics

~20 s

Out-of-order arrival is a property of the transport, not a fault: buffered devices, retries and parallel inputs merged together all deliver an earlier-stamped record after a later one. It costs nothing by itself until a time-based grouping has to be treated as finished.

open as a page

A daily order count runs over a continuous stream. Which three moments could each order be counted under, and which should the count use?

level: juniorimportance: must knowfreq 76%

basics

~20 s

Three moments compete: when the order was placed, stamped in the payload; when the receiving system wrote the record down; and the wall clock of the worker that handled it. A daily count belongs to the payload moment.

open as a page

A job carries a timestamp asserting nothing older than it is still expected. Why is that a guess rather than an observed fact?

level: middleimportance: must knowfreq 62%

basics

~20 s

Nothing in an endless input reports what has not yet been sent. The job infers the claim from the moments it has seen plus a standing assumption about how far behind a record may run, so it can be too early or too late.

open as a page

A record turns up after its group's answer was already published — what three things can a pipeline do with it?

level: middleimportance: must knowfreq 66%

basics

~20 s

Drop it, divert it to a separate repair channel, or reopen the group and publish a corrected answer. These are three different promises to everyone downstream about whether a published number can still change, not three settings.

open as a page

Before choosing how long a job waits for out-of-order records, what do you measure on the live source, and in what shape?

level: middleimportance: must knowfreq 54%

basics

~20 s

Measure, per record, how far behind the newest stamped moment already seen on that input it arrived. Report the result as a distribution with percentiles and a maximum, broken down per input and per hour, rather than as a single number.

open as a page

A continuous job buckets orders by hour, but no payload field was ever nominated as the record's moment. What is it actually grouping by?

level: middleimportance: must knowfreq 58%

basics

~20 s

Something, and probably arrival. Which moment a record is grouped by is assigned by the pipeline from a payload field or from source metadata; nominate nothing and the job either refuses to run or falls back to a clock nobody chose.

open as a page

A job's claim that nothing older than some timestamp is still coming: at what points can it advance, and what does that decide?

level: middleimportance: should knowfreq 48%

basics

~20 s

Where it can advance depends on the runtime: between any two records, usually damped by a short emission interval, where records flow one at a time; only at each job boundary where continuous work runs as repeated small finite jobs. That sets the finest step of any time-based decision.

open as a page

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%

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.

open as a page

A team halves how far behind the newest moment seen their pipeline holds its completeness claim - its assertion that nothing older is still coming - to freshen a dashboard. What did they trade?

level: seniorimportance: should knowfreq 55%

basics

~20 s

Latency for correctness, directly. Periods become eligible to close sooner, and more records now arrive after their period was already declared finished. The gain is paid on every result; the loss falls only on the tail of the source's behaviour.

open as a page

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%

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.

open as a page

Your team wants records arriving after a published daily total to correct it rather than be dropped — what must every downstream destination be able to do first?

level: seniorimportance: should knowfreq 54%

basics

~20 s

Each destination must be able to replace a value it already published for that group, keyed by the group's identity. An append-only destination turns a correction into a second value beside the old one, which reads as a double count.

open as a page

A team wants to raise a streaming job's out-of-order wait from thirty seconds to five minutes to catch more records — what does that cost and buy?

level: seniorimportance: should knowfreq 47%

basics

~20 s

It adds four and a half minutes of latency to every result the job publishes, not just to affected groups, and holds each group's retained state that much longer. It buys only the share of records the measured distribution places between the two durations.

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

Yesterday's records are pushed through the same job again this morning, and a whole day of history lands in one hour's bucket. Which moment was it grouping by, and which would have been immune?

level: seniorimportance: should knowfreq 54%

basics

~20 s

It was grouping by the wall clock of the worker that handled each record, which is read fresh on every pass, so the second pass re-dates everything to now. The moment stamped in the payload is immune because it travels inside the record.

open as a page

Phones stamp each order with their own clock and some records arrive dated a week in the future. What does grouping by the receiving system's stamp buy, and what does it cost?

level: seniorimportance: should knowfreq 46%

basics

~20 s

It buys a moment set on your side of the boundary that no producer can invent, and it costs the offline gap: an order placed in a tunnel and uploaded four hours later is now counted four hours late, and the delay you were measuring disappears.

open as a page

Three teams build reports on your hourly figures — what must your lateness policy state so that they know when a number is final?

level: principalimportance: should knowfreq 40%

basics

~20 s

Which of drop, divert or restate applies; the point at which a figure stops being provisional; whether it can change afterwards and how a change is signalled; and that the long tail beyond the policy is reconciled separately rather than waited for.

open as a page

Measured disorder shows 99.9% of a source's records within three hours and a tail reaching nine days — how long should the job wait, and what covers the rest?

level: principalimportance: should knowfreq 34%

basics

~20 s

Choose the wait at the shoulder of the bulk population and treat the nine-day tail as a separate problem: buying it with waiting would make every result nine days old. The tail belongs to a later reconciliation pass, not to the streaming hold.

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