## TL;DR
Design for lateness from the start: process on event time with a watermark, allow a grace period for late events, and run a periodic reconciliation job that recomputes recent windows. Late data is normal (networks lag, devices batch uploads), so the pipeline should correct aggregates when it arrives rather than silently dropping it or corrupting results.

```text
how to handle late-arriving data in pipelines
```

## Use this when
- Events arrive after their time window already closed
- Aggregates need to be corrected when late data lands
- You are designing a streaming or windowed batch pipeline

## Not for this skill when
- Data never arrives at all (missing data, different problem)
- The schema of late events differs (schema drift)
- You need exactly-once delivery guarantees (related but separate)

## Steps

1. Write down the lateness contract: how late is acceptable, and what happens after:

```text
Events up to 24h late are incorporated automatically.
Events later than 24h go to a quarantine table for manual backfill.
```
Expected output: a documented policy. Without it, every late-data incident becomes an ad-hoc debate.

2. Implement watermarks or allowed lateness in the processing engine:

```python
# spark structured streaming example
(df.withWatermark("event_time", "24 hours")
   .groupBy(window("event_time", "1 hour"))
   .agg(sum("amount").alias("total")))
```
Expected output: events within the watermark update their windows; the engine tracks how far behind the event-time frontier it will still accept data.

3. Add a reconciliation job that recomputes recent windows on a schedule:

```sql
-- daily, recompute the last 3 days of aggregates from raw events
DELETE FROM hourly_totals WHERE window_start >= CURRENT_DATE - 3;
INSERT INTO hourly_totals SELECT ... FROM raw_events
WHERE event_time >= CURRENT_DATE - 3;
```
Expected output: corrected aggregates that include late arrivals, converging to the right numbers even if the streaming layer dropped something.

4. Make the output sink idempotent so reprocessing never double-counts:

```sql
-- upsert keyed on the window, not append
INSERT INTO hourly_totals (...) VALUES (...)
ON CONFLICT (window_start) DO UPDATE SET total = EXCLUDED.total;
```
Expected output: rerunning the reconciliation is safe. Idempotent sinks are what make "just recompute" a viable strategy.

5. Monitor the late-event rate and alert on spikes:

```sql
SELECT DATE(event_time), COUNT(*) FILTER (WHERE arrival_time - event_time > INTERVAL '1 hour') AS late_events
FROM raw_events GROUP BY 1 ORDER BY 1 DESC LIMIT 14;
```
Expected output: a daily late-event count. A spike means an upstream device, region, or job started lagging, and you want to know before the aggregates drift.

## Variant phrasings

### handling late data in streaming pipelines
Watermarks plus allowed lateness plus a reconciliation job is the standard trio. The watermark bounds state, the lateness window catches the stragglers, and reconciliation fixes whatever slipped through.

### event time vs processing time
Event time is when the thing happened; processing time is when your pipeline saw it. Late data is the gap between them. Always aggregate on event time, and treat processing time as the thing you monitor.

### late arriving events in kafka consumers
Same principles at the consumer level: keyed upserts into the sink, a reprocessing path for corrections, and monitoring on the event-time to processing-time skew per partition.

## Why it happens
Event time and processing time are different clocks. Devices buffer and retry, networks partition, batch uploads run nightly, and timezones confuse everyone. Any pipeline that assumes "data for hour H arrives during hour H" will eventually produce wrong aggregates. Lateness isnt an edge case; it is the normal operating condition of distributed data collection.

## Edge cases
- Very late data (days past the grace period) usually needs a manual backfill path; decide this in the contract, not during the incident.
- Watermarks that are too aggressive drop legitimate data; too lenient and state grows unboundedly. Tune from measured lateness, not guesses.
- Exactly-once semantics interact with reprocessing: make sure the reconciliation job and the streaming job cant both write the same window concurrently.
- Daylight saving transitions create duplicate or missing hours in local-time windows; prefer UTC event times everywhere.
- Quarantined late events still need retention and access policy; "manual backfill" without tooling becomes "never backfilled."

## Provenance

Resolved from the public thread: https://vectle.com/posts/pst_GC8Pp1OLZMpN8kawPBbYzw
