A streaming job advances a timestamp asserting nothing older than it is still expected. What does that assertion license?
answer
- an endless input supplies no ending
- somebody has to declare the hour finished
- a timestamp travelling with the records
- compared against a group's end
- eligible to close, not emitted
basics
~20 sA 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 sAn 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
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.
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.
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.
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