A Flink job's keyed state grows without bound across millions of keys. How do you bound it?
answer
- Flink never retires a key by itself
- Two levers: a config and a callback
- One of them ignores watermarks entirely
- Best effort, not a scheduled delete
- clear() inside onTimer for the current key
basics
~20 sAttach a StateTtlConfig to the state descriptors so entries expire, and back it with explicit cleanup: a KeyedProcessFunction timer that calls state.clear() when a key goes idle. Flink never garbage-collects a key just because it stopped appearing.
solid answer
~50 sFlink creates keyed state on first touch and never removes it on its own, so an unbounded key space means unbounded state. Two mechanisms bound it, and mature jobs use both. **State TTL**: build a `StateTtlConfig` and call `descriptor.enableTimeToLive(ttlConfig)`. TTL is *processing-time only*; `UpdateType.OnCreateAndWrite` (the default) refreshes the clock on writes, `OnReadAndWrite` also on reads; `StateVisibility.NeverReturnExpired` (the default) hides expired values even before they are physically gone. Cleanup is best-effort: expired values are dropped on read, plus background cleanup — incremental iteration on the heap backend, a compaction filter on RocksDB. If a key is never read and no records flow, its expired state persists. **Explicit timers**: in a `KeyedProcessFunction`, register an event-time or processing-time timer and call `clear()` on every state handle in `onTimer`. That is deterministic and event-time aware, which TTL is not.
code
java · 10 linesStateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Duration.ofHours(24))
.setUpdateType(StateTtlConfig.UpdateType.OnReadAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.cleanupInRocksdbCompactFilter(1000, Duration.ofHours(1))
.build();
ValueStateDescriptor<Profile> descriptor =
new ValueStateDescriptor<>("profile", Profile.class);
descriptor.enableTimeToLive(ttlConfig);go deeper
Recall that keyed state is created per key and never removed automatically, and that StateTtlConfig plus descriptor.enableTimeToLive is the built-in way to expire it.
Explain the mechanics: the update type and visibility settings, the fact that TTL runs on processing time, and how background cleanup differs between the heap backend's incremental iteration and RocksDB's compaction filter.
Diagnose from evidence — climbing checkpoint size, key cardinality, which operator owns the growth — then combine TTL with KeyedProcessFunction timers, and state plainly that TTL reclaims disk on a best-effort basis rather than on schedule.
Own retention as a policy question: what the business retention deadline actually is, whether the key space itself should be redesigned, how state size feeds checkpoint SLAs and cluster cost, and where a privacy requirement forces NeverReturnExpired semantics.
## Why the state grows Keyed state is created lazily, the first time a record for a key touches a state handle, and Flink has no notion of a key being "finished". If your key space is unbounded — session ids, order ids, device ids that churn — then every key that ever appeared keeps its entry forever. On `HashMapStateBackend` this shows up as heap pressure and eventually OOM; on `EmbeddedRocksDBStateBackend` it shows up as growing local disk use, growing checkpoint size and slowly degrading access latency as the LSM tree deepens. The diagnosis is usually straightforward: checkpoint size climbing monotonically with no plateau, while throughput is flat. What matters is the fix. ## Mechanism one: state TTL A time-to-live can be assigned to keyed state of any type. You build a `StateTtlConfig` and enable it on the descriptor: ```java StateTtlConfig ttl = StateTtlConfig .newBuilder(Duration.ofHours(24)) .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired) .build(); descriptor.enableTimeToLive(ttl); ``` Four properties decide whether it behaves the way you expect. **It is processing-time only.** TTL is measured against wall clock, not against watermarks. A job replaying two years of history in an hour will not expire anything, because only an hour of processing time has passed. If your retention rule is genuinely in event time, TTL is the wrong tool. **The update type sets what refreshes the clock.** `OnCreateAndWrite` (default) refreshes on creation and writes only; `OnReadAndWrite` also refreshes on reads, turning the TTL into an idle-timeout rather than an age limit. **Visibility is separate from cleanup.** `NeverReturnExpired` (default) makes expired state behave as if it does not exist even while the bytes remain — the right choice when data must become unreadable after a retention deadline. `ReturnExpiredIfNotCleanedUp` will hand back an expired value that has not been collected yet. **Collections expire per entry.** All state collection types support per-entry TTLs: list elements and map entries expire independently. That is a strong argument for `MapState` over a `ValueState` holding a map. ## How the cleanup actually happens This is where people get burned. By default, expired values are removed *on read* — `ValueState#value` drops them — and collected in the background if the configured backend supports it. Background cleanup can be turned off with `disableCleanupInBackground()`. The two backends do background cleanup differently: - **Heap backend: incremental cleanup.** The backend keeps a lazy global iterator over the state's entries and advances it on state access and, optionally, on record processing. `cleanupIncrementally(10, true)` checks 10 entries per trigger and additionally triggers per record. The default background cleanup for the heap backend checks 5 entries without cleanup per record processed. Note this is heap-only — setting it for RocksDB has no effect. It also adds latency to record processing, and with synchronous snapshotting the iterator keeps a copy of all keys. - **RocksDB: compaction filter.** A Flink-specific compaction filter runs during RocksDB's own asynchronous compactions and drops entries whose expiry timestamp has passed. `cleanupInRocksdbCompactFilter(queryTimeAfterNumEntries, periodicCompactionTime)` controls how often the filter refreshes its notion of "now" from Flink — the default queries the current timestamp every 1000 entries — and how old a file must be before it is picked for compaction anyway, defaulting to 30 days. Refreshing more often speeds cleanup but costs JNI calls into native code; periodic compaction speeds up removal of rarely-accessed entries at the cost of more compaction work. - **Full snapshot cleanup.** `cleanupFullSnapshot()` excludes expired entries when a full snapshot is taken, shrinking the snapshot. It does not clean local state, and it does not apply to RocksDB incremental checkpointing. The consequence to state out loud in an interview: **if no access happens to the state and no records are processed, expired state persists.** TTL is not a guarantee that disk shrinks by the deadline. It is a guarantee about visibility plus a best-effort reclaim. ## Mechanism two: explicit timers The deterministic tool is a `KeyedProcessFunction`. From `processElement` you register a timer on `ctx.timerService()` — `registerEventTimeTimer(t)` or `registerProcessingTimeTimer(t)` — and implement `onTimer(long timestamp, OnTimerContext ctx, Collector<OUT> out)`, where all state is again scoped to the key the timer was created for. Clearing every handle there physically retires the key. Event-time timers fire when the watermark advances to or beyond their timestamp, so this is the mechanism that respects event time and behaves identically on replay. The `TimerService` deduplicates timers per key and timestamp, so the standard "idle timeout" pattern — store the last-modified time in a `ValueState`, register a timer for `lastModified + gap`, and on fire check whether the stored timestamp still matches before clearing — is safe. Flink synchronises `onTimer()` and `processElement()` on the same key, so no concurrent-modification handling is needed. Under RocksDB, timers are themselves stored in RocksDB by default (`state.backend.rocksdb.timer-service.factory`, default `rocksdb`), which scales but is not free — a timer per key is state too. ## Costs and caveats to raise TTL is not free storage-wise: the backend stores a last-modification timestamp alongside each user value. RocksDB adds 8 bytes per stored value, list entry or map entry; the heap backend adds an object reference plus a primitive long. The TTL configuration is not part of checkpoints or savepoints — it is how the currently running job treats the data — and it is not recommended to restore state while lengthening a short TTL, which can produce data errors. Since Flink 2.2.0, switching TTL on or off for existing state migrates cleanly rather than throwing a `StateMigrationException` on restore. ## What a strong answer sounds like Diagnose first — checkpoint size trend, key cardinality, which operator owns the growth. Then bound it: TTL for coarse retention and for the privacy requirement, explicit timers for the semantic "this session is over" rule, `MapState` rather than blobs so entries expire individually, and a hard look at whether the key itself should be coarser. Finish with the honest caveat that TTL is best-effort reclamation, not a scheduled delete.
- Why does state TTL not shrink disk usage by the deadline you configured?Because cleanup is best-effort. Expired values are removed when the state is read, and otherwise only by background cleanup — incremental iteration on the heap backend, a compaction filter during RocksDB's own compactions. If a key is never read again and no records are processed for it, its expired entry stays on disk until a compaction happens to touch it. Visibility is immediate; physical reclamation is not.
- Why can't you use state TTL to implement an event-time retention rule?TTL is measured only against processing time — the wall clock — never against watermarks. During a historical replay, a job can process two years of event time in an hour of wall clock, so nothing expires. Conversely a stalled job expires state that is still semantically current. For event-time retention, use a `KeyedProcessFunction` with `registerEventTimeTimer` and clear the state in `onTimer`, which advances with the watermark.
- What does enabling TTL cost in state size?The backend stores the last-modification timestamp alongside each user value. `EmbeddedRocksDBStateBackend` adds roughly 8 bytes per stored value, list entry or map entry; `HashMapStateBackend` adds an extra Java object holding a reference to the user value plus a primitive long. On a job with hundreds of millions of small entries that overhead is real, so TTL is a trade rather than a free win.
- How do you clean up state for a key without holding one timer per key?Coalesce the timers. Round the timer timestamp to a coarse bucket — the next full minute, say — so many keys share a firing instant, and store the precise deadline in state; on fire, re-check the stored deadline and either clear the state or register the next coalesced timer. The `TimerService` deduplicates timers per key and timestamp, so coalescing collapses repeated re-registrations for the same key into one.
saying these in an interview costs you the question
- Thinks Flink drops a key's state when it stops appearing
- Says state TTL expires against event time or watermarks
- Expects TTL to delete entries exactly at the deadline
- Believes enabling TTL is free in state size
- Says a savepoint restore re-applies the old TTL configuration