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.