Growing Flink state rarely shows up as an error. The Full Checkpoint Data Size on the Flink web UI’s Checkpoints tab is a little larger every day, each checkpoint takes longer to write, the local disks under the RocksDB directories fill up, and TaskManagers run short of memory or get killed by their container limits. The job keeps running until one of those limits is reached.
Flink state is the data an operator keeps between events, such as a running count, a buffered window or a timer, held by the state backend and copied into every checkpoint.
To fix it, find which operator and which named state hold the growth, decide whether that growth is a bug or the job doing what it was built to do, then change one setting at a time. Config keys and defaults on this page were checked against the Apache Flink 2.3 documentation, and they can differ in other releases, so check the documentation for the version in use. The large state tuning guide covers the operational side of large state, and this page follows it where it applies.
This page is about state that grows. For a different problem, start here instead:
- Checkpoints failing or expiring now: how to fix Flink checkpoints that fail or time out.
- A job cycling between RESTARTING and RUNNING: how to fix a Flink job stuck in a restart loop.
- State stores in Kafka Streams rather than Flink jobs: Kafka Streams.
For where state fits in Flink as a whole, see the complete Flink guide. The steps below run in this order.
- Confirm the state is really growing
- Find the operator and the state that hold it
- Decide whether the growth is a bug
- Apply fixes, one change at a time
Confirm the state is really growing
Rule out these three causes before changing the job.
Read the full size, not the delta
With incremental checkpoints on, each checkpoint writes only what changed since the previous one. The state backends documentation says that once incremental checkpoints are enabled, the Checkpointed Data Size in the web UI represents only the delta of that checkpoint, not the full state size. The checkpoint monitoring documentation lists Full Checkpoint Data Size beside it as the accumulated size over all acknowledged subtasks. Trend the full figure.
The same split exists in the job’s checkpoint metrics, which the metrics documentation lists as gauges in bytes.
lastCheckpointFullSize
lastCheckpointSize
The delta figure can differ from the full size when incremental checkpoints or changelog are enabled. A flat lastCheckpointSize next to a rising lastCheckpointFullSize is still growth.
Compare live data with files on disk
A full local disk under RocksDB is not proof that the state grew. RocksDB keeps data in SST files and merges them through compaction, so files on disk can hold more than the live state. The Flink configuration reference describes the live data estimate as usually smaller than the SST files size due to space amplification. Both are RocksDB native metrics, off by default.
state.backend.rocksdb.metrics.estimate-live-data-size: true
state.backend.rocksdb.metrics.total-sst-files-size: true
The same reference warns that the SST files size metric may slow down online queries when there are many files. If the live data estimate is flat while SST files grow, look at compaction and disk, not at the job’s logic.
Check what is filling checkpoint storage
Checkpoint storage can fill for reasons unrelated to the running job’s state. The checkpoints documentation says that with RETAIN_ON_CANCELLATION, the checkpoint is kept when the job is cancelled and has to be cleaned up manually. Old retained checkpoints and savepoints from earlier deployments add up in the bucket while the live job’s full checkpoint size stays flat. The retention setting is this key, and its value decides whether a cancelled job leaves its checkpoint behind.
execution.checkpointing.externalized-checkpoint-retention: DELETE_ON_CANCELLATION
Find the operator and the state that hold it
In the Flink web UI’s checkpoint History tab, the More details link on a completed checkpoint gives a minimum, average and maximum summary over all its operators, and the detailed numbers for every subtask, as the checkpoint monitoring documentation describes. Compare two checkpoints a day or more apart. The operator whose size grew between them is where to look.
Within that operator, compare subtasks. Flink splits keyed state into key groups, the units it uses to spread keys across parallel subtasks, as the stateful stream processing documentation explains. If one subtask carries far more state than the rest, the state is concentrated in the keys that subtask owns, which can mean a few hot keys.
An operator can hold several named states. The large state tuning guide says that with RocksDB each state corresponds to one column family. Exposing the column family as a metric variable breaks the RocksDB metrics down by state, and the estimated key count shows which state has a key space that keeps rising.
state.backend.rocksdb.metrics.column-family-as-variable: true
state.backend.rocksdb.metrics.estimate-num-keys: true
state.backend.rocksdb.metrics.estimate-live-data-size: true
Flink can also sample the key and value size of keyed state, for the standard backends and custom ones. The metrics documentation says this is disabled by default, samples every 100 accesses unless changed, and may impact performance.
state.size-track.keyed-state-enabled: true
A rising key count points at the key space. A flat key count with a rising size points at values that keep getting bigger, such as a list or map that is appended to and never trimmed.
Decide whether the growth is a bug
The table below sets the two cases side by side.
Four patterns that make state grow without end
Four patterns documented by Apache Flink cause growth that never levels off.
Keyed state with no TTL and an unbounded key space
Keyed state is scoped to the key of the current element. The working with state documentation describes ValueState as possibly holding one value for each key the operation sees. Without a TTL or explicit clearing, an entry stays for as long as the job runs. If the key includes a value that is rarely or never seen twice, such as a session, request or order ID, the number of entries grows with every new value. The estimated key count for that state rises and never falls back.
Windows that never close
The windows documentation says a window is completely removed when time passes its end timestamp plus the allowed lateness, and that Flink guarantees removal only for time-based windows. A global window has no natural end, so its state stays until a custom trigger clears it.
An event-time window also stays open while the watermark is held back. The watermark documentation explains that the watermark is the minimum over all parallel inputs, so one idle Kafka partition or source split holds it back for the whole operator. Windows pile up behind it, and the operator’s watermark stops moving while its input keeps arriving. An idle partition can come from skewed partition keys on the Kafka side, where a few partitions carry most of the records.
Timers that are never deleted
Timers are state too. The process function documentation says timers are fault tolerant and checkpointed along with the state of the application, and that large numbers of timers can increase checkpointing time. A function that registers a new timer on every event and never deletes the old one, or registers timers at millisecond resolution, can carry far more timers than keys.
SQL queries with no state TTL
Flink SQL keeps state for joins, aggregations and deduplication. The table configuration reference gives the idle state retention a default of 0, which means state is never cleaned up. A streaming join or GROUP BY over a key with unbounded cardinality grows for as long as the query runs.
table.exec.state.ttl: 0
When the growth is legitimate
Some growth is the job working as designed. The windows documentation says Flink keeps one copy of each element per window it belongs to, so sliding windows store several copies of each element. The window function sets how much each window holds.
- A ReduceFunction or AggregateFunction stores one value per window.
- A ProcessWindowFunction on its own accumulates every element.
- An Evictor prevents pre-aggregation, so the window keeps every element.
More customers, more devices or a longer window also mean more state. Legitimate growth levels off once the window or retention period is full, or follows the business volume.
That case calls for a backend and checkpoint setup that can carry it, covered below. For examples of production Flink workloads at scale, see the collected Flink use cases.
Fixes, one change at a time
Change one thing, then watch the same Full Checkpoint Data Size and key count trends for at least one full retention period before the next change.
Add state TTL to keyed state
State TTL gives each keyed state entry an automatic expiry, for example one hour after it was last written. The working with state documentation says a TTL can be assigned to keyed state of any type, list elements and map entries expire individually, and expired values are cleaned up on a best effort basis.
StateTtlConfig ttlConfig = StateTtlConfig
.newBuilder(Duration.ofHours(1))
.setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
.setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
.build();
stateDescriptor.enableTimeToLive(ttlConfig);
The same page sets the limits a reader needs before relying on it.
- Only TTLs in reference to processing time are supported, so a TTL cannot express “an hour of event time”.
- On RocksDB, expired entries are removed by a Flink compaction filter during RocksDB compaction. Periodic compaction makes rarely accessed files pass through that filter, with a default of 30 days.
- On the heap backend, incremental cleanup runs on state access and record processing. If no access happens to the state or no records are processed, expired state persists.
- Cleanup in full snapshot does not apply to incremental checkpointing in the RocksDB state backend.
- The TTL configuration is not stored in checkpoints or savepoints, and the documentation does not recommend restoring with a TTL changed from a short value to a long one.
Clear state and delete timers explicitly
Where the job knows when a key is finished, such as an order that completes or a session that ends, clearing the state at that moment is more precise than waiting for a TTL. Timers are the usual tool for it. The Flink event-driven applications guide walks through clearing a key’s state when its timer fires, and deleting the registered timer and the state that remembers it when the key is cleared early.
ctx.timerService().deleteProcessingTimeTimer(timestampOfTimerToStop);
Timer resolution sets how many timers there are. The process function documentation says the TimerService keeps at most one timer per key and timestamp, so rounding target times coalesces timers. At one timer per key per second, that is up to 86,400 timers per key per day. At one per hour, it is at most 24.
long coalescedTime = ((ctx.timestamp() + timeout) / 1000) * 1000;
ctx.timerService().registerProcessingTimeTimer(coalescedTime);
Pick the time domain to match the purpose. The Flink notions of time documentation says processing time gives the best performance and lowest latency but does not provide determinism, while event time can give consistent and deterministic results that hold when the job replays data. Processing-time timers suit cleanup and timeouts. Event-time timers suit results that must be the same on replay, such as window output and late data handling.
Make windows close
For a global window, add a trigger that fires and purges. For event-time windows behind an idle source, mark the source idle so the watermark can advance without it. The watermark documentation shows the helper.
WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(20))
.withIdleness(Duration.ofMinutes(1));
Keep the allowed lateness as short as the use case permits, since a window is only removed after its end plus that lateness. Where a ProcessWindowFunction is used only to compute an aggregate, combining it with a ReduceFunction or AggregateFunction stores one value per window instead of every element.
Set a state TTL for SQL jobs
For Flink SQL, set an idle state retention that is longer than the longest gap between updates for a key that the query must still remember.
table.exec.state.ttl: 1 h
The table configuration reference notes that cleaning up state requires additional bookkeeping overhead.
Bound the key space
A key with bounded cardinality bounds the state. Where the key carries a unique value only to tell events apart, keying by the entity the state describes, such as the customer or device, keeps one entry per entity instead of one per event. Where the business really does have an unbounded key space, the key still needs a TTL or an explicit end so old keys leave state.
Move large state to RocksDB with incremental checkpoints
When the growth is legitimate, the state has to live somewhere that can hold it.
The state backends documentation says the HashMapStateBackend keeps state as objects on the Java heap, limited by available memory, and points jobs with very large state, long windows and large key or value states to the EmbeddedRocksDBStateBackend, whose limit is the available disk space. It also says each RocksDB state access goes through serialization and may read from disk, with average performance an order of magnitude slower than the heap backends.
state.backend.type: rocksdb
execution.checkpointing.incremental: true
Incremental checkpoints reduce checkpoint time. The large state tuning guide says incremental checkpoints should be one of the first considerations when reducing checkpoint time, because they record only the changes since the previous completed checkpoint.
RocksDB performance also depends on the memory it has. The same guide says it uses Flink’s managed memory by default, and that the default managed memory fraction of 0.4 can often be raised on TaskManagers with multi-GB process sizes.
taskmanager.memory.managed.fraction: 0.6
For state that exceeds the local disk of the TaskManagers, the guide points to the disaggregated ForStStateBackend, which keeps state in a separate storage system such as S3 or HDFS. Switching backend changes memory use and checkpoint cost together, so plan it as a migration through a savepoint and keep the fixes above in place.
RocksDB raises the limit to disk space but does not remove expired keys, so a leak continues until the disk fills.
Where Flex fits
This page is published by Factor House, which makes Flex, its Flink job management and monitoring product, so weigh this section with that in mind. Flex is one of several ways to watch a Flink job, and the comparison of tools for self-managed Apache Flink sets it beside the Flink web UI, the REST API and metrics stacks. Teams on Ververica Platform can start from the Ververica Platform tool comparison.
Diagnosis. Flex documents these views for each job in its job inspection documentation.
- The Checkpoints tab shows Avg checkpoint size and a history table with the Size of each checkpoint, so the trend across checkpoints is visible for one job. The size is the figure Flink reports, so the incremental checkpoint caveat above applies to it.
- The Configuration view on that tab shows the job’s state backend, checkpoint storage, interval, timeout and minimum pause as Flink reports them.
- The Topology tab shows the job graph, metrics for each subtask of a selected task, and the current watermark for each subtask, which is where a watermark held back by an idle input shows up.
- The Overview tab’s assigned TaskManagers summary shows heap and off-heap memory for each TaskManager running the job.
What stays in Flink. The job documentation does not describe the following, so they stay in the Apache Flink web UI, REST API or a metrics system.
- Checkpoint size broken down by operator.
- RocksDB native metrics such as the estimated key count and live data size for each state.
- State TTL settings, which live in the job’s code or SQL configuration.
Fixes. Flex documents these actions.
- The Inspect view has quick actions to Stop, Cancel, take a Savepoint and trigger a Checkpoint, which covers the savepoint before a TTL change or a backend migration.
- Submitting a job lets the operator choose a savepoint path, a restore mode and the Allow Non Restored State option, which applies when a fix removes or renames a state.
To see each job’s watermark by subtask, TaskManager memory and checkpoint history, and to take a savepoint before a TTL change, see watermarks and checkpoints for each job in Flex.
FAQ
Does Flink state TTL work with event time?
No. The Flink state documentation says only TTLs in reference to processing time are currently supported. Expiry based on event time needs event-time timers that clear the state when they fire.
Do incremental checkpoints make Flink state smaller?
No. They make each checkpoint write only what changed, which cuts checkpoint time, but the state backends documentation says the Checkpointed Data Size in the web UI then shows only the delta. Read the full size from Full Checkpoint Data Size or lastCheckpointFullSize.
How do you find which Flink operator holds the most state?
Open More details on a completed checkpoint in the web UI’s History tab, which breaks the checkpoint down by operator and by subtask, as the checkpoint monitoring documentation describes. To see which named state within that operator is growing, turn on the RocksDB native metrics with the column family as a variable and compare the estimated key count for each state.
Will switching to RocksDB fix growing state?
It removes the heap as the limit, so the job can hold much more state on disk, as the state backends documentation describes. It does not remove expired keys or close windows, so state from a leak keeps growing until disk becomes the limit. Fix the cause with TTL, timers or key design.