skip to content

In Flink, what does a watermark with timestamp T assert about the stream?

level: middleimportance: must knowfreq 74%

answer

  1. a promise, not a measurement
  2. it travels with the records
  3. two inputs, take the pessimistic one
  4. windows fire when it passes the end
  5. break the promise and you are late

basics

~20 s

A watermark of T is the pipeline's declaration that no record with an event timestamp at or before T is expected any more. It travels in-band with the records, advances event time, and fires event-time windows and timers whose deadline it passes.

solid answer

~50 s

A watermark is a special marker that flows through the stream alongside records. A watermark of value `T` asserts: *no further record with an event timestamp less than or equal to `T` should arrive*. It is a claim about completeness, not about progress. That assertion is what makes event-time computation possible. When the watermark reaches a window's last timestamp — its end minus one millisecond — the window is considered complete and fires. Event-time timers registered in a `ProcessFunction` fire on the same rule. Watermarks are monotonic per channel: Flink never moves a watermark backwards. An operator with several input channels takes the **minimum** of its inputs' watermarks as its own current watermark, because it cannot claim completeness beyond its slowest input. Crucially the assertion is a heuristic supplied by *you*, via the `WatermarkStrategy`. If a record shows up whose timestamp is behind the current watermark, the assertion was wrong: that record is **late**, and by default a window drops it.

code

text · 5 lines
text
input channel A:  ... r(11:59:58)  r(11:59:55)  W(11:59:40)
input channel B:  ... r(11:52:10)  W(11:51:50)

operator current watermark = min(11:59:40, 11:51:50) = 11:51:50
  -> window [11:50, 11:55) has NOT fired yet

go deeper

for a junior

Recall that a watermark is a marker in the stream saying 'no more records older than this are expected', and that windows fire when it passes their end.

for a middle

Explain the mechanics precisely: monotonic, in-band, minimum across input channels, generated by your WatermarkStrategy rather than inferred by the engine, and what makes a record late.

for a senior

Demonstrate you have operated this. Talk about watermark lag as a health metric, the silent drop of late records and the numLateRecordsDropped counter, and how one slow input channel stalls an entire job.

for a principal

Frame the watermark as an explicit contract between the data producers and the pipeline. Own where that bound is decided, how it is validated against real out-of-orderness, and what the business agrees to when it accepts a completeness heuristic.

