What does WatermarkStrategy.forBoundedOutOfOrderness do to a Flink job's watermarks?
answer
- max seen, minus a grace
- the grace is your guess about the data
- it is also your latency budget
- zero grace has its own strategy name
- off by one millisecond on purpose
basics
~20 sIt tracks the largest event timestamp seen and periodically emits a watermark of that maximum minus the configured bound (minus one millisecond). The bound is how long the job agrees to wait for out-of-order records, trading result latency against completeness.
solid answer
~50 s`WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(20))` builds a periodic generator that remembers the maximum event timestamp it has seen and, on each tick of `pipeline.auto-watermark-interval`, emits `maxTimestamp - 20s - 1ms`. The extra millisecond exists because a watermark of `T` claims completeness *inclusive* of `T`. The bound is a statement about your data: "records may arrive up to twenty seconds out of order". It is also directly a latency knob — every event-time window result is delayed by roughly the bound, because the watermark trails the newest record by that much. The sibling strategy `forMonotonousTimestamps()` is the same generator with a zero bound: watermark equals the largest timestamp seen minus one millisecond. Use it when a source genuinely emits in timestamp order, such as a single Kafka partition written by one ordered producer. Setting the bound too small drops real records as late; too large delays every result and holds window state longer. There is no correct value in the abstract — you measure the actual out-of-orderness distribution and pick a percentile.
code
java · 6 linesWatermarkStrategy<Order> strategy = WatermarkStrategy
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((order, recordTs) -> order.getPlacedAtMillis())
.withIdleness(Duration.ofMinutes(1));
DataStream<Order> orders = env.fromSource(kafkaSource, strategy, "orders");go deeper
Recall that this strategy tells Flink how far out of order records may be, and that a bigger value means waiting longer before window results are emitted.
Be able to state the arithmetic — maximum timestamp seen, minus the bound, minus a millisecond — and explain the periodic emit tick and how forMonotonousTimestamps is the zero-bound case.
Show how you would choose the bound from measured lateness rather than guesswork, and explain the silent-drop failure mode when the bound is too tight and the state cost when it is too loose.
Own it as a service-level decision: the bound is a published freshness-versus-completeness contract for downstream consumers, and it needs monitoring, a documented tail policy, and a story for what happens when an upstream producer degrades.
## What the strategy is made of A `WatermarkStrategy` bundles two things: a `TimestampAssigner`, which pulls the event timestamp out of a record, and a `WatermarkGenerator`, which decides what watermark to emit and when. The generator interface has two callbacks: - `onEvent(event, eventTimestamp, output)` — invoked for every record. - `onPeriodicEmit(output)` — invoked on a timer whose interval is the `pipeline.auto-watermark-interval` configuration (also settable through `ExecutionConfig#setAutoWatermarkInterval`). A generator that emits inside `onEvent` is called **punctuated** — appropriate when the stream itself carries explicit "batch complete" markers. A generator that emits inside `onPeriodicEmit` is **periodic**, which is what both built-in strategies do: `onEvent` merely records the maximum timestamp, and the timer emits. ## The arithmetic `forBoundedOutOfOrderness(bound)` keeps `maxTimestamp` = the largest event timestamp observed so far. On each periodic emit it outputs: ``` watermark = maxTimestamp - bound - 1 ms ``` The `- 1 ms` is not a rounding artefact. A watermark of `T` asserts completeness for timestamps *less than or equal to* `T`, so subtracting one millisecond makes the guarantee "everything strictly older than `maxTimestamp - bound`", which is what the bound is meant to express. `forMonotonousTimestamps()` is the identical generator with `bound = 0`, so the watermark is `maxTimestamp - 1 ms`. It asserts that timestamps never go backwards within the stream it is attached to. A third option, `WatermarkStrategy.noWatermarks()`, generates none at all — useful when a stream is only consumed by processing-time or non-windowed logic, or when you assign watermarks elsewhere. ## The dial you are actually turning Because windows fire when the watermark passes their end, the bound shows up directly as end-to-end latency: a five-minute window with a twenty-second bound emits its result roughly twenty seconds after the newest record for that window has been seen. Raising the bound to five minutes multiplies that delay and keeps every open window's state alive five minutes longer. Lowering the bound does the opposite and is the more dangerous direction, because the failure is silent: records that arrive behind the watermark are simply dropped by the window operator (counted only in the `numLateRecordsDropped` metric). A pipeline can lose a measurable slice of its input for months without anyone noticing. The honest way to choose is empirical. Run a job — or a simple stateless one — that computes `processingTimestamp - eventTimestamp` per record and look at the distribution. Pick a high percentile (p99 is a common starting point) as the bound, and let `allowedLateness` plus a late side output catch the tail rather than trying to make the bound cover the worst case ever observed. ## Where you attach it Two places, and they are not equivalent: ```java // preferred: watermarks generated inside the source, per split DataStream<Order> s = env.fromSource(kafkaSource, strategy, "orders"); // after the source: one generator over the interleaved stream DataStream<Order> s = env.fromSource(kafkaSource, WatermarkStrategy.noWatermarks(), "orders") .assignTimestampsAndWatermarks(strategy); ``` With the first form the connector runs one generator **per split** (per Kafka partition) and takes the minimum across the splits a subtask owns. With the second form, a subtask reading three partitions sees one interleaved stream and computes a single maximum — so a partition that is far behind gets its records declared late by the progress of the partitions that are ahead. Per-split watermarking in the source is the correct default in modern Flink. ## Parallelism and the WatermarkStrategy's other builders The strategy is fluent, and the other builders matter operationally: - `.withTimestampAssigner((event, recordTs) -> ...)` — required unless the source already attaches usable timestamps. - `.withIdleness(Duration)` — marks a stream idle after a period of silence so downstream operators exclude it from the watermark minimum. - `.withWatermarkAlignment(group, maxDrift)` — throttles splits or sources that race ahead of the rest of the group. ## Writing your own You implement `WatermarkGenerator` when the built-ins do not model your data. Two common cases: a punctuated generator that emits when a control record signals the end of an upstream batch, and a *self-adjusting* generator that measures observed lateness over a sliding period and widens or narrows the bound automatically. Custom generators are also where you can emit a watermark based on wall clock when a keyed source goes quiet, though `withIdleness` covers the common version of that need. ## Version note This answer assumes Flink 2.3. The `WatermarkStrategy` API arrived in Flink 1.11 and superseded `AssignerWithPeriodicWatermarks` and `AssignerWithPunctuatedWatermarks`; Flink 2.0 removed those interfaces, the `assignTimestampsAndWatermarks` overloads that accepted them, and Flink's own `Time` class, so every bound is now a `java.time.Duration`. If you meet code calling `assignTimestampsAndWatermarks(new BoundedOutOfOrdernessTimestampExtractor<>(Time.seconds(20)))`, it predates the `WatermarkStrategy` API and does not compile on Flink 2.x.
- Why is the emitted watermark maxTimestamp minus the bound minus one millisecond, rather than just minus the bound?Because a watermark of `T` claims completeness for timestamps `<= T`, inclusive. Subtracting one extra millisecond makes the claim cover only timestamps strictly older than `maxTimestamp - bound`, which is precisely what "records may be up to `bound` out of order" means. Without it the generator would over-claim by one millisecond.
- When would you write a custom WatermarkGenerator instead of using the built-ins?When the stream carries explicit completeness signals — an end-of-batch marker from an upstream system — a punctuated generator emitting from `onEvent` is exact rather than heuristic. The other case is a self-adjusting generator that measures observed lateness over a sliding period and widens the bound automatically when the source degrades.
- How do you pick the bound in practice?Measure it. Emit `processingTime - eventTime` per record and look at the distribution, then set the bound at a high percentile such as p99. Cover the long tail with `allowedLateness` and a late side output instead of inflating the bound, so that the common case stays fast while rare stragglers are still accounted for.
The bound is how long you hold the lift door open. A couple of seconds catches most stragglers cheaply; holding it for five minutes makes everyone already inside late.
saying these in an interview costs you the question
- Thinks the bound delays only late records, not all results
- Says the watermark equals the largest timestamp seen
- Sets a huge bound to be safe without mentioning latency or state cost
- Confuses the out-of-orderness bound with allowedLateness
- Assumes assignTimestampsAndWatermarks after the source behaves like per-split watermarking