Everything interesting Flink does, from a windowed count to a fraud rule that remembers a card's last fifty transactions, depends on state: data the job keeps between events. Flink promises that state is local, fast, consistent after failures and movable when you rescale. Each property comes from specific machinery with limits you can hit in production.

This article is about the state itself: which primitives exist, how keyed state is partitioned into key groups, how the heap and RocksDB backends store it, how large it gets, how it expires, and how its serialized format survives code changes. The checkpoint protocol that makes state durable is covered in Apache Flink architecture in depth, and portable snapshots for upgrades in Flink savepoints; we refer to both rather than repeat them. Code targets Flink 1.19 and later, where open(OpenContext) and java.time.Duration replace the older signatures.

Advertisement

What state is, and who owns it

Flink has two families of state. Keyed state exists only after a keyBy. Every record carries a key, the runtime routes all records with the same key to the same parallel subtask, and any state you declare is automatically scoped to the current key. When processElement runs for user 42, session.value() returns user 42's value; you never write a lookup by key. Operator state is scoped to a parallel instance instead of a key. Sources use it to remember offsets; a sink might use it to hold a buffer of uncommitted writes.

A subtask can touch only the state for keys routed to it: no global map, no cross-key query. That restriction makes state local, lock-free and partitionable. A computation that needs two keys at once must re-key or join two keyed streams on a common key, as stateful stream joins shows.

Where keyed state lives in a Flink jobSourceKafka partitionskeyBy(userId)murmur(hash) % maxPSubtask 0key groups 0-63Subtask 1key groups 64-127Sinkfeatures topicHeap backendJava objects, JVM heapRocksDB backendserialized, local SSDstate reads/writesFull snapshotcopy-on-write heapIncrementalnew SST files onlyCheckpoint storageS3 / HDFS / GCSRescaling moves whole key groups between subtasks; the number of key groups (max parallelism) never changes.
Records are hashed to key groups, key groups are assigned to subtasks, and each subtask stores its key groups in a backend that snapshots to durable storage.

Keyed state primitives and choosing among them

Flink offers five keyed primitives, and the choice matters more on RocksDB than on the heap, because on RocksDB every access is a serialization round trip.

PrimitiveShapeUse it forCost on RocksDB
ValueState<T>one value per keya counter, a small POJO, a last-seen timestampone get or put of the whole value
ListState<T>append-only list per keybuffered events awaiting a triggeradd() is a merge with no read; get() deserializes the whole list
MapState<K,V>map per keyper-key sub-entities: counts by categoryeach entry is its own RocksDB key; point get and put are cheap
ReducingState<T>one value, folded on addrunning sum or max of the same typeread, reduce, write on each add
AggregatingState<IN,OUT>accumulator, different out typerunning average, sketchesread, add, write on each add

The most common mistake is storing a collection inside ValueState: a ValueState<List<Event>> or ValueState<HashMap<...>>. On the heap it looks harmless. On RocksDB, every update deserializes and re-serializes the entire collection, so the cost of one event grows with the size of the key's history. A user with 10,000 buffered events pays for 10,000 events on every click. Use ListState when you only append and occasionally drain, and MapState when you update individual entries.

Timers are keyed state too, stored per key and timestamp and checkpointed with the rest.

Advertisement

Key groups: how keyed state is partitioned and rescaled

Flink does not assign keys directly to subtasks, because that would make rescaling a key-by-key reshuffle. It inserts an indirection. Each key is hashed into one of a fixed number of key groups, and contiguous ranges of key groups are assigned to subtasks. The number of key groups equals the job's max parallelism. With max parallelism 128 and parallelism 2, subtask 0 owns groups 0 to 63 and subtask 1 owns 64 to 127. Rescale to 4 and each subtask receives a range of 32 groups, read directly from the checkpoint files that hold them.

Max parallelism is therefore the one state setting you cannot change casually: it fixes the hash layout of every checkpoint (changing it means rewriting state with the State Processor API), and it caps parallelism. Left unset, Flink derives it from the initial parallelism (roughly 1.5 times, rounded up to a power of two, minimum 128), which is how a job started at parallelism 8 discovers it can never pass 128. Set it explicitly; 720 divides evenly by many parallelisms.

Key skew passes straight through this design: a key with 20 percent of the traffic puts at least 20 percent of the work on one subtask, whatever the parallelism. The fix is salting the hot key into sub-keys and merging downstream.

Operator state and broadcast state

Operator state is declared through CheckpointedFunction, and because there is no key to hash, redistribution on rescale is explicit. Even-split concatenates the lists from all old instances and deals them out evenly, which suits Kafka partition offsets. Union gives every new instance the full list; with large lists and high parallelism, every instance downloads everything and restores blow up.

