In Flink, what problem does WatermarkStrategy.withWatermarkAlignment solve?
answer
- backfill, not steady state
- one partition sprints, one crawls
- the minimum rule protects results, not memory
- the fast reader gets told to wait
- named group, allowed drift
basics
~20 sIt stops one source split from racing far ahead of the others in event time. Flink pauses reading from splits whose watermark exceeds the group's minimum by more than the configured drift, bounding the state that accumulates while slow splits catch up.
solid answer
~50 sWhen a job starts from the earliest offsets — a backfill, or recovery after a long outage — the sources read as fast as each partition allows. One partition may reach yesterday afternoon while another is still on last week. The operator watermark is the minimum, so it sits on the slow split, while records from the fast split keep opening windows and populating keyed state that cannot be cleaned up. State explodes, checkpoints slow, and the job may die before it catches up. `.withWatermarkAlignment("orders", Duration.ofSeconds(30), Duration.ofSeconds(1))` puts the sources in a named alignment group with a maximum allowed drift and an update interval. Sources whose watermark is more than the drift ahead of the group's minimum stop consuming until the rest catch up. Flink pauses individual splits, so the connector must implement split pause/resume; in Flink 2.3 a connector that does not fails the job unless you opt out. It is a backfill and recovery tool, not something a steady-state pipeline usually needs.
code
java · 5 linesWatermarkStrategy<Order> strategy = WatermarkStrategy
.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withTimestampAssigner((o, ts) -> o.getPlacedAtMillis())
.withIdleness(Duration.ofMinutes(1))
.withWatermarkAlignment("orders-group", Duration.ofSeconds(30), Duration.ofSeconds(1));go deeper
Know only that Flink sources can advance through event time at very different speeds, and that there is a setting to keep them roughly together.
Explain why the minimum-watermark rule leaves windows open when one split runs ahead, and that alignment pauses the fast reader until the group minimum catches up.
Recognise the incident from its symptoms — a backfill whose state and checkpoint duration climb while the job appears healthy — and configure a group name, drift and update interval that fit the state budget.
Decide the catch-up strategy at platform level: whether backfills run in batch execution mode, in bounded time chunks, or as aligned streaming jobs, and what the state and checkpoint budget per TaskManager actually allows.
## The skew that alignment addresses Event-time correctness is protected by taking the minimum watermark across inputs. That protects the results, but it says nothing about resource usage while inputs are unevenly advanced. Consider a job restarted from the earliest Kafka offsets across twelve partitions. Partition throughput is uneven — different sizes, different broker load, different deserialization cost — so after ten minutes of catch-up one reader is at event time `T` and another at `T - 3 days`. The operator watermark is `T - 3 days`. Every window between `T - 3 days` and `T` opened by records from the fast reader is *live*, because the watermark has not reached its end. So is every keyed timer. All of it sits in the state backend. The result is a job that appears to be making progress but whose state grows super-linearly with the skew: checkpoint duration climbs, RocksDB compaction saturates disk, and eventually the job fails or the checkpoint times out. This is the classic "the backfill worked in staging with one day of data and died on the real topic" incident. ## What alignment does `withWatermarkAlignment` gives the sources a shared notion of how far apart they may drift: ```java WatermarkStrategy<Order> strategy = WatermarkStrategy .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(20)) .withTimestampAssigner((o, ts) -> o.getPlacedAtMillis()) .withWatermarkAlignment("orders-group", Duration.ofSeconds(30), Duration.ofSeconds(1)); ``` The three arguments are the **watermark group** name, the **maximum allowed drift**, and the **update interval** at which readers report their watermark to the coordinator and receive the group's current minimum. The mechanism: source readers report their watermark to the `SourceCoordinator`; the coordinator computes the group-wide minimum and broadcasts it back. A reader whose watermark exceeds `groupMinimum + maxDrift` stops emitting records — it pauses consumption — until the minimum advances enough to bring it back inside the window. Readers that share a group name align together, so you can put several different sources in one group when they feed a common join or union. Since Flink 2.3 the pause decision is deliberately delayed. A split is compared not against its latest watermark but against the oldest of its last few sampled watermarks, held in a ring buffer sized by `pipeline.watermark-alignment.buffer-size` (default 3). The redesign (FLINK-37399) stops the announcement round trip from capping backlog speed at roughly one drift per update interval; the price is that a split can run for about three update intervals before it is paused. A buffer size of 0 restores the earlier behaviour. ## Granularity, and why it matters The original FLIP-182 implementation aligned at **source-subtask** granularity: if a subtask owned four splits and one of them was far ahead, the whole subtask paused, which could stall the slow splits it also owned — exactly the wrong outcome. Flink 1.17 added **split-level** alignment, where individual splits are paused via `SourceReader#pauseOrResumeSplits`. Connectors must implement that method; the Kafka source does. In Flink 2.3 there is no silent fallback: if a connector does not implement it, the first attempt to pause one of a reader's splits throws `UnsupportedOperationException` and fails the job. Setting `pipeline.watermark-alignment.allow-unaligned-source-splits` to `true` (default `false`, documented as due for removal) opts back into subtask-level-only alignment, which works properly only with one split per subtask. ## What it is not - **Not a fix for idle sources.** Idleness is a silent input pinning the minimum from *below*; alignment throttles inputs that are too far *above*. They solve opposite skews and are frequently configured together on the same strategy. - **Not a correctness feature.** Results are already correct without it, because the minimum rule protects them. Alignment protects state size and checkpoint health. - **Not free throughput.** Pausing the fast reader deliberately caps total ingestion at roughly the slowest split's rate. During a backfill that is the point; in steady state, a badly chosen drift can throttle a healthy pipeline for no benefit. ## Choosing the drift The drift needs to be comfortably larger than the normal event-time spread between partitions — otherwise readers thrash between paused and resumed and throughput suffers — and small enough that the state accumulated across the drift window fits your backend. A useful way to size it: estimate the keyed state produced per unit of event time, multiply by the candidate drift, and compare with the memory or disk budget per TaskManager. Backfills typically want a tight drift (seconds to a minute); steady-state pipelines often need no alignment at all. ## Alternatives and companions If alignment is unavailable or insufficient, the usual levers are: run the backfill in `RuntimeExecutionMode.BATCH`, where input is sorted by key and time and no unbounded window state accumulates; stage the backfill in bounded chunks by time range; reduce the parallel span by processing fewer partitions per run; or shrink per-key state with `StateTtlConfig` and smaller windows. Observing the problem is easier than fixing it after the fact — watch the spread between the minimum and maximum `currentOutputWatermark` across source subtasks, and alert when it widens during catch-up.
- Why does event-time skew across splits cause state growth rather than wrong results?Because the operator watermark is the minimum across inputs, so it stays on the slow split and results remain correct. But records from the fast split keep opening windows and registering timers whose deadlines the watermark has not reached, so none of that state can be cleaned up. Correctness is preserved by paying in memory.
- How does withWatermarkAlignment differ from withIdleness?They handle opposite skews. Idleness removes a silent input that is pinning the watermark minimum from below, so the job can progress. Alignment throttles an input that is racing ahead of the minimum, so state does not accumulate. Many production strategies configure both on the same source.
- What has to be true of a connector for split-level alignment to work?It must implement `SourceReader#pauseOrResumeSplits`, which arrived with the 1.17 split-level alignment work; the Kafka source does. In Flink 2.3 a connector without it fails the job with `UnsupportedOperationException` the first time a split must be paused, unless `pipeline.watermark-alignment.allow-unaligned-source-splits` is set to true, and that fallback aligns only whole subtasks, so keep one split per subtask.
- When would you deliberately not use alignment?In a steady-state pipeline where partitions track each other closely. Alignment caps throughput at roughly the slowest split's rate and adds coordinator round-trips, and a drift set too tight makes readers thrash between paused and resumed. It earns its keep during backfills and recovery from a long outage, not on a healthy live stream.
saying these in an interview costs you the question
- Thinks alignment makes results more correct
- Confuses it with withIdleness
- Assumes it exists in every Flink version and connector
- Believes it speeds up a backfill rather than throttling the fast readers
- Sets a drift of milliseconds and expects normal throughput