Every stateful Flink job makes one decision that shapes its latency, its memory bill, its checkpoint duration and how long it takes to come back after a failure: the state backend. Flink 2.x offers three. The hashmap backend keeps state as Java objects on the heap. The RocksDB backend keeps serialized state in an embedded RocksDB instance on local disk. The ForSt backend, still marked experimental, keeps its files on remote storage such as S3 and uses local disk as a cache.
This article compares them operationally. It covers what each backend costs per access, what each does at checkpoint time, how fast each recovers and rescales, how to configure memory, and how to switch backends safely when a job outgrows its choice. The primitives themselves, such as ValueState, MapState, key groups and TTL, are covered in Flink state architecture. Here the question is which backend to run and how to run it well.
Backend versus checkpoint storage
Two settings are easy to confuse. The state backend, set with state.backend.type, decides where working state lives while the job runs. Checkpoint storage, set with execution.checkpointing.storage and execution.checkpointing.dir, decides where snapshots of that state are written. For any production job the storage should be a durable filesystem such as S3 or HDFS, whatever the backend. If no backend is configured, Flink uses hashmap.
# config.yaml (Flink 2.x)
state.backend.type: rocksdb # hashmap | rocksdb | forst
execution.checkpointing.storage: filesystem
execution.checkpointing.dir: s3://my-bucket/flink/checkpoints/orders-enricher
execution.checkpointing.incremental: true # RocksDB: upload only new SST files
execution.checkpointing.interval: 60 s
# RocksDB memory: keep it inside Flink's managed memory budget
taskmanager.memory.managed.fraction: 0.4
state.backend.rocksdb.memory.managed: true
state.backend.rocksdb.memory.write-buffer-ratio: 0.5
state.backend.rocksdb.memory.high-prio-pool-ratio: 0.1
state.backend.rocksdb.timer-service.factory: rocksdb
The three backends at a glance
| hashmap | rocksdb | forst | |
|---|---|---|---|
| Where state lives | JVM heap objects | local RocksDB files, memory cache | remote SST files, local cache |
| Per-access cost | object lookup | serialize and deserialize every access | as RocksDB, plus remote reads on cache misses |
| State size limit | heap per TaskManager | local disk | remote storage |
| Checkpoints | full snapshot | full or incremental, always asynchronous | always asynchronous and incremental |
| State API | synchronous | synchronous | designed for the asynchronous State API V2 |
| Maturity | production | production | experimental in the current docs |
The hashmap backend is fastest per access because a get is a hash lookup returning a live object. The cost is that all state must fit on the heap, garbage collection pressure grows with state, and every checkpoint writes all of it. RocksDB moves state off-heap and onto disk, so it scales to state far larger than memory, at the price of serialization on every read and write. ForSt goes further and puts the files on remote storage, so a TaskManager's local disk no longer limits state size and recovery does not need to download everything first. Remote reads are slow, which is why the documentation pairs ForSt with asynchronous state access: the operator issues many state requests at once instead of blocking on each.
The object-reuse trap when switching
Code written against the hashmap backend can depend, without anyone noticing, on getting the stored object back. Mutating it changes the state. On RocksDB, value() returns a freshly deserialized copy, so the mutation is lost unless you call update(). The documentation also warns the other way round: since hashmap returns the stored objects, it is unsafe to reuse or emit them. Write code that is correct on both.
// Works on hashmap, silently broken on rocksdb.
public class CountItems extends KeyedProcessFunction<String, Order, Summary> {
private transient ValueState<Summary> summary;
@Override
public void processElement(Order o, Context ctx, Collector<Summary> out) throws Exception {
Summary s = summary.value();
if (s == null) { s = new Summary(); summary.update(s); }
s.count += o.items; // heap: mutates the stored object, so the change "sticks"
// rocksdb: mutates a deserialized copy, the change is lost
// Correct on every backend: write back explicitly after mutating.
summary.update(s);
out.collect(s.copy()); // never emit the state object itself
}
}Test every job against RocksDB in CI even if production runs hashmap. This one rule removes the most common surprise of a backend switch.
Checkpoints: full, incremental and changelog
With hashmap, every checkpoint is a full copy of the state. At 10 GB per TaskManager that is 10 GB uploaded every interval. Checkpoint duration grows linearly with state, and once it approaches the interval, checkpoints back up and time out.
RocksDB with execution.checkpointing.incremental: true uploads only the SST files created since the last checkpoint and refers to the earlier ones. Upload volume tracks the rate of change rather than total size, which is the main reason large-state jobs choose RocksDB. Incremental checkpoints are not uniformly small: a background compaction rewrites files, and the next checkpoint uploads the rewritten files, so expect occasional large checkpoints. Checkpoint storage holds a chain of shared files, so delete it only through Flink's retention, never by age.
The changelog option, enabled with state.changelog.enabled: true plus state.changelog.storage and state.changelog.dstl.dfs.base-path, continuously uploads state changes so that each checkpoint has little left to do. It shortens and stabilizes checkpoint duration, which matters for transactional sinks that commit on checkpoint completion, as explained in exactly-once semantics. The documented costs are more files on the remote filesystem, more IO bandwidth, more CPU for serialization and more TaskManager memory for buffering. It also requires execution.checkpointing.max-concurrent-checkpoints: 1.
Recovery and rescaling
Recovery time is part of your availability, so measure it per backend. A hashmap job restores by reading the full snapshot and rebuilding every object on the heap, so restore time grows with state size and deserialization speed. A RocksDB job downloads its SST files and opens the database; with incremental checkpoints this is still all the files, but no per-record rebuild is needed. Rescaling a RocksDB job must also split or merge key-group ranges, so plan for it taking longer than a same-parallelism restart. Task-local recovery, which keeps a local copy of the latest snapshot, lets a failed task that is rescheduled onto the same TaskManager skip the download. ForSt's design target is that a restored task can start working against remote files without first downloading them, which is attractive for very large state but should be measured in your environment, given its experimental status.
Memory configuration
With hashmap, size the TaskManager heap for the state plus headroom for garbage collection, typically keeping live state under half the heap. Use a garbage collector suited to large heaps and watch pause times, because a long pause stalls processing and can make checkpoints time out.
With RocksDB, state memory comes from Flink's managed memory, set by taskmanager.memory.managed.fraction. With state.backend.rocksdb.memory.managed: true, RocksDB's write buffers and block cache share that budget per slot. The write-buffer ratio (default 0.5) splits it between writes and the block cache, and the high-priority pool ratio (default 0.1) reserves cache for index and filter blocks. If you disable managed memory, RocksDB allocates outside Flink's accounting, and containers start being killed for exceeding their memory limit.
Timers deserve their own line. RocksDB stores timers in RocksDB by default. Setting state.backend.rocksdb.timer-service.factory: heap is faster when there are few timers, but the documentation notes that RocksDB with heap timers does not snapshot timer state asynchronously, so a job with many timers then pays for them synchronously at every checkpoint.
A worked sizing example
An order-enrichment job keeps per-customer state for 30 days: 200 million keys at about 300 bytes serialized, roughly 60 GB, on 20 TaskManagers, so 3 GB each. About 2 percent of keys change per minute, and checkpoints run every 60 seconds.
On hashmap, 3 GB of serialized state is often 2 to 4 times larger as Java objects, so 6 to 12 GB of live heap per TaskManager, and a full snapshot of 3 GB per TaskManager, 60 GB in total, every minute. That is feasible on paper and fragile in practice: garbage collection pauses grow, and 60 GB per minute of checkpoint upload is roughly 1 GB per second sustained.
On RocksDB with incremental checkpoints, the heap is small, state sits on local SSD, and each checkpoint uploads the changed data: about 2 percent of 60 GB, about 1.2 GB per minute plus compaction output. The price is per-access serialization, perhaps tens of microseconds per record, which this job's throughput tolerates. RocksDB is the clear choice. Hashmap would win for a job with 200 MB of state and a tight latency budget.
Switching backends
Since Flink 1.13, savepoints use a unified binary format, so a savepoint taken on one backend can be restored on another. Canonical-format savepoints are the portable kind; native-format savepoints, like checkpoints, are backend-specific. The procedure is short, but the operator uids must be stable, or state cannot be matched to operators. The savepoints article covers uids and compatibility.
# 1. Stop the job with a savepoint (canonical format is portable across backends)
flink stop --savepointPath s3://my-bucket/flink/savepoints/orders-enricher <jobId>
# 2. Change the backend in config.yaml (for example hashmap -> rocksdb), keep operator uids unchanged
# 3. Restart from the savepoint
flink run -s s3://my-bucket/flink/savepoints/orders-enricher/savepoint-xxxx orders-enricher.jarRehearse the switch on a copy of production state first and compare record counts and key metrics between old and new jobs before cutting over.
Benchmarking your own job
- Replay a fixed slice of production input at the target rate against each candidate backend, with identical parallelism and hardware.
- Record throughput, p99 end-to-end latency, checkpoint duration and size, and alignment time from the Flink web UI or metrics.
- Kill a TaskManager mid-run and measure time to resume processing; repeat with a parallelism change to measure rescale time.
- Enable RocksDB native metrics, which are off by default, for memtable, cache and compaction behaviour.
- Run long enough for compaction and garbage collection to reach steady state; ten-minute benchmarks flatter both.
Failure modes
- Heap exhaustion on hashmap: state growth from a missing TTL ends in long garbage collection pauses and then OutOfMemoryError.
- Value size limit on RocksDB: keys and values are limited to 2^31 bytes, and the documentation warns that ListState built by merge operations can silently grow past it and fail on the next read. Bound list sizes or use MapState.
- Local disk full: RocksDB working files plus compaction can exceed the disk; monitor it like any USE resource.
- Unbounded memory: disabling managed memory leads to container kills that look like random crashes.
- Checkpoint storage cleanup: deleting old incremental checkpoint files by age breaks newer checkpoints that still reference them.
- Production use of an experimental backend: ForSt is documented as experimental, so treat it as an evaluation target, not a default.
Trade-offs
Choose hashmap when state per TaskManager is comfortably within the heap and latency matters most; it is the simplest and fastest. Choose RocksDB with incremental checkpoints when state is large, grows over time or must checkpoint quickly; it is the production default for large state. Add the changelog option when checkpoint duration itself is the bottleneck and you can pay for the extra IO. Evaluate ForSt when state is too large for local disks or recovery downloads dominate your downtime, and be ready to adopt the asynchronous State API to get its benefit. For the wider runtime picture, see Apache Flink architecture in depth.
What to do next
- Record each job's state size per TaskManager, checkpoint duration and size, and restore time.
- Run your test suite against RocksDB and fix any code that mutates state without calling update().
- For RocksDB jobs, enable incremental checkpoints and confirm managed memory is on.
- Set TTL on every state that should expire, and bound any ListState.
- Benchmark the alternative backend on replayed input, including a kill-and-restore test.
- If switching, take a canonical savepoint, change the backend, restore and verify outputs before cutting over.