04 / 05

How would you choose a windowing strategy for joins with late events?

Difficulty: 8/10
Stateful stream processing, State stores, Windowing

Choosing a Windowing Strategy for Joins with Late Events

The core problem with late events is that stream time and event time diverge. A record's event time is when the event actually occurred; stream time is when it arrives at the processor. If a record arrives after its event-time window has closed, it is late. In a windowed join, you need to decide how long to keep a window open to wait for the matching record on the other side, and what to do with records that arrive after that. Kafka Streams uses the concept of a grace period: the window remains open for a configurable amount of time after the window end, and records arriving within the grace period are still processed. Records arriving after the grace period are dropped. The choice of window type and grace period is a trade-off between completeness and latency.

Kafka Streams offers several window types. Tumbling windows are fixed-size, non-overlapping windows (e.g., every 5 minutes). Hopping windows are fixed-size but overlapping (e.g., 5-minute window every 1 minute). Sliding windows are aligned to the record's timestamp and are useful for joins where you want to match records that are close in time regardless of fixed boundaries. Session windows are dynamic and based on activity gaps, useful for user sessions. For joins with late events, the key decisions are: which window type matches the business semantics, how large the window is, and how long the grace period is. A longer grace period handles more late events but increases state size and delays results. A shorter grace period reduces state and latency but drops more late events.

A common mistake is to set the grace period to zero, which means any record that arrives after the window end is dropped. This is almost never what you want in a real pipeline with network delays, producer retries, or out-of-order events. Another mistake is to make the window huge to avoid dropping late events; this increases state size and can cause restoration to take a long time. The trade-off is between completeness, latency, and state size. A good approach is to measure the distribution of event-time lateness in your data and set the grace period to cover, say, the 99th percentile. For joins, you also need to consider whether both sides are co-partitioned and whether the join is a stream-stream, stream-table, or table-table join. Stream-stream joins require windows; stream-table joins do not, because the table side is treated as current state. Version note: Kafka Streams introduced grace period in 2.1; before that, late events were handled differently. The suppression operator, also introduced around that time, allows you to emit results only after the window closes, which is useful for finalizing aggregates.

javascript
  1. 1

    Event time vs stream time: late events arrive after their window has closed.

  2. 2

    Grace period keeps the window open for late events; records after it are dropped.

  3. 3

    Tumbling, hopping, sliding, and session windows have different semantics for joins.

  4. 4

    Stream-stream joins require windows; stream-table joins do not.

  5. 5

    Longer grace period = more completeness but more state and latency.

  6. 6

    Measure lateness distribution and set grace to cover the 99th percentile.

  7. 7

    Grace period added in Kafka Streams 2.1; suppress operator emits after window closes.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.