skip to content

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%

answer

  1. valid record, no error raised
  2. short totals, green run
  3. the claim already passed it
  4. count it or never know
  5. compare moment to claim yourself

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.

solid answer

~50 s

A continuous job carries a **completeness claim** alongside the records — a *watermark*: a timestamp asserting that nothing older than it is still expected to arrive. Once that claim has passed a group's end, the group is eligible to close. A record stamped earlier that turns up afterwards is *late*, and on most runtimes the default is that it is thrown away with no error, no failed run and no rejected row — the output looks fine, it is just short. The thing that makes it visible is a **late-record counter**: an explicit count of records discarded for arriving after their group. How much you get for free varies: some runtimes publish a count per grouping, some only per-run progress figures, some nothing. Where nothing is published, compare each record's moment against the claim in a step before the grouping and count it yourself.

go deeper

for a junior

Know that a record which misses its group is usually discarded quietly rather than rejected loudly, and that a successful run says nothing about whether records were dropped.

for a middle

Explain where the drop happens — the completeness claim had already passed the record's moment — and name the explicit count that turns an invisible loss into a number with a trend.

for a senior

Show how you would make the number operational: alert on the rate rather than the total, and know that a zero can mean a clean source, an unwired counter, or a grouping that absorbs the record instead.

for a principal

Treat silent loss as a reporting question. If no published figure carries a dropped-record count beside it, the organisation cannot tell a genuine decline from data the pipeline threw away.

## The three things an arriving record can be Every record in a continuous job has a moment attached to it: the moment the thing being recorded actually happened, written into the payload by whatever produced it — often called **event time**, and here the **occurrence time**. A job that groups by that moment has to decide, *while it runs*, when a span of moments is finished. It does that with a **completeness claim** carried alongside the records — a *watermark*: a timestamp travelling with the data that asserts nothing older than it is still expected to arrive. The claim is a running guess the job makes, not a fact it reads off the data. Measured against that claim, an arriving record is one of three things: - **On time** — its moment is newer than the claim, so it lands in its group and there is nothing to discuss. - **Out of order** — it arrived after a record carrying a later moment, but the claim has not yet passed its own moment. It still lands in its own group. This is normal and costs nothing. - **Late** — the claim had already passed its moment, so the group it belongs to was already eligible to close before it showed up. Every late record arrived out of order; most out-of-order records are never late. Only the third case is a loss. ## Why the loss is silent A late record is not an invalid record. It parses, it satisfies the schema, its fields are sane — the one thing wrong with it is that it is *behind the job's own guess about time*. So none of the machinery that shouts at you fires: | Signal you already have | What it actually tells you | What it hides | |---|---|---| | Run status is green | No worker died and no step threw | Discarded records are not errors | | Output row count | How many groups were published | Nothing about records that missed one | | Input record count | How many records the job read | That some were read and then dropped | | Destination rejections | Which rows the destination refused | Records that never reached it | | Late-record counter | How many valid records missed a group | Which records, and what they were worth | The symptom, if anyone ever sees one, arrives weeks later from outside the pipeline: a finance or operations figure computed some other way is a fraction of a percent higher than yours, every day, and nobody can say since when. ## The counter, and getting one when you are not given one The **late-record counter** is an explicit count of records discarded for arriving after their group. Where the runtime publishes one, wire it to a dashboard and an alert on the *rate*, not the cumulative total — a total that only ever rises tells you nothing about today. Where it does not, you can produce the same number yourself: 1. Assign the record's moment explicitly, so you know which field the grouping is actually using rather than inheriting a default. 2. In a step *before* the grouping, compare that moment with the claim the job is currently carrying, where the runtime lets a per-record step read it; where it does not, compare against the newest moment seen so far minus the wait the pipeline has agreed to hold back. 3. Count every record that fails the comparison, and emit the count next to the group totals so a reader of the number can see both. The counter answers *how many*. It does not answer *which*, or *what they were worth* — a hundred dropped sensor readings and a hundred dropped payments are the same integer. ## What varies between engines This is one of the places where a sentence true of one engine is false of the next, so state the mode you mean: - **A grouping that republishes an updated running result** for as long as it retains the group will absorb the record rather than drop it. The counter then legitimately stays at zero while a previously published answer quietly changes — a different problem, not an absence of one. - **A grouping that emits once and releases the group** has nothing left to update, so the record has nowhere to go but away. - **A single pass over a finished bounded input** has no running claim at all, so during the run nothing is late by this definition; the same record simply turns up in the input some later run reads. - **Where the claim is readable from inside a per-record step** you can count lateness yourself; where it is not, your own comparison is an approximation of the runtime's. ## What the number is for A counter nobody watches is the same as no counter. The point of having it is that it converts an invisible, unbounded, unfalsifiable worry — *are our numbers short?* — into a quantity with a trend, which is the precondition for deciding what the pipeline should do about late records at all.

  • How would you count late records when the runtime publishes no such number?
    Assign the record's moment explicitly, then in a step before the grouping compare it with the completeness claim the job carries; count every record the claim has already passed. Where a per-record step cannot read the claim, compare against the newest moment seen minus the wait the pipeline holds back. It approximates the runtime's own arithmetic, which is enough to see a trend.
  • Does a record that arrives out of order show up on the late-record counter?
    Usually not. Out of order means it arrived after a record carrying a later moment but before the claim passed its own — it still lands in its group and costs nothing. Only a record the claim has already passed is late. A source can be wildly out of order and produce no lateness at all.
  • The counter has read zero for a month. What are the possible explanations?
    Three, and they look identical from the dashboard: the source genuinely produces nothing the claim has passed; the counter was never wired to anything; or the grouping keeps republishing an updated result rather than dropping, so nothing is ever counted while published answers change underneath. Confirm which before treating zero as good news.

saying these in an interview costs you the question

  • Says a late record fails the job or raises an error somewhere visible.
  • Assumes late records are automatically retried into their group later.
  • Treats a green run and a plausible row count as proof no data was lost.
  • Calls any out-of-order record late, erasing the distinction that matters.
  • Expects the destination's rejected-row count to reveal records dropped for lateness.