skip to content

How do you choose between Flink's HashMapStateBackend and EmbeddedRocksDBStateBackend?

level: middleimportance: must knowfreq 78%

answer

  1. Memory versus disk, speed versus scale
  2. One stores objects, one stores bytes
  3. Ask what bounds total state size
  4. Only one can do incremental checkpoints
  5. state.backend.type defaults to hashmap

basics

~20 s

HashMapStateBackend keeps state as Java objects on the TaskManager heap: fastest access, but state must fit in heap and only full snapshots are possible. EmbeddedRocksDBStateBackend keeps serialised state on local disk: far larger state and incremental checkpoints, at the cost of serialisation on every access.

solid answer

~50 s

It is a performance-versus-scalability call. `HashMapStateBackend` — the default when `state.backend.type` is unset — stores key/value state and window contents as objects on the JVM heap, so access involves no de/serialisation and is very fast. Its limits are hard: total state per TaskManager is bounded by heap, restoring needs enough heap to hold that share, and it only supports full snapshots, so checkpoint duration and recovery time grow with state size. `EmbeddedRocksDBStateBackend` holds state as serialised bytes in an embedded RocksDB in the TaskManager's local data directories, so state is bounded by disk rather than memory, snapshots are always asynchronous, and, of the two, it is the one that offers incremental checkpoints. The price is de/serialisation on every read and write, which is roughly an order of magnitude slower per access. Choose heap for small state and tight latency; choose RocksDB for large state and long windows.

code

yaml · 5 lines
yaml
# cluster-wide default state backend
state.backend.type: rocksdb

# keep RocksDB's native memory inside Flink's managed memory budget
state.backend.rocksdb.memory.managed: true

go deeper

for a junior

Recall the two bundled names and the one-line trade: HashMapStateBackend keeps state in memory and is fast but limited; EmbeddedRocksDBStateBackend keeps state on local disk and scales much further.

for a middle

Explain the mechanics behind the trade — heap objects versus serialised bytes, heap-bounded versus disk-bounded, full snapshots versus incremental checkpoints — and name the config key state.backend.type with its hashmap default.

for a senior

Demonstrate operating judgment: sizing RocksDB's managed memory budget, knowing incremental restore can be slower when the bottleneck is network, and migrating a live job between backends through the unified savepoint format.

for a principal

Own the platform default and the exception policy: which workloads are allowed the heap backend, what state-size threshold forces RocksDB, whether disaggregated state is worth piloting, and how backend choice interacts with checkpoint SLAs and cluster cost.

