Explain event-time versus processing-time and the role of watermarks in stateful windowing. How do Kafka Streams, Flink, and Spark differ in handling late and out-of-order data?
answer
- event-time = happened; processing-time = observed
- watermark = 'seen everything up to t'
- Flink: side-output + allowed lateness
- Streams: stream time + grace period, drop after
- watermark/grace also bounds state size
basics
~20 sEvent-time is when an event actually happened; processing-time is when the engine sees it. Watermarks estimate how far event-time has progressed so windows know when to close. Flink uses explicit watermarks; Spark uses withWatermark; Kafka Streams uses a grace period instead of true watermarks.
solid answer
~50 sEvent-time is the timestamp embedded in the record (when it happened); processing-time is the wall clock when the engine handles it. Because records arrive out of order and late, an engine needs a notion of progress to decide when a time window is complete. Flink and Spark use watermarks: a watermark W(t) asserts no events with timestamp <= t should still arrive, so windows up to t can fire; events past the watermark are 'late'. Flink lets you choose drop, side-output, or allowed-lateness handling; Spark drops events older than the watermark and bounds state with withWatermark. Kafka Streams does not use watermarks — it tracks 'stream time' (the max observed event timestamp per partition) and closes windows after a configurable grace period (ofSizeAndGrace / window grace); records arriving after grace are dropped. The practical difference is how explicitly you control lateness and how state is bounded.
go deeper
Define event-time vs processing-time and that windows need a way to know they are complete.
Explain watermarks as 'progress markers' and that late data past them gets dropped; know withWatermark and grace period.
Contrast Flink's flexible lateness handling (side outputs, allowed lateness) with Streams' stream-time + grace and Spark's withWatermark, including state bounding.
Reason about idle-source stalls, watermark-vs-state-size trade-offs, and choosing engines by required lateness control for SLAs.
## Three notions of time - **Event-time**: the timestamp of when the event occurred at the source (e.g., the sensor reading time). In Kafka this is usually the record's timestamp field (set to CreateTime). - **Ingestion-time**: when the record entered the system (broker append time, LogAppendTime). - **Processing-time**: the wall-clock time when the operator processes the record. Real streams are **out of order** (a phone offline for an hour sends old events later) and have **skew** across partitions. To produce correct windowed aggregates (e.g., 'count per 1-minute window'), the engine must know event-time, and must decide *when a window is done* despite stragglers. ## Watermarks (Flink and Spark) A **watermark** is a moving marker in event-time. A watermark of value `t` is an assertion: *'I believe I have seen all events with timestamp <= t; later ones are late.'* It is usually computed as `max observed event-time - allowed out-of-orderness`. When the watermark passes a window's end, that window **fires** (emits its result) and its state can be cleaned up. - **Flink**: you attach a `WatermarkStrategy` (e.g., `forBoundedOutOfOrderness(Duration.ofSeconds(10))`). Late events past the watermark can be **dropped**, sent to a **side output** (`OutputTag`) for separate handling, or tolerated via **allowed lateness** that keeps the window state around and re-fires. Watermarks also drive **timers** and event-time processing functions. This is the most explicit, flexible model. - **Spark Structured Streaming**: you call `.withWatermark("eventTime", "10 minutes")`. This bounds how long window/aggregation state is retained and **drops** events older than `(max event-time - watermark threshold)`. Stream-stream joins require watermarks on both sides to bound state. Spark's model is simpler but less configurable (no per-event side-output of lateness out of the box). ## Kafka Streams — stream time + grace, not watermarks Kafka Streams deliberately avoids a global watermark. Each task tracks **stream time** = the maximum event timestamp observed so far on that partition (it never goes backward). Windowed operations use `TimeWindows.ofSizeAndGrace(size, gracePeriod)` (or `ofSizeWithNoGrace`). A window's result can keep updating until stream time exceeds `window end + grace`; after that the window is **closed** and later records for it are **dropped** (counted in the `late-record-drop` metric). Because stream time advances only when records arrive, a quiet partition can stall window closing — there is no idle-source watermark advancement by default the way Flink offers `withIdleness`. Suppression (`Suppressed.untilWindowCloses`) lets you emit only the final result per window instead of a stream of intermediate updates. ## How they differ in practice | Concern | Kafka Streams | Flink | Spark Structured Streaming | |---|---|---|---| | Progress mechanism | stream time (per-partition max ts) | watermarks (configurable strategy) | watermarks (withWatermark) | | Late data handling | dropped after grace | drop / side-output / allowed-lateness | dropped past threshold | | Idle source handling | window may stall | withIdleness | handled via watermark/processing | | Control granularity | coarse (grace + suppress) | fine | medium | ## Edge cases - **Stalled windows on idle partitions**: if one partition stops producing, stream time (Streams) or the min-across-partitions watermark (Flink/Spark) can freeze, delaying output. Flink/Spark mitigate with idleness detection. - **State growth**: without a watermark/grace bound, windowed state would grow forever. The grace/watermark is also a **state-retention** control, not just a correctness control. - **Choosing the bound**: too small drops legitimate late data; too large inflates state and latency. It is a tunable trade-off, not a fixed value.
- Why does a Kafka Streams windowed aggregation stop emitting when one input partition goes idle?Stream time is the max observed event timestamp; if a partition stops sending records its contribution to stream time freezes, so windows whose closing depends on advancing stream time never reach window-end + grace. You mitigate with heartbeat records or by ensuring traffic; Flink/Spark offer explicit idleness handling.
- In Flink, what are the options for an event that arrives after the watermark has passed its window?Drop it (default), route it to a side output via OutputTag for separate processing, or use allowed-lateness to keep the window state alive and re-fire the updated result for a bounded extra period.
saying these in an interview costs you the question
- Saying Kafka Streams uses watermarks like Flink — it uses stream time plus a grace period, not watermarks.
- Claiming late data is always recovered — past the watermark/grace, it is typically dropped unless you configured side-output/allowed-lateness.
- Confusing event-time with processing-time, or assuming records arrive in order.
- Forgetting that the watermark/grace also bounds state retention, not just correctness.