Lesson 3 of 4 · 25 min

Handle late data with explicit window rules

Compute event-time results with a stated lateness policy.

Mechanism and reasoning

Event time says when the event happened. Processing time says when the pipeline observed it. The difference matters when devices go offline, producers retry or a source delivers a delayed batch. A processing-time window is easy to operate, but it can assign business activity to the wrong period.
A watermark estimates progress in event time. It is not a proof that no earlier event can ever arrive. The pipeline needs a policy for data that arrives after the watermark and for data beyond allowed lateness. It might update prior results, send corrections, or place records in a review stream. The business should know whether a published value is provisional or final under that policy.
Define window boundaries precisely. A half-open interval includes its start and excludes its end. This avoids counting a boundary event in two adjacent windows. State the timezone and handling of ambiguous local times. UTC storage helps, but business reporting may still require a local calendar.
A late event can change more than one number. A session window may merge previously separate sessions. A corrected dimension can reclassify historical facts. Retractions or upserts must use stable output keys so downstream consumers do not append a second result as if it were new activity.
In an interview, work through a small stream with event and arrival times. State when an early result appears, when it updates and when the system stops accepting corrections automatically. Do not promise both immediate finality and complete inclusion of arbitrarily late data. The policy has a cost in retained state, processing and user expectations.

Follow the same window through multiple publications

Separate the logical window from each publication of its result. The logical key can be metric, tenant and window start. A publication adds a revision or update time. Downstream readers must know whether a new row replaces the previous aggregate or adds a delta. If the producer publishes complete sums and the consumer adds them, every late update inflates the result.
code
1Illustrative complete-value update stream2key=(tenant7, revenue, 10:00)3revision1: sum=10, status=provisional4revision2: sum=13, status=provisional5revision3: sum=13, status=closed-under-policy6Consumer rule: retain highest revision per logical window key7Incorrect consumer rule: add10 +13 +13
The correct current value is thirteen, not thirty-six. If the system instead publishes deltas, the messages and deduplication contract must say so. A delta of three can update ten to thirteen, but replaying that delta twice is unsafe without stable update identity. Complete values and deltas are both useful; mixing their semantics is the error.
A watermark is normally driven by source and runner behavior, not simply by the machine clock reaching the window end. With multiple inputs, one slow or idle partition can delay progress under the framework's combination rules. Some systems support idle-input handling, but the exact behavior is implementation-specific. Do not invent a universal watermark formula from the wall-clock delay observed in a small test.
RecordEvent timeArrival conditionWindow action
A: value 410:01Before first publicationInclude
B: value 610:04Before first publicationInclude
C: value 310:03Window still accepts late updatesRevise to 13
D: value 810:05Any ordinary arrivalNext window
E: correction−210:02Beyond automatic horizonExplicit correction path
The final record makes the business decision visible. Dropping E silently produces a number that omits known information. Applying it silently changes a supposedly closed report. A correction ledger can preserve both the original published value and the later adjustment with a reason. Some consumers may restate history; others may book an adjustment in a later period. The data pipeline must follow the agreed reporting policy rather than choose whichever is easiest to code.
Allowed lateness has a state cost. Retaining more windows permits more automatic revisions but keeps aggregation state longer and can increase downstream update traffic. A short horizon reduces state but moves more records to manual or batch reconciliation. The right horizon depends on measured lateness and business consequences, not a universal setting. Plot the observed lateness distribution while keeping exceptional outages visible; a median delay says little about the rare records that alter closed results.
Calendar boundaries need explicit timezone rules. A UTC timestamp is an unambiguous instant, but a business day in a local timezone is a calendar interval. Daylight-saving transitions can make a local day shorter or longer than twenty-four hours. Avoid deriving local days by dividing epoch seconds by a constant. Use a specified timezone conversion and test transition dates if that reporting calendar is in scope.
Late session events are more complex than fixed-window sums. An event between two existing sessions can connect them under a gap rule, requiring the system to retract or replace earlier session outputs. A stable session output strategy must account for merged identity and downstream corrections. Merely appending another session row can double-count activity. If the interview only requires fixed windows, state that this additional behavior is outside the simplified example.
Test the pipeline with on-time, boundary, accepted-late and beyond-horizon records. Then replay the same updates out of order at the downstream sink and verify that revision handling preserves the intended current value. The source window policy and sink update policy are two halves of the result. Correct window assignment alone does not guarantee that the final dashboard is correct.

