Command Palette

Search for a command to run...

Hectal
PHASE 11Advanced ~8 min· topic 4 of 5

Topic 11.4

Event Time, Late Events, and Grace Periods

In one line

Stream processors should usually group by event time (when it happened), not processing time (when it arrived). Events arrive late and out of order; a grace period keeps windows open for late arrivals, and events later than the grace are dropped (and counted) or handled by a correction path.

0/5 · 0%

Think of it like this

Counting votes by the hour they were cast. Postal votes cast at 10:01 might arrive at 10:10. If you count by arrival time, they land in the wrong hour; if you wait a while before declaring each hour's total, you catch most of them, but at some point you must publish.

Key ideas

  1. 01

    Times: event time (the record timestamp or a payload field, extracted with a TimestampExtractor), ingestion/log-append time (when the broker stored it), processing time (when your app handles it). Kafka Streams uses record timestamps by default.

  2. 02

    Stream time: Kafka Streams tracks, per task, the maximum event timestamp seen so far. Windows close when stream time passes window end + grace. This is similar to watermarks in Flink or Beam, but simpler: it only moves forward with data.

  3. 03

    Grace period: ofSizeAndGrace(5m, 2m) accepts records up to 2 minutes after the window end (in stream time). Later records are dropped and counted in the dropped-records metric.

  4. 04

    Trade-off: longer grace catches more late data but delays final results and holds more state. Choose from the observed lateness distribution (for example the 99.9th percentile of arrival delay).

  5. 05

    Handling very late data: route it to a side topic for correction jobs, or recompute affected aggregates in a batch job; don't silently lose it if it matters (billing).

Code & diagrams

PayloadTimestampExtractor.javajava
public class PayloadTimestampExtractor implements TimestampExtractor {
    @Override
    public long extract(ConsumerRecord<Object, Object> record, long partitionTime) {
        if (record.value() instanceof Purchase p && p.occurredAtMs() > 0) {
            return p.occurredAtMs();          // business event time
        }
        return partitionTime;                 // fall back to the partition's stream time
    }
}
// props.put(StreamsConfig.DEFAULT_TIMESTAMP_EXTRACTOR_CLASS_CONFIG, PayloadTimestampExtractor.class);
late-event.mermaiddiagram
Rendering diagram…

Interview problem

The problem

An event from 10:01 arrives at 10:10

An event occurred at 10:01 but arrived at 10:10. Your 10:00–10:05 window has already closed. What should happen?

Explain it without notes

01

What is stream time in Kafka Streams and how does it close windows?

Practice

01

In TopologyTestDriver, send an on-time event, advance stream time past the grace, then send a late one; check the dropped-records metric.

Trade-offs

  • ↔

    Longer grace means more complete, later, more state-heavy results.

Done when you can

  • I use event time, choose grace periods from data, and handle very late events deliberately.