Broadcast state is a map replicated to every instance: a low-volume control stream such as fraud rules is broadcast, and a high-volume stream is processed against it. Only the broadcast side may write, and every instance must apply identical updates, because Flink assumes the checkpointed copies are the same.

The backends: heap, RocksDB and the disaggregated model

The heap backend (HashMapStateBackend) keeps state as Java objects in the TaskManager's heap. Access involves no serialization, so it is the fastest option for small state. Snapshots use copy-on-write so processing continues, but every checkpoint writes the entire state, and object overhead plus garbage collection limit how large it can grow.

The RocksDB backend (EmbeddedRocksDBStateBackend) stores serialized key-value pairs in an embedded LSM tree on local disk, with block cache and memtables budgeted from Flink's managed memory. State can exceed memory, but every access pays for serialization and possibly a disk read. Its key feature is incremental checkpoints: RocksDB's sorted files are immutable, so a checkpoint uploads only new files and references the rest.

Flink 2.0 adds ForSt (FLIP-427 and the VLDB 2025 paper on disaggregated state): a RocksDB-derived store whose primary copy lives on a distributed file system such as S3, with local disk as a cache, paired with an asynchronous state API that keeps many state requests in flight. The paper reports much shorter checkpoints and faster recovery and rescaling. It is the newest option and its operational knobs are still evolving.

# Flink 1.x key names; several checkpoint keys moved under execution.checkpointing.*
# in later releases, so check the configuration page for your version.
state.backend.type: rocksdb
state.backend.incremental: true
state.backend.rocksdb.memory.managed: true
taskmanager.memory.managed.fraction: 0.4
pipeline.max-parallelism: 720
pipeline.generic-types: false      # fail at job build time instead of silently using Kryo

A worked example: per-user session features

Suppose a clickstream of 40 million active users feeds a feature store. For each user we want sessions, closed after 30 minutes of inactivity, plus a rolling count of clicks by product category that forgets categories untouched for a week. The function below uses a ValueState for the open session, a MapState with a TTL for the category counts, and one event-time timer per user that is moved forward on every click.

public class SessionFeatures
        extends KeyedProcessFunction<String, Click, SessionSummary> {

    private static final long GAP_MS = 30 * 60 * 1000L;   // 30-minute inactivity gap

    private transient ValueState<SessionAgg> session;       // one small POJO per user
    private transient MapState<String, Long> categoryCounts; // category -> clicks, 7-day TTL

    @Override
    public void open(OpenContext ctx) {                       // Flink 1.19+ signature
        session = getRuntimeContext().getState(
            new ValueStateDescriptor<>("session", SessionAgg.class));

        StateTtlConfig ttl = StateTtlConfig.newBuilder(Duration.ofDays(7))
            .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
            .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
            .cleanupInRocksdbCompactFilter(1000)
            .build();
        MapStateDescriptor<String, Long> cc =
            new MapStateDescriptor<>("category-counts", String.class, Long.class);
        cc.enableTimeToLive(ttl);
        categoryCounts = getRuntimeContext().getMapState(cc);
    }

    @Override
    public void processElement(Click c, Context ctx, Collector<SessionSummary> out)
            throws Exception {
        SessionAgg s = session.value();
        if (s == null) {
            s = new SessionAgg(c.ts);
        } else {
            ctx.timerService().deleteEventTimeTimer(s.lastTs + GAP_MS);
        }
        s.clicks += 1;
        s.lastTs = Math.max(s.lastTs, c.ts);
        session.update(s);
        ctx.timerService().registerEventTimeTimer(s.lastTs + GAP_MS);

        Long n = categoryCounts.get(c.category);   // one point lookup, not a full-map read
        categoryCounts.put(c.category, n == null ? 1L : n + 1);
    }

    @Override
    public void onTimer(long ts, OnTimerContext ctx, Collector<SessionSummary> out)
            throws Exception {
        SessionAgg s = session.value();
        if (s != null && ts == s.lastTs + GAP_MS) {
            out.collect(new SessionSummary(ctx.getCurrentKey(), s.startTs, s.lastTs, s.clicks));
            session.clear();                       // the session ends; its state must too
        }
    }
}

// wiring: a stable uid is what ties this state to the operator on restore
clicks.keyBy(Click::userId)
      .process(new SessionFeatures())
      .uid("session-features")
      .name("session-features");

Three details carry the design: the old timer is deleted before a new one is registered, so each user has one pending timer; session state is cleared when the session closes; and a click touches one map entry, not the whole map. The uid is the name under which this operator's state is stored in checkpoints, and changing it orphans the state.

Sizing the state: the arithmetic