Worked example

Teaching window is [10:00, 10:05) UTC with two minutes of allowed lateness under the chosen runner policy. Events worth 4 and 6 occur at 10:01 and 10:04 and arrive on time, giving ten. A value of 3 with event time 10:03 arrives while the window remains eligible for late updates; the result becomes thirteen under the same window key. An event exactly at 10:05 belongs to the next window. A later arrival beyond the retention policy goes to a correction path rather than silently changing a finalized number.

Exercise

For [12:00, 12:10), values are 2 at 12:00, 5 at 12:09:59 and 7 at 12:10. A late value of 4 at 12:03 is accepted. Calculate the updated sum and identify the next-window value.

Model solution and rubric

The first window contains 2 + 5 + 4 = 11. The value seven belongs to [12:10, 12:20). Write an upsert for the first window's stable key rather than append eleven next to the earlier seven. Verify downstream readers interpret updates correctly. The exercise does not specify a particular runner's trigger API, so the answer must keep policy separate from implementation syntax.
Score out of four: one point for the correct result, one for showing the intermediate reasoning, one for identifying the stated failure case, and one for a verification that could disprove the answer. Do not award the reasoning point for a tool name alone.

Failure modes and misconceptions

“Closed under policy means no earlier event can exist.” It means automatic processing follows a stated horizon. Later evidence can still require correction.
“A revised complete sum is another contribution.” It replaces the earlier value for the window key. Adding complete revisions double-counts the same activity.

Interview probe

Evidence class: recommended. Original practice.
Can a watermark guarantee that a day's revenue is final?
Strong answer: It represents event-time progress under a source and runner policy. Finality requires an explicit lateness and correction contract. I would expose provisional status and route arrivals beyond the automatic correction horizon for reconciliation.
Follow-up: What extra state is needed if late events can merge sessions?
Weak answer indicators: Treating watermark as wall-clock time; double-counting the boundary; appending corrected aggregates without stable keys.

Sources

Technical references: Apache Beam programming guide; Kafka 4.1 design. Sources support the documented mechanisms. The numbers, decisions, rubrics and interview prompts in this lesson are original teaching examples, not measurements or employer question claims.
docsApache Beam programming guidebeam.apache.orgdocsKafka 4.1 designkafka.apache.org

Checkpoint

A complete window result changes from 10 to 13. What should a latest-value sink do?

AAppend both as separate business totalsBAdd them to 23CKeep 10 because it was firstDReplace the same logical window key with the newer revision's value of 13
Sign up free to answer and see why

Checkpoint

An event at exactly 10:05 belongs to which half-open interval?

A[10:05,10:10)B[10:00,10:05)CBothDNeither
Sign up free to answer and see why

Checkpoint

What does a watermark establish?

ANo older event can ever arriveBProgress under the source/runner policyCThe exact current wall-clock timeDThe final accounting policy for every consumer
Sign up free to answer and see why

Checkpoint

A late event arrives beyond the automatic correction horizon. What should happen?

ASilently mutate a closed reportBSilently discard it in every business domainCApply the explicit correction or review policyDAssign its event time to arrival time without disclosure
Sign up free to answer and see why

Checkpoint

Automatic late correction grows from one hour to seven days. What operational change should be planned?

ARetain aggregate state longer and account for additional downstream revisionsBKeep the one-hour state TTL because watermark progress is unchangedCTreat all previously published values as permanently finalDRemove revision keys because older corrections now arrive on time
Sign up free to answer and see why

Explain how you would compute event-time results with a stated lateness policy without reading the solution. State one assumption that could change your answer, and one observation that would make you revise it.

Not yetGetting thereConfident

Wrap-up

  • Specify event time, boundaries, lateness and correction behavior together. A number is incomplete without its finality rule.

Sources

Free to read · better with Enzo

Learn it with Enzo

Save your progress, answer the checkpoints, and let Enzo quiz you on what you just read.