## The problem watermarks solve Under event time, a window's contents are defined by the records' own timestamps. But a stream is unbounded and the network reorders things, so the engine can never *prove* that no more records belong to `[12:00, 12:05)`. It has to decide to stop waiting. A watermark is exactly that decision, made explicit and made schedulable. ## The contract A watermark is an element in the stream — not a message on a side channel — carrying a single long timestamp. When an operator receives `Watermark(T)`, it may assume that **no further element with a timestamp `<= T` will arrive on that channel**. From this the operator derives that any event-time work whose deadline is `<= T` is now complete and can be released. Two properties matter: - **Monotonicity.** Watermarks never decrease. Flink's built-in generators track the maximum timestamp seen and only ever emit forward. - **In-band flow.** Because they travel with the records, a watermark can never overtake the records emitted before it on the same channel. Watermarks are not themselves part of a checkpoint: after recovery the sources regenerate them from the replayed records. ## What fires, and when A tumbling event-time window `[start, end)` registers a timer at `end - 1` (its `maxTimestamp()`). When the operator's watermark reaches or passes that value, the window fires and, absent allowed lateness, its state is purged. A `KeyedProcessFunction` timer registered at `t` fires when the watermark passes `t`; inside `onTimer` you can read `ctx.timerService().currentWatermark()`. Note the direction of causality: the *watermark* moves event time forward, not the data. A stream that stops producing records also stops producing watermarks, and everything downstream freezes — which is why source idleness is a real operational concern. ## The minimum rule Most operators have more than one input channel: a `keyBy` shuffle gives every downstream subtask a channel from every upstream subtask, and a union or connect gives an operator two logical inputs. Flink tracks the last watermark per channel and sets the operator's current watermark to the **minimum** across them. It then forwards that minimum downstream. This is the only safe rule — claiming completeness at `T` while one input has only reached `T - 10 minutes` would silently drop ten minutes of records. It is also the source of the most common event-time incident: one silent or slow input channel pins the whole operator's watermark and no window in the job ever fires. ## Where the assertion comes from The engine does not invent watermarks. You supply a `WatermarkStrategy`, which combines a `TimestampAssigner` (how to read a record's time) with a `WatermarkGenerator`. The generator has two hooks: `onEvent(event, timestamp, output)`, called for every record, and `onPeriodicEmit(output)`, called on a timer whose interval is `pipeline.auto-watermark-interval`. A generator that emits from `onEvent` is *punctuated*; one that emits from `onPeriodicEmit` is *periodic*, which is what the built-in strategies do. So the watermark's assertion is only as good as your model of the data. `forBoundedOutOfOrderness(Duration.ofSeconds(20))` says "I believe records are never more than twenty seconds out of order". Reality is under no obligation to comply. ## When the assertion is violated A record whose timestamp is at or before the current watermark is **late**. Its fate depends on configuration: 1. In a window with no allowed lateness, it is dropped — silently, but counted by the `numLateRecordsDropped` metric on the window operator. 2. With `allowedLateness(Duration)`, the window's state is kept past the watermark deadline, and the late record triggers an *additional* firing that emits an updated result. 3. With `sideOutputLateData(outputTag)`, records too late even for that grace period are routed to a side output stream instead of being discarded, so you can log, alert, or reprocess them. Because dropping is the default and it is silent from the application's point of view, alerting on the dropped-records metric is one of the highest-value things you can do to an event-time job. ## Observing watermarks Every operator exposes `currentInputWatermark` and `currentOutputWatermark` gauges (two-input operators also expose `currentInput1Watermark` / `currentInput2Watermark`), and the Flink Web UI shows the low watermark per subtask. Watermark lag — wall clock minus watermark — is the single best health signal for an event-time pipeline: it tells you both how far behind the job is and how long window results will be delayed. At the end of a bounded stream Flink emits `Watermark.MAX_WATERMARK` (`Long.MAX_VALUE`), which flushes every remaining window and timer; a job running in `RuntimeExecutionMode.BATCH` does not rely on intermediate watermarks at all.

  • Why does an operator take the minimum watermark across its input channels rather than the maximum?
    Because the watermark is a completeness claim for the operator as a whole. If one channel has only reached 12:00 and another 12:10, advancing to 12:10 would declare windows complete while ten minutes of records are still in flight on the slow channel — and those records would then be late and dropped. The minimum is the only value the operator can honestly assert.
  • How does a window operator know exactly when to fire?
    The event-time trigger registers a timer at the window's `maxTimestamp()`, which is `end - 1` millisecond. When the operator's current watermark reaches or exceeds that value the timer fires, the window function runs, and — unless allowed lateness keeps it alive — the window state is purged.
  • What happens to watermarks when a bounded source finishes, or when the job runs in BATCH execution mode?
    Flink emits `Watermark.MAX_WATERMARK` (`Long.MAX_VALUE`), which flushes every open window and pending event-time timer. In `RuntimeExecutionMode.BATCH` intermediate watermarks are not needed at all: input is sorted by key and time, and the final max watermark fires everything at the end.
  • How would you monitor whether watermarks are healthy in production?
    Track watermark lag — wall clock minus `currentOutputWatermark` — per subtask, and alert when it grows. Pair it with the window operator's `numLateRecordsDropped` counter: a rising lag means results are delayed, and rising drops mean the out-of-orderness bound is too optimistic for the real data.

A watermark is the gate agent announcing that boarding is closed. Stragglers may still show up at the desk, but the aircraft has committed to stop waiting for them.

saying these in an interview costs you the question

  • Says a watermark means all records up to T have been processed
  • Thinks watermarks are sent on a control channel outside the data stream
  • Claims an operator averages the watermarks of its inputs
  • Believes Flink computes watermarks automatically from the data
  • Says a late record rolls the watermark backwards

context