Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Choosing a watermark delay

Batch & Streaming · Streaming Design

Choosing a watermark delay

Mediumbatch-streaming-35
watermarksevent-timelagstreaming

Question

How do you pick the watermark delay for a streaming job?

Solution

You pick a watermark delay by measuring the statistical distribution of lag between event creation time and ingestion processing time, typically setting the boundary at the 99th percentile. A longer watermark delay produces more complete aggregations but increases output latency and expands state store memory requirements. Conversely, a shorter delay emits results quickly but forces the pipeline to drop or misplace late-arriving events. Watermark strategies must account for distinct per-source characteristics and should be revisited after major operational incidents.

Measuring skew distribution

Before selecting an arbitrary delay value, query raw historical data in your object lake or message bus to calculate the difference between ingestion time and event creation time.

Event generation (mobile app)  ---> Ingestion (Kafka broker)
t = 12:00:00                       t = 12:04:30
Lag = 4 minutes 30 seconds

Plotting this lag across percentiles reveals real client behavior:

  • Web browser traffic over stable broadband typically exhibits a tight p99 lag of two to five seconds.
  • Mobile client applications or IoT sensors traversing cellular networks frequently demonstrate long tail distributions, with p99 lag reaching ten to fifteen minutes due to intermittent signal drops and local client queuing.

The trade-off between latency and state size

Setting the watermark delay dictates how long window operators hold intermediate calculations inside state backends like RocksDB:

  • Long delays: Retaining records for thirty minutes guarantees that nearly all stragglers are captured. However, downstream consumers wait thirty minutes past the window close before receiving outputs, and the cluster must retain substantial memory state.
  • Short delays: Configuring a five-second delay produces snappy dashboards. However, any mobile device uploading data on a fifteen-second retry will have its events rejected as late data.

Source differentiation and operational adjustments

Different input streams require customized watermark policies rather than one global setting. When joining web telemetry with mobile telemetry, apply separate bounded out-of-orderness generators to each source before performing unified transformations.

Watermark configurations are not static. Revisit your delay settings after major platform incidents, such as upstream gateway outages or mobile release regressions, because unexpected publisher behavior can dramatically shift lag distributions.

PreviousNext