Flink Event-Time Watermarks
MSc Data Science · Big Data coursework (Task 4) · Coventry University
Events never arrive in the order they happened. Watermarks let Flink count them by when they happened, not when they showed up, and still produce results on time.
- Data Engineer
- 2025
- Stream processing
- 6 technologies
- Apache Flink 1.18
- Java 17
- Apache Kafka
- Kafka UI
- Maven
- Docker Compose
15 s
Tumbling event-time windows
5 s
Allowed out-of-orderness
4,652
Posts streamed (1,545 + 3,107)
2
Kafka partitions and parallelism
10 s
Idle-partition timeout
4
Flink jobs
The problem
A tweet sent at 10:00:14 can reach the system at 10:00:16 because of mobile signal, retries, buffering, or Kafka batching. If windows use the processing machine's clock, that tweet lands in the wrong 15-second window and the counts change on every run. The job needs correct, repeatable counts without waiting forever for stragglers.
Overview
The stack runs in Docker Compose: ZooKeeper, a Kafka broker with two listeners (kafka:9092 inside the Docker network, localhost:29092 from the host), a Flink 1.18 JobManager and TaskManager with two task slots, and Kafka UI for browsing topics and messages.
Two social-media CSV datasets are published line by line into the social_facebook (1,545 messages) and social_twitter (3,107 messages) topics with kafka-console-producer inside the broker container, which avoids Windows path issues. Two verifier jobs first confirm that Flink reads both topics.
The hashtag counters read each topic with a KafkaSource, keep lines that contain the chosen hashtag (case-insensitive, with or without #), and count them in 15-second tumbling event-time windows. The event time is parsed from each post's date_posted field, and a bounded out-of-orderness watermark of 5 seconds lets late posts still land in the right window. Each window emits a small JSON result (hashtag, window start and end, count) to an output Kafka topic through a KafkaSink.
How it works
Hashtag counting job
- 01
KafkaSource
Reads social_twitter or social_facebook as plain strings with its own consumer group.
- 02
Watermark strategy
forBoundedOutOfOrderness(5 s) with a timestamp assigner for date_posted, plus withIdleness(10 s). Applied on the source so each Kafka partition gets its own watermark.
- 03
Filter
Keep only posts containing the hashtag, matched case-insensitively and with or without the # sign, passed in as -Dhashtag=…
- 04
Window
TumblingEventTimeWindows groups posts by when they were written. A window closes only once the watermark passes its end.
- 05
Count
Counts the posts in each closed window and formats {hashtag, windowStart, windowEnd, count} as JSON.
- 06
KafkaSink
Writes the results to twitter-hashtag-counts or facebook-hashtag-counts, where they can be consumed or viewed in Kafka UI.
Implementation
Screenshots from the running Docker stack, Kafka UI, and the Flink dashboard. Tap any figure to zoom.
The stack is up
Kafka, ZooKeeper, JobManager, TaskManager, and Kafka UI all running. Before any job, docker compose ps, the service logs, the Flink UI on port 8081, and a Kafka produce/consume smoke test confirm the infrastructure is healthy.
Datasets in Kafka
Each CSV line became one message: 1,545 in social_facebook (709 KB) and 3,107 in social_twitter (986 KB), one partition each for a single-broker setup.
Flink reads both topics
Two verifier jobs consume each topic and print the lines to the TaskManager logs, proving the Kafka-to-Flink connection before adding windowing logic.
Hashtag counters running
TwitterHashtagCounter and FacebookHashtagCounter running with event-time windows and watermarks. The completed list keeps the earlier verifier and test runs.
Results back in Kafka
The window results land in twitter-hashtag-counts and facebook-hashtag-counts, ready for any downstream consumer.
Scaling out to two partitions
Both input topics altered to 2 partitions with kafka-topics --alter, so two Flink subtasks can each read one partition.
Partition-aware job graph
Source → filter → map runs at parallelism 2, one subtask per partition, each tracking its own watermark. The global windowAll operator runs at parallelism 1 and receives the minimum watermark across partitions.
Watermark types
Periodic watermarks
Emitted on a timer (e.g. every 200 ms) as the latest timestamp minus an allowed lateness. Easy to configure and predictable, but needs tuning and isn't flexible during traffic spikes.
Punctuated watermarks
Emitted only when a special marker event arrives (e.g. end_of_batch). Strict and accurate for structured logs, but social media has no such boundary events.
Event-time watermarks
Advance with the timestamps actually observed. They adapt to network delays, bursts, uneven arrival, and slow partitions, so fewer events are wrongly marked late.
Processing time vs event time
- Accuracy: processing-time windows use the Flink machine's clock, so a delayed tweet falls into the next window and counts change between runs. Event time with watermarks always puts it in the window it was written in, so results are correct and reproducible.
- Latency: processing-time windows close the moment the clock passes the boundary. Event-time windows deliberately wait for the watermark, adding up to the 5-second lateness bound.
- Throughput: processing time is simpler and handles slightly more events per second. Event time must read each timestamp, compare it with the watermark, and buffer out-of-order events.
- CPU and memory: event time holds late events and timers in state until each window is complete, so it uses more resources.
- Verdict: when accuracy matters (hashtag trends, billing, monitoring, anomaly detection) event time with watermarks is the industry standard and worth the extra cost.
Per-partition watermarks
- Before: one global watermark meant a slow or idle partition held back event time for the whole job, so data from faster partitions waited.
- After: the WatermarkStrategy is applied directly on the KafkaSource, which tracks a watermark per partition and forwards the minimum to downstream operators. withIdleness(10 s) stops an idle partition from blocking progress.
- Accuracy was unchanged between parallelism 1 and 2, because timestamps, not arrival order or task assignment, decide which window an event belongs to.
- Throughput rose and latency improved under uneven traffic: ingestion, timestamp extraction, filtering, and mapping now run in parallel.
- Limitation: windowAll is a single, non-parallel operator, so it becomes the bottleneck. Full scaling needs keyed windows, keyBy(hashtag), so each hashtag aggregates in parallel.
What I'd add next
- Switch from windowAll to keyBy(hashtag) windows to remove the single-operator bottleneck.
- Enable checkpointing with transactional Kafka sinks for exactly-once results across restarts.
- Send events later than the watermark to a side output so they're counted separately instead of dropped.