Most Flink jobs that misbehave in production do not fail because the cluster is wrong. They fail because a stateful function was written with the wrong mental model: state that is never cleaned up, a list rewritten on every record, a timer per event instead of per key, or an HTTP call inside processElement that the next checkpoint cannot take back. This article is about the code you write against Flink state and the habits that keep it correct and cheap.

The architecture underneath, meaning keyed versus operator state, key groups, backends, TTL and schema evolution, is covered in Flink State Architecture and the backend choice in Flink State Backends Compared. Here we take one realistic function, a card-velocity fraud check, and use it to explain the current-key contract, the cost of each access, timers, checkpoint semantics, testing, production inspection and the newer asynchronous state API. API names were checked against the Flink 2.0 documentation; on 1.18 or earlier, use open(Configuration) where this page uses open(OpenContext), added in 1.19.

Advertisement

The contract: state is scoped to the current key

Keyed state looks like a field on your function, but it is not one. When a record arrives on a KeyedStream, Flink sets the operator's current key from your key selector before calling processElement. Every read or write to a ValueState, ListState or MapState handle then resolves to that key's slot only. A handle is a view, not a container.

Three consequences follow. First, you cannot iterate across keys from inside a keyed function; if you need a global view, that is a different operator or a batch read of a savepoint. Second, timers are also keyed: when onTimer fires, the current key is restored to the key that registered it, so state you touch there belongs to that key. Third, the key selector must be deterministic. If it depends on anything other than the record, such as a random salt or the wall clock, records for one logical entity scatter across key groups and the state silently splits.

A worked example: velocity check with buckets and timers

The requirement: alert when one card makes more than five transactions within three minutes of event time. A naive version keeps a ListState<Long> of timestamps, appends each one, reads the whole list back, filters out old entries and rewrites it. That works in a unit test and falls over in production for heavy keys, because every record pays to deserialise and reserialise the entire list.

The version below keys by card and keeps a MapState<Long, Integer> from minute-start to count. Each record touches at most three map entries. Cleanup is done with one event-time timer per key per minute, which leans on a useful guarantee: Flink deduplicates timers by key and timestamp, so registering the same minute boundary for the sixth time costs nothing extra.

public class VelocityCheck
        extends KeyedProcessFunction<String, Txn, Alert> {

    private static final long MINUTE = 60_000L;
    private static final int LIMIT = 5;          // more than 5 txns in 3 minutes
    private transient MapState<Long, Integer> perMinute;   // minute start -> count

    @Override
    public void open(OpenContext ctx) {
        perMinute = getRuntimeContext().getMapState(
            new MapStateDescriptor<>("perMinute", Types.LONG, Types.INT));
    }

    @Override
    public void processElement(Txn t, Context ctx, Collector<Alert> out) throws Exception {
        long minute = t.eventTime - (t.eventTime % MINUTE);
        Integer n = perMinute.get(minute);                 // touches ONE entry
        perMinute.put(minute, n == null ? 1 : n + 1);

        int total = 0;
        for (long m = minute - 2 * MINUTE; m <= minute; m += MINUTE) {
            Integer v = perMinute.get(m);
            total += v == null ? 0 : v;
        }
        if (total > LIMIT) out.collect(new Alert(ctx.getCurrentKey(), minute, total));

        // One timer per key per minute: registering the same timestamp twice is a no-op.
        ctx.timerService().registerEventTimeTimer(minute + 3 * MINUTE);
    }

    @Override
    public void onTimer(long ts, OnTimerContext ctx, Collector<Alert> out) throws Exception {
        perMinute.remove(ts - 3 * MINUTE);                 // that bucket has left the window
        if (perMinute.isEmpty()) perMinute.clear();        // key fully gone from state
    }
}

Walk through one key. Transactions for card_42 at 12:00:10, 12:00:40 and 12:01:05 produce map entries {12:00 to 2, 12:01 to 1} and two timers, 12:03 and 12:04. When the watermark passes 12:03, onTimer removes the 12:00 bucket. When it passes 12:04, the 12:01 bucket goes and the key's map is empty, so nothing about this card remains in state. That last point matters more than any tuning: state size is the number of live keys times bytes per key, and a function that never removes keys grows until the job is restarted from scratch.

Advertisement

What each state access actually costs

On the heap backend, state is Java objects in a hash table, so access is a lookup and a reference. On RocksDB, and on the ForSt backend introduced with Flink 2.0, state lives in serialised form, and every read deserialises and every write serialises. That one fact explains most performance surprises.

OperationHeap backendRocksDB / ForSt backend
ValueState.value()hash lookup, returns the live objectpoint read plus full deserialisation of the value
ValueState.update(v)store referencefull serialisation of v, write to memtable
ListState.add(x)append to Java listmerge operation; serialises only x, no read
ListState.get()returns the listreads and deserialises every element
MapState.get(k)nested hash lookupone point read; each map entry is its own key
MapState iterationiterate Java mapprefix scan over this key's entries

