A Flink job with SlidingEventTimeWindows of 24 hours and a 1-minute slide exhausts state — why?
answer
- divide the size by the slide
- how many windows are open at once
- each record joins every one of them
- the window function decides what gets copied
- sometimes the right answer is no window at all
basics
~20 sA 24-hour window sliding every minute keeps 1440 windows open per key at once, and every record is assigned to all of them. State is multiplied roughly 1440-fold per key, so checkpoints balloon and the job runs out of memory or disk.
solid answer
~50 sSliding windows overlap, and the overlap factor is `size / slide` — here 1440 minutes of window divided by a 1-minute slide. Flink's window operator keeps **one entry per window per key**, so at any instant each key has 1440 open windows. With a `ProcessWindowFunction` that means every record is buffered 1440 times; with a `ReduceFunction` or `AggregateFunction` it means 1440 accumulators per key rather than 1440 copies of every element — far cheaper, but still 1440× a tumbling window. Multiply by the key cardinality and the state, checkpoint size and checkpoint duration all follow. The first fix is to stop buffering: use `aggregate(...)`, optionally combined with a `ProcessWindowFunction` for window metadata. If that is still too much, the sliding assigner is the wrong tool — a `KeyedProcessFunction` holding per-minute partial sums per key with a per-minute timer, or a coarser slide, gives the same metric at a fraction of the cost.
code
java · 9 lines// 1440 open windows per key, every element buffered in all of them
input.keyBy(Event::getUserId)
.window(SlidingEventTimeWindows.of(Duration.ofHours(24), Duration.ofMinutes(1)))
.process(new TrailingStats());
// still 1440 windows per key, but one accumulator each
input.keyBy(Event::getUserId)
.window(SlidingEventTimeWindows.of(Duration.ofHours(24), Duration.ofMinutes(1)))
.aggregate(new TrailingAgg(), new EmitWithBounds());go deeper
Be able to compute how many windows overlap from the size and slide, and to say that each record is assigned to all of them.
Explain why a ProcessWindowFunction multiplies element copies while an AggregateFunction multiplies only accumulators, and estimate state from key cardinality times the overlap factor.
Walk a concrete diagnosis: read the assigner, check per-operator checkpoint size in the UI, identify the window function, then give a ranked set of fixes including abandoning the assigner for a KeyedProcessFunction.
Own the upstream question — whether the business metric really needs a trailing 24 hours refreshed each minute, what that costs in cluster capacity and recovery time, and what refresh contract you are prepared to guarantee.
## Reading the symptom The complaint usually arrives as one of: checkpoints that take minutes and then time out, a RocksDB state directory growing to hundreds of gigabytes, `TaskManager` heap exhaustion on a `HashMapStateBackend`, or backpressure that never clears. The assigner configuration is the first thing to look at, because sliding windows are the most reliable way to multiply state by a large constant without noticing. ## The arithmetic A sliding window assigner takes a **size** and a **slide**. A new window starts every slide interval and each window lasts for the size, so at any instant the number of simultaneously open windows per key is `ceil(size / slide)`. For size 24 hours and slide 1 minute that is 1440. A `WindowAssigner` assigns each incoming record to *every* window it belongs to. Flink's own documentation puts it plainly: Flink creates one copy of each element per window to which it belongs. So a record that would occupy one slot in a tumbling window occupies 1440 slots here. If you have 200,000 active keys and each key sees a handful of records a minute, you now have 288 million window entries live at once, and every one of them is checkpointed. ## The window function changes the constant, not the exponent The multiplier applies to whatever the operator stores per window: - With a **`ProcessWindowFunction` alone**, the operator must buffer every element until firing, so the state is `records_per_key_per_day × 1440 × keys`. This is catastrophic. - With a **`ReduceFunction` or `AggregateFunction`**, the operator holds one value or accumulator per window, so the state is `1440 × keys × sizeof(accumulator)`. Much smaller, but still 1440 times a tumbling window of the same size, and if your accumulator is a fat object — a set of distinct IDs, a sketch, a map — that constant hurts. So the first diagnostic question is: which window function is attached? Switching from `.process(...)` to `.aggregate(agg, processFn)` frequently reduces state by two or three orders of magnitude on its own, and it is a small code change. ## Confirming the diagnosis - Read the assigner from the job graph and compute `size / slide`. - Look at checkpoint size per operator in the Flink web UI — the window operator will dominate. - Check whether the window function is buffering (`process`) or incremental (`reduce` / `aggregate`). - Estimate key cardinality; state is linear in it and the assigner multiplier compounds on top. - If the state backend is `EmbeddedRocksDBStateBackend`, look at disk usage and compaction; if it is `HashMapStateBackend`, look at heap and GC. RocksDB relocates the problem to disk, it does not shrink it. ## Remedies, roughly in order 1. **Make the aggregation incremental.** Replace a standalone `ProcessWindowFunction` with `aggregate(...)`, or the combined `aggregate(agg, processFn)` form if you need window bounds in the output. 2. **Widen the slide.** Ask what refresh interval the consumer actually needs. Moving from a 1-minute to a 5-minute slide divides state by five and is very often invisible to the dashboard consuming it. 3. **Shrink the accumulator.** If the accumulator holds a set for distinct counting, an approximate sketch turns an unbounded per-window structure into a fixed-size one. 4. **Abandon the assigner.** For "trailing N, refreshed every M", a `KeyedProcessFunction` holding per-minute partial aggregates in a `MapState` plus a per-minute timer stores at most about 1440 *small numbers* per key instead of 1440 *windows* per key, and emits by summing the partials still inside the trailing 24 hours while deleting the older ones. This is the standard rewrite for large trailing metrics and is usually a large win. 5. **Reconsider the metric.** If the requirement is really cumulative-since-midnight rather than trailing-24-hours, a tumbling daily window with early firing — or the `CUMULATE` windowing table-valued function in Flink SQL — expresses it directly with a fraction of the state. ## What not to reach for Raising parallelism does not reduce total state; it spreads the same state across more subtasks, which helps heap pressure but not checkpoint volume. Switching state backends does not reduce state either — `EmbeddedRocksDBStateBackend` lets state exceed memory by spilling to local disk, which buys headroom rather than solving the design problem. Increasing checkpoint timeouts hides the symptom until recovery time becomes the new outage. ## Answering well Do the arithmetic out loud (`1440 open windows per key`), name the copy-per-window behaviour as the mechanism, distinguish the buffering from the incremental function, then give a ranked set of fixes ending with "and if the numbers still do not work, this is not a window, it is a `KeyedProcessFunction`." That last move — knowing when to leave the built-in abstraction — is what the question is really testing.
- Would switching from HashMapStateBackend to EmbeddedRocksDBStateBackend fix this?It buys headroom, not a fix. RocksDB keeps state on local disk rather than on the JVM heap, so the job stops dying of heap exhaustion and can hold state larger than memory. The volume is unchanged, so checkpoint size, checkpoint duration and recovery time stay bad, and RocksDB adds serialization and compaction overhead per access.
- Does increasing the operator's parallelism reduce the state problem?No. Parallelism redistributes the same keyed state across more subtasks, which relieves per-TaskManager memory or disk pressure and can improve throughput. Total state across the cluster is unchanged, so aggregate checkpoint size stays the same. It also does nothing at all if the state is concentrated on a few hot keys, since one key never splits.
- How would you tell whether the fix should be a wider slide or a rewrite to KeyedProcessFunction?Start from what the consumer needs. If a 5- or 15-minute refresh is acceptable, widening the slide is a one-line change that divides state proportionally and should be tried first. Move to a KeyedProcessFunction when the refresh interval is genuinely fine-grained, when the accumulator is large, or when you need a bound the assigner cannot express — such as capping per-key state explicitly.
It is like photocopying every incoming invoice into 1440 different folders so that each folder holds a different trailing day. The filing cabinet fills up long before anyone reads a folder.
saying these in an interview costs you the question
- Says each record is stored once regardless of overlap
- Proposes more parallelism as a way to shrink total state
- Claims RocksDB reduces the amount of state rather than relocating it
- Ignores which window function is attached when estimating state
- Treats raising the checkpoint timeout as the fix