## What a state backend actually decides A state backend determines two things: how state is represented while the job runs, and how it is written out when a checkpoint or savepoint is taken. It does *not* decide what your code can express — the same `ValueState`, `MapState` and window code runs on either backend. It decides how much state you can hold and what each access costs. Flink 2.3 bundles `HashMapStateBackend`, `EmbeddedRocksDBStateBackend` and the experimental `ForStStateBackend`, added in Flink 2.0. If you configure nothing, you get `HashMapStateBackend`. ## HashMapStateBackend State lives on the Java heap of the TaskManagers as ordinary objects: key/value state and window operators hold hash tables of the values, triggers and so on. Reading `ValueState.value()` returns the stored object directly. Nothing is serialised on the access path, which makes this the low-latency choice. The limitations follow from that design: - **State size is bounded by JVM heap.** Not just at steady state — restoring a checkpoint or savepoint requires each TaskManager to have enough heap for its share of the state. - **Only full snapshots.** Incremental checkpoints are not available, so every checkpoint captures the complete state. As state grows, checkpoint duration and recovery time grow with it. - **Objects are shared, not copied.** Because values are handed back as heap objects, it is unsafe to reuse or mutate an object you have read out without writing it back. It is the right choice for jobs whose state comfortably fits in heap and that care about latency more than scale. When you use it, it is recommended to set Flink's *managed memory* to zero, so the maximum amount of memory goes to the JVM heap for user code and state. ## EmbeddedRocksDBStateBackend RocksDB is an embedded log-structured key/value store. Flink runs it embedded in the TaskManager, keeps its files in the TaskManager's local data directories, and stores every state entry as a serialised byte array. Key comparisons are byte-wise rather than through `hashCode()`/`equals()`. What that buys you: - **State bounded by disk, not memory.** You can hold very large keyed state, long windows and large per-key collections. - **Always asynchronous snapshots.** The snapshot does not block record processing. - **Incremental checkpoints.** Rather than writing a full self-contained backup, an incremental checkpoint records only what changed since the last completed one, building on previous checkpoints and self-consolidating over time through RocksDB's own compaction. This is the single biggest lever on checkpoint duration for large state, and it must be enabled deliberately: the option `execution.checkpointing.incremental` defaults to `false`. - **Safe object reuse**, precisely because everything round-trips through serialisation. What it costs: - **Every access de/serialises**, and may read from disk. Average state access is roughly an order of magnitude slower than the heap backend, so maximum throughput is lower. - **A hard 2^31-byte limit per key and per value**, because the RocksDB JNI bridge is `byte[]`-based. States that use RocksDB merge operations, such as `ListState`, can silently accumulate past that and then fail on the next retrieval. - **A memory budget to manage.** RocksDB allocates native memory outside the JVM. By default Flink sizes it from the TaskManager's *managed memory* on a per-slot basis (`state.backend.rocksdb.memory.managed`, default `true`), using a shared block cache and write-buffer manager so the total stays inside the budget. RocksDB is the recommended choice for jobs with very large state, long windows and large key/value state, and for high-availability setups. ## Recovery-time nuance Incremental checkpoints do not always restore faster. If network bandwidth is the bottleneck, restoring can take longer because more deltas must be fetched. If CPU or IOPS is the bottleneck, restoring is faster, because RocksDB's local table files are copied back rather than rebuilt from Flink's canonical key/value snapshot format (which is what savepoints and full checkpoints use). ## Configuring it The default for a cluster is set with the configuration key `state.backend.type`, whose default is `"hashmap"`; recognised shortcut names are `hashmap`, `rocksdb` and `forst`, or you can give the class name of a `StateBackendFactory`. A job can override it programmatically on the `StreamExecutionEnvironment` by setting `StateBackendOptions.STATE_BACKEND` in a `Configuration` and calling `env.configure(config)`; incremental checkpoints are switched on the same way, with `CheckpointingOptions.INCREMENTAL_CHECKPOINTS`. Flink 2.0 removed `StreamExecutionEnvironment.setStateBackend(...)`, so in Flink 2.3 a job picks its backend through options, not by handing over a backend object. ## Switching backends on a live job Savepoints come in two formats. The **canonical** format — the default when you trigger a savepoint — is unified across state backends, so you can take a canonical savepoint with one backend and restore it with another: the standard migration path when a job outgrows heap. A **native**-format savepoint (for RocksDB, its SST files) is faster to take and restore but cannot change backend, and `ForStStateBackend` cannot produce canonical savepoints at all. One more boundary: state compatibility between Flink 1.x and 2.x is not guaranteed, so a savepoint from a 1.x job is not a supported starting point for a switch on 2.3. ## The third option `ForStStateBackend`, new in Flink 2.0, is also an LSM-tree store (built on RocksDB) but designed for *disaggregated* state: its SST files live on a remote filesystem such as HDFS or S3, with the TaskManager's local disk used only as a cache. That lifts state size beyond local disk capacity and makes checkpoint and recovery lighter, at the cost of network latency on state access — which is why it pairs with Flink's asynchronous State API V2. It is still experimental and not fully production-ready, and it does not support canonical savepoints, full snapshots, changelog or file-merging checkpoints. Know it exists; do not default to it. ## How to answer in an interview Say it is a memory-versus-disk trade, then get concrete: heap objects and no serialisation versus serialised bytes on local disk; heap-bounded versus disk-bounded; full snapshots only versus incremental checkpoints; and note that a canonical savepoint lets you switch later, so the decision is reversible.

  • What is the default state backend if nothing is configured?
    `HashMapStateBackend`. The configuration key `state.backend.type` defaults to `"hashmap"`; the other recognised shortcut names are `rocksdb` for `EmbeddedRocksDBStateBackend` and `forst` for `ForStStateBackend`, and you may also give the fully qualified class name of a `StateBackendFactory`. A job can override the cluster default programmatically by setting the option in a `Configuration` and calling `env.configure(...)` on the `StreamExecutionEnvironment`.
  • Can you switch a running job from HashMapStateBackend to EmbeddedRocksDBStateBackend?
    Yes, via a savepoint in the canonical format, which is the default and is unified across backends. Take a canonical savepoint, change the backend configuration, restart from the savepoint. A native-format savepoint cannot change backend, and `ForStStateBackend` does not produce canonical savepoints. On Flink 2.3 the other boundary is major-version state compatibility: state from a 1.x job is not guaranteed to restore on 2.x at all.
  • Why does enabling RocksDB not automatically make checkpoints incremental?
    Incremental checkpointing is an opt-in feature, not RocksDB's default. You enable it with `execution.checkpointing.incremental: true` in the configuration file, or by setting `CheckpointingOptions.INCREMENTAL_CHECKPOINTS` on the job's `Configuration`; the option defaults to `false`. Once on, the web UI's Checkpointed Data Size shows only the delta for that checkpoint rather than full state size, which surprises operators who expect it to report total state.

The heap backend is a desk covered in open folders — instant to consult, but only as much as the desk holds. RocksDB is a filing cabinet beside the desk: far more capacity, but every lookup costs a walk and a page turn.

saying these in an interview costs you the question

  • Says RocksDB is faster than heap for state access
  • Claims HashMapStateBackend supports incremental checkpoints
  • Thinks RocksDB memory comes out of the JVM heap
  • Believes changing backend always requires reprocessing from scratch
  • Still names MemoryStateBackend or FsStateBackend as current options

context