Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Out-of-order events

Batch & Streaming · Streaming Design

Out-of-order events

Mediumbatch-streaming-43
out-of-orderevent-timewatermarksordering

Question

Why do events arrive out of order, and how do systems cope?

Solution

Events arrive out of order due to mobile devices operating offline, automatic network retries, varying transport latencies, and distributed routing across multiple messaging partitions. Message systems like Apache Kafka guarantee ordering only within an individual partition, meaning parallel processing can scramble message arrival sequences. Streaming engines cope with out-of-order data by using event-time windows bounded by watermarks, tracking entity sequence numbers, and buffering records for in-window sorting before computing outputs.

Root causes of event disorder

Real-world distributed systems frequently disrupt the chronological order of data:

  • Mobile and edge disconnects: Mobile applications cache interaction telemetry locally while devices travel through tunnels or airplane mode. When connectivity returns, hundreds of cached events upload simultaneously, arriving hours after their real creation time.
  • Network retries: A temporary gateway timeout may cause a publisher to retry sending packet A after packet B has already been acknowledged, inverting their sequence.
  • Partitioning across multiple channels: In Kafka, messages with different keys land in different partitions. Consumers read partitions at varying speeds, causing events generated at the same time to be processed out of order.

Event-time windows and watermarks

Processing engines avoid relying on arrival time (processing time) by embedding event timestamps directly into record schemas.

Engines group data using event-time windows paired with watermarks:

Events arriving: [t=10:05] [t=10:02] [t=10:08] [t=10:01]
Watermark advances to t=10:05
Window [10:00 - 10:05] collects [10:01, 10:02, 10:05] -> Sort & Process

The watermark establishes a bounded delay, allowing straggler events to arrive and collect in operator state before the window triggers.

Entity sequence numbers and in-window sorting

When business logic requires strict monotonic state transitions, relying on timestamps alone can fail due to clock skew across client devices.

In these pipelines, source systems append monotonically increasing sequence numbers or version counters per entity. When records enter a window, the streaming operator buffers the records in an internal priority queue, sorts them by sequence number, and executes transformations in exact order. If an event arrives with an older sequence number than the entity currently committed state, the operator rejects it as an obsolete revision.

PreviousNext