Estimate state before choosing a backend, and measure it afterwards. For this job, assume a serialized session of about 60 bytes plus a 16-byte key, and around 20 active categories per user, at roughly 50 bytes per MapState entry once the key, user key prefix, value and RocksDB overhead are counted.

  • Category maps: 40,000,000 users x 20 entries x 50 bytes = 40 GB.
  • Open sessions: at most 40,000,000 x about 100 bytes = 4 GB, far less in practice because most sessions are closed.
  • Total around 44 GB. At parallelism 32 that is about 1.4 GB per subtask on RocksDB, which is comfortable.
  • On the heap, object overhead of three to five times would put 4 to 7 GB of live objects in each slot, and every checkpoint would write all 44 GB.

That decides it: RocksDB with incremental checkpoints, each uploading roughly what changed plus compaction output. Without the seven-day TTL, the map would grow with every category a user ever touched, for the life of the job.

State TTL and cleanup

TTL attaches a last-modified timestamp to each value or map entry. The update type decides what refreshes it, the visibility whether an expired but not yet removed value may be returned, and the cleanup strategy when expired data actually leaves storage, which is what controls size.

By default expired values disappear only when read, so a value never read again stays forever. Pick a cleanup strategy per backend:cleanupFullSnapshot() drops expired entries from full snapshots; cleanupIncrementally(n, false) lets the heap backend check a few entries on each state access; and cleanupInRocksdbCompactFilter(queryTimeAfterNumEntries) drops expired entries during RocksDB compaction. Two further constraints: TTL is based on processing time, not event time, so replaying a month of history expires state by wall clock; and restoring state saved without TTL through a TTL-enabled descriptor, or the reverse, has failed with a StateMigrationException (FLINK-32955 tracks lifting this), so check your version before toggling TTL on live state.

Serializers and schema evolution

Every value in RocksDB and every value in a checkpoint is bytes produced by a serializer that Flink chose from the type. Flink's own serializers cover primitives, tuples, rows and POJOs. Avro types use Avro's serializer. Anything else falls back to Kryo, a generic serializer that works until you change the class.

Schema evolution means restoring state written by an older class version. POJOs support adding and removing fields; Avro follows Avro's compatibility rules; Kryo supports essentially nothing, so renaming a field can make a savepoint unreadable. Make state types valid POJOs or Avro records, set pipeline.generic-types: false so any Kryo fallback fails at build time, never change a key type, and restore the production savepoint into each new build before rollout.

Failure modes

  • Unbounded growth. Keys that never recur, such as session IDs or request IDs, with no TTL and no clear(). State size climbs linearly for months until checkpoints time out. Chart state size per operator from day one.
  • Collections in ValueState. Throughput degrades as keys age, because each update re-serializes a growing blob. The symptom is rising CPU with flat input.
  • Timer leaks. Timers registered per event and never deleted. Timer state grows even though value state looks bounded.
  • RocksDB memory overrun. Native memory outside Flink's budget, often from many state descriptors (each a column family with its own memtables), gets containers killed. Keep managed memory on.
  • Checkpoint size cliffs. A large compaction makes one incremental checkpoint upload gigabytes. Watch size and duration, not just success.
  • Missing uid. Generated operator IDs change with unrelated graph edits, and the state becomes unreachable on restore.

Trade-offs

ChoiceGainsCosts
Heap backendlowest latency, no serializationstate bounded by heap, GC pauses, full-size checkpoints
RocksDB backendstate larger than memory, incremental checkpointsserialization on every access, native memory tuning
ForSt (Flink 2.x)small local footprint, fast recovery and rescaleremote-read latency, needs async state API, newer
Large max parallelismroom to scale outfixed forever; slight per-group overhead
Short TTLbounded statecorrectness depends on processing-time expiry

For end-to-end guarantees across sources and sinks, state consistency is only one half; the other is transactional or idempotent output, covered in exactly-once processing.

What to do next

  1. List every state descriptor in your job with its primitive, key, value type and expected entries per key. Replace any collection held in ValueState.
  2. Set max parallelism explicitly on every job, choosing a value that leaves room for five to ten times growth.
  3. Give every stateful operator a stable uid.
  4. Estimate state size with the arithmetic above and pick heap or RocksDB from the result, not habit.
  5. Add a TTL and a backend-appropriate cleanup strategy to every state whose keys can stop recurring, and clear state explicitly when an entity's lifecycle ends.
  6. Set pipeline.generic-types to false and move state types to POJOs or Avro.
  7. Dashboard checkpoint size, duration and state size per operator; alert on steady growth.
  8. In CI, restore the latest production savepoint into each new build before deploying.
Key takeaway: Flink state is local, key-scoped and partitioned into a fixed number of key groups, which is what lets it be fast, consistent and rescalable. Choose the primitive that matches the access pattern, never a collection inside ValueState. Fix max parallelism and uids when the job is born, because both are baked into every checkpoint. Size the state before choosing heap or RocksDB, bound it with TTL, explicit clears and deleted timers, keep its types evolvable, and restore a real savepoint before every release.