Two rules fall out. Never store a growing collection inside a ValueState; on RocksDB every access is O(size of the collection). Prefer MapState when you need random access and ListState.add when you only append and read rarely. And on the heap backend, never mutate an object you got from value() without calling update(): it appears to work on heap because you are holding the live object, then silently loses writes on RocksDB, where you were holding a copy. Write code that is correct on the serialised backend and it will be correct on both.

Timers are state too

Timers are stored per key in the state backend and are checkpointed with it. That has costs and consequences that are easy to miss.

  • Count. A timer per event is a timer per event in state. Round timestamps, as the example does, so timers coalesce per key.
  • Event time fires on watermarks. If a source goes idle and the watermark stops, cleanup timers never fire and state grows. Configure idleness on the watermark strategy for partitions that can go quiet.
  • Restore storms. After a long outage, recovering and replaying a backlog can advance the watermark by hours at once, firing every overdue timer in a burst. Processing-time timers whose time has passed fire right after restore. Make onTimer cheap and idempotent.

What a checkpoint restores, and what it cannot undo

One keyed operator subtask: where state, timers and checkpoints meetInput recordkey = card_42processElement()current key set to card_42keyByMapState: minute to countonly card_42's entries visibleTimer servicecard_42 @ 12:03:00 (deduplicated)get / putregisteronTimer()same key restoredwatermark passes timestampState backendheap objects or ForSt / RocksDBstored inCheckpoint barrier narrives from upstreamsnapshot state + timersDurable checkpoint storagerestored together on failureRecords, state and timers are consistent with each other at a barrier; external side effects are not.
Barriers flow with the records. When an operator has seen barrier n on its inputs, it snapshots its state and timers, so state reflects exactly the records before the barrier.

On failure, Flink rewinds every operator to the last completed checkpoint and rewinds the sources to the matching offsets. Records after the barrier are processed again. For state this is exactly-once: the map counts above will be right after recovery. For anything your function does outside Flink state, it is at-least-once.

So an alert that was emitted and sent to a non-transactional sink after checkpoint 41 will be emitted again after a restore from 41. A call to an external service in processElement is repeated. A counter you increment in a static field or a database is double-counted. The fixes are the usual ones: make side effects idempotent with a deterministic id, such as card plus minute here, or send them through a transactional sink as described in exactly-once stream processing.

One more trap: state is matched to operators by uid on restore. If you do not set .uid("velocity-check") on the operator, Flink generates one from the job graph, and an innocent refactor upstream can change it. Savepoint restores then fail, or with the allow-non-restored-state option set, start the operator with empty state. Set uids on every stateful operator from day one; the savepoints article covers the rest of the upgrade workflow.

Testing stateful functions with the harness

Stateful logic deserves unit tests that exercise timers and recovery, not just one input and one output. Flink ships operator test harnesses in flink-test-utils (and the streaming test-jar), with processElement, processWatermark to drive event time, setProcessingTime for processing-time timers, and extractOutputValues to assert on output. The documentation is explicit that these harness classes are not part of the public API and can change between versions, so pin the dependency to your Flink version and expect small edits on upgrade.

@Test
void alertsOnSixthTxnAndForgetsOldMinutes() throws Exception {
    var harness = new KeyedOneInputStreamOperatorTestHarness<>(
        new KeyedProcessOperator<>(new VelocityCheck()),
        (Txn t) -> t.card, Types.STRING);
    harness.open();

    for (int i = 0; i < 6; i++) harness.processElement(new Txn("c1", 1_000L + i), 1_000L + i);
    assertEquals(1, harness.extractOutputValues().size());      // sixth one fires

    var snap = harness.snapshot(1L, 2_000L);                    // checkpoint mid-stream
    harness.processWatermark(4 * 60_000L);                      // fire the cleanup timer
    assertEquals(0, harness.numKeyedStateEntries());

    var restored = new KeyedOneInputStreamOperatorTestHarness<>(
        new KeyedProcessOperator<>(new VelocityCheck()),
        (Txn t) -> t.card, Types.STRING);
    restored.initializeState(snap);                             // recovery path
    restored.open();
    restored.processElement(new Txn("c1", 1_010L), 1_010L);
    assertEquals(1, restored.extractOutputValues().size());     // count survived restore
}

The valuable part is the snapshot and restore. Taking a snapshot, building a fresh harness and initialising it from that snapshot runs the same serialisation path as a real recovery, so it catches state that does not serialise and logic that only works on live heap objects. Add a test that drives the watermark far forward and asserts that the keyed state entry count returns to zero; that single assertion prevents the most common production incident with stateful jobs, unbounded growth.

