- Published on
Late-Arriving Records Without Rewinding the Whole Pipeline
- Authors

- Name
- Mehdi Akiki
Article · Derived state
A record is late when it arrives after the system has already made a decision about the time where that record belongs.
This is not only a streaming-framework problem. An API export can arrive tomorrow with an event from last week. A mobile client can reconnect after three days. A provider can finish a backfill after the daily report was sent.
I do not rewind all history for every late record. I find the smallest derived state the record can affect, recompute or compensate that state, and publish the correction with an explicit version.
Four times that should not be collapsed
For one business event I may record:
event time: 2026-09-24 23:58 when the purchase happened
source write time:2026-09-25 00:03 when the source stored it
ingestion time: 2026-09-26 09:12 when my pipeline received it
processing time: 2026-09-26 09:13 when this stage handled it
The purchase belongs to the September 24 revenue window even though the pipeline first sees it on September 26.
If I group by ingestion time because it is convenient, the pipeline is operationally consistent and semantically wrong.
A watermark is an estimate, not a fact from the future
A watermark represents progress in event time: the system estimates that it has seen the relevant events before some point. When the watermark passes a window's end, the system may emit a result. An older event that appears afterward is late.
Apache Beam's model documentation describes a watermark as an estimate of completeness and uses triggers to control early and late output. Google Cloud's streaming pipeline guide similarly distinguishes event time from processing time and explains that data can arrive out of order.
The estimate lets work progress without waiting forever. It also creates a product question: what happens when the estimate is wrong?
One late record, one affected window
Suppose daily aggregates currently contain:
2026-09-24 / tenant-7 / revenue = 1,200.00 / version 8
A verified purchase of 35.00 arrives two days late. Its identity is new and its event time is September 24.
I map it to the affected key:
window = day(2026-09-24)
partition = tenant-7
Then I recompute that partition-window from durable input facts or apply an idempotent delta:
old aggregate: 1,200.00 / version 8
late delta: 35.00 / event purchase-991
new aggregate: 1,235.00 / version 9
The remaining days and tenants do not move.
Choose recompute or delta from the operation
An incremental delta is safe for an operation with a clear inverse and deduplication identity. Sum of immutable purchases is a good example.
Full affected-window recomputation is safer when:
- the aggregation is non-associative or depends on order;
- an input can be edited or deleted;
- business rules changed;
- several late records may interact;
- previous intermediate state is not trustworthy;
- the result contains top-N, sessions or attribution logic.
“Recompute” still means one bounded partition and window, not the whole pipeline.
I retain raw or normalized facts long enough to rebuild every window inside the accepted correction period. If input evidence expires earlier than correction obligations, the recovery promise is false.
Published results need versions
Downstream consumers must distinguish an initial result from its correction:
{
"aggregateId": "tenant-7/revenue/2026-09-24",
"version": 9,
"supersedes": 8,
"value": "1235.00",
"status": "corrected",
"reason": "late_input",
"computedThrough": "2026-09-26T09:13:00Z"
}
A database projection can update version 8 to 9 conditionally. An event consumer can keep the greatest version per aggregate identity. A generated report may require a correction notice rather than silent replacement.
This is where technical and product semantics meet. Updating a dashboard is easy. Correcting an invoice, email or external filing may require a compensating business operation.
Finality is a product contract
I define states such as:
provisional → settled → corrected
Provisional results change freely within the common lateness period. Settled results crossed an operational threshold but may still receive a correction. Some regulated outputs become externally final and need a separate amendment workflow.
“Final” must not mean “our timer fired.” It describes what downstream users may now rely on and which correction process remains available.
Allowed lateness is a cost decision
Keeping window state forever is expensive. Dropping every late event is inaccurate. I choose an allowed-lateness period from observed data and business obligations:
0–2 hours: update active state immediately
2 hours–7 days: reopen affected window and publish a new version
7–90 days: enqueue reviewed correction/backfill
over 90 days: retain evidence and follow explicit business policy
These are example bands, not universal defaults. Payment, telemetry and collaborative editing products have very different acceptable delay and correction costs.
The policy includes what happens after the easy automatic window. “Dropped as too late” is a product outcome that needs a metric and owner.
Ingestion overlap and downstream correction are separate
An overlap-window reader tries to discover source changes that became visible late. Once discovered, this article's problem begins: derived state may already have been emitted.
ingestion repair: did I eventually observe the source fact?
derivation repair: did every affected output incorporate it?
Closing the first gap does not automatically repair cached aggregates, search indexes, notifications or exported reports.
Keep an impact index
For expensive pipelines, I store enough lineage to locate affected outputs:
input identity → event-time partition → derived window IDs → published version
This can be implicit when the mapping is deterministic, such as tenant plus calendar day. Complex joins may need explicit dependency records or a bounded query over the input keys.
The index does not need field-level lineage for every value. It needs to answer the recovery question: which durable outputs can this late or corrected fact change?
Test lateness as a normal path
My pipeline tests deliver the same logical input in different orders:
- all records on time;
- one record before and one after the watermark;
- the late record duplicated;
- a correction to an already-late record;
- a deletion after settlement;
- lateness just before and just after the policy boundary;
- process crash between new aggregate storage and correction publication;
- downstream receiving version 9 before version 8;
- rule version changing before recomputation.
For each order I assert the final durable facts and published aggregate version. When an external effect cannot be silently replaced, I assert the compensating workflow.
Measure the tail
I watch:
- lateness distribution, not only average delay;
- late records by provider, tenant and record type;
- windows reopened and time to correction;
- corrections rejected by policy;
- oldest unprocessed correction;
- output versions delivered out of order;
- downstream effects awaiting compensation;
- input evidence approaching retention expiry.
A rising late-data tail may reveal provider trouble, offline clients, clock errors or an ingestion backlog. One global percentage hides which repair promise is failing.
My practical rule
Late data is normal in distributed systems. The design decision is how much state and evidence I keep, when I emit a provisional answer, and how I correct it.
I preserve event time, map each input to the smallest affected output, recompute or apply an idempotent delta, version the result and define finality as a product contract. This repairs the truth that changed without treating the entire history as one transaction that must start again.