Looking inside production state

When a job's checkpoints grow and you need to know why, metrics give you total size but not which keys. The State Processor API reads a savepoint or retained checkpoint as a bounded input, using the same descriptors your job uses. The reader function registers descriptors in open and is called once per key.

// Batch job: dump per-card bucket counts from a savepoint to find the hot keys.
SavepointReader sp = SavepointReader.read(env, "s3://ckpt/savepoints/sp-123", new EmbeddedRocksDBStateBackend());

DataStream<String> rows = sp.readKeyedState(
    OperatorIdentifier.forUid("velocity-check"),
    new KeyedStateReaderFunction<String, String>() {
        MapState<Long, Integer> perMinute;
        @Override public void open(OpenContext ctx) {
            perMinute = getRuntimeContext().getMapState(
                new MapStateDescriptor<>("perMinute", Types.LONG, Types.INT));   // same name and types
        }
        @Override public void readKey(String card, Context ctx, Collector<String> out) throws Exception {
            int entries = 0;
            for (Long ignored : perMinute.keys()) entries++;
            out.collect(card + "," + entries + "," + ctx.registeredEventTimeTimers().size());
        }
    });

Sort the output by entry count and you will usually find the answer quickly: a handful of keys, such as a test card, a null-key default or a bot user, holding most of the state, or millions of keys whose timers never fired. The same API can write a new savepoint, through SavepointWriter and OperatorTransformation.bootstrapWith, which is how you backfill state for a new operator from historical data or repair state without replaying months of input. Run it on a copy and validate in staging first.

The asynchronous State V2 API

Flink 2.0 added a second keyed state API in org.apache.flink.api.common.state.v2, enabled per stream with keyBy(...).enableAsyncState(). Accessors return a StateFuture that you compose with thenAccept, thenApply, thenCompose and thenCombine, using methods such as asyncGet, asyncPut, asyncValue and asyncUpdate.

// Flink 2.x State V2: keyBy(...).enableAsyncState(), descriptors from ...state.v2
perMinute.asyncGet(minute)
    .thenCompose(n -> perMinute.asyncPut(minute, n == null ? 1 : n + 1))
    .thenAccept(ignored -> ctx.timerService().registerEventTimeTimer(minute + 3 * MINUTE));

The point is disaggregated state. When state lives on remote storage, a synchronous read stalls the task thread for a network round trip; the async API lets the runtime keep many keys' accesses in flight while preserving per-key ordering. The documentation recommends the ForSt backend for it and notes that other backends execute these accesses synchronously, so on heap or RocksDB the rewrite buys you nothing yet. Treat it as a newer API that is still settling: adopt it for jobs whose state does not fit on local disk, and keep the synchronous API elsewhere.

Failure modes

  • Unbounded growth. Keys are created and never removed. Fix with cleanup timers or TTL, and a harness test that asserts state returns to zero.
  • Collection in a ValueState. Latency rises with key size on RocksDB. Move to MapState or ListState.
  • Mutating without update. Correct on heap, silently wrong on RocksDB. Always write back.
  • Missing uids. Savepoint restore fails or drops state after a refactor.
  • Side effects in processElement. Duplicated after every recovery. Make them idempotent or transactional.
  • Hot keys. One key bottlenecks one subtask. Find it with a savepoint read, then pre-aggregate or split it.

Trade-offs

Minute buckets trade exactness for bounded cost; per-event timestamps give exact windows at the price of more entries. Cleanup timers are precise and testable but add timer state; TTL is simpler to declare but cleans up lazily. A built-in window operator, covered in Flink Windowing, saves you this code when the logic fits its model; a hand-written KeyedProcessFunction gives you control when it does not, and makes you responsible for cleanup.

What to do next

  1. List every stateful operator in your job and confirm each has an explicit uid.
  2. For each one, write down how a key's state is eventually removed: timer, TTL or never. Fix every never.
  3. Search for growing collections stored in ValueState and move them to MapState or ListState.
  4. Search processElement and onTimer for external calls and make each idempotent or route it through a transactional sink.
  5. Add a harness test per stateful function that snapshots, restores and drives the watermark forward until keyed state returns to zero.
  6. Run a State Processor API read on a recent savepoint and look at the top 20 keys by entry count.
  7. Configure idleness on sources with partitions that can go quiet, so cleanup timers keep firing.
Key takeaway: Flink state is a per-key view backed by a serialising store, so write every function as if it ran on RocksDB: choose the primitive by access pattern, write back every change, coalesce timers and make every key eventually disappear. Checkpoints make state exactly-once but not your side effects, so make those idempotent. Then prove it with harness tests that snapshot, restore and drain state, and look inside production state with the State Processor API before guessing at why checkpoints grow.