In Flink, why prefer MapState over a ValueState holding a HashMap?
answer
- Ask what the backend actually stores
- One shape is a single opaque blob
- Serialised bytes, not Java objects
- Cost per record scales with map size
- Per-entry TTL only works on one of them
basics
~20 sMapState stores each entry separately in the state backend, so one get or put touches one entry. A ValueState holding a HashMap is one opaque blob that must be fully deserialised and re-serialised on every access.
solid answer
~40 sThe difference only really bites on `EmbeddedRocksDBStateBackend`, which stores state as serialised byte arrays. With `MapState<UK,UV>`, Flink maps each user key to its own backend entry, so `get(uk)` and `put(uk, uv)` deserialise or write just that one entry. With `ValueState<Map<K,V>>` the whole map is a single value: reading one field means deserialising every entry, and updating one field means re-serialising and rewriting all of them. For a map with thousands of entries per key that turns an O(1) access into O(n) work on every record. `MapState` also gives you `entries()`, `keys()`, `values()`, `isEmpty()` and `remove(uk)` without materialising the whole map, and — when state TTL is configured — map entries expire independently, whereas a `ValueState` blob expires as a unit. On `HashMapStateBackend`, which keeps objects on the heap, the gap is much smaller.
code
java · 9 lines// Blob form: every access touches the entire map
private transient ValueState<Map<String, Long>> lastSeenBlob;
public void processElement(Event e, Context ctx, Collector<Event> out) throws Exception {
Map<String, Long> m = lastSeenBlob.value(); // deserialises ALL entries
if (m == null) m = new HashMap<>();
m.put(e.productId, e.timestamp);
lastSeenBlob.update(m); // re-serialises ALL entries
}go deeper
Recall that MapState exposes put, get, remove, contains, entries, keys, values and isEmpty, and that a ValueState holds exactly one value. Knowing the API surface of each is enough at this stage.
Explain the mechanics: RocksDB stores serialised bytes, so a ValueState holding a map is one blob read and rewritten in full on every access, while MapState composes a separate backend entry per map key.
Show the production consequence — per-record cost proportional to map size, RocksDB write amplification, no per-entry TTL — and be ready to say how you would migrate a live job that already picked the wrong shape.
Own the state-modelling standard: which collections are allowed to grow per key, what bounds them, and how state shape choices interact with backend selection, TTL policy and future schema evolution across the whole platform.
## The two shapes of "a map per key" You need to keep, for each keyed entity, a set of sub-keyed facts: per-user, a map of product id to last-seen timestamp. Flink lets you express that two ways. You can declare `MapState<String, Long>` with a `MapStateDescriptor`, or you can declare `ValueState<Map<String, Long>>` with a `ValueStateDescriptor` and put a `HashMap` inside it. They look interchangeable in the code. They are not interchangeable at runtime, and the reason is how each state backend physically stores a value. ## Why the backend decides the answer Flink bundles `HashMapStateBackend` and `EmbeddedRocksDBStateBackend`. `HashMapStateBackend` holds state internally as ordinary Java objects on the TaskManager heap. Reading a `ValueState` hands you back the very object you stored — no serialisation on the access path at all. So a `HashMap` inside a `ValueState` costs roughly the same as `MapState` for point access; the difference is mostly about API convenience and snapshot size. (Because objects are shared rather than copied, this backend is also the one where mutating a value you read out is unsafe.) `EmbeddedRocksDBStateBackend` is different in kind. It keeps in-flight state in an embedded RocksDB instance in the TaskManager's local data directories, and everything is stored as **serialised byte arrays**, with key comparisons done byte-wise rather than through `hashCode()`/`equals()`. Every read has to deserialise; every write has to serialise. That is the whole reason RocksDB scales past heap, and it is also the reason the granularity of a state entry matters. ## The cost model, concretely Under RocksDB with `ValueState<Map<String, Long>>` containing 5,000 entries: - `map = state.value()` deserialises 5,000 pairs to build the `HashMap`. - You mutate one entry. - `state.update(map)` re-serialises all 5,000 pairs and writes one large value. Do that once per record and per-record work is proportional to the size of the map, not to the size of the change. Throughput collapses as the map grows, and RocksDB is doing large-value writes that inflate its write amplification. With `MapState<String, Long>`, Flink composes the backend's storage key from the keyed-state key *and* the user map key, so each pair is its own RocksDB entry: - `state.get("p-99")` deserialises one value. - `state.put("p-99", ts)` writes one small entry. Per-record work is now independent of map size. `remove(uk)` deletes one entry, `contains(uk)` probes one entry, and `isEmpty()` avoids pulling anything into the JVM. Only the iterating methods — `entries()`, `keys()`, `values()` — walk the whole map, and they stream lazily rather than materialising a copy. ## Two more differences that matter in production **TTL granularity.** All of Flink's state collection types support per-entry TTLs: list elements and map entries expire independently. So with `MapState` plus a `StateTtlConfig`, a stale product entry disappears on its own while the rest of the user's map survives. With `ValueState<Map<...>>` the blob is one state entry: either the whole map is still live or the whole map expires. If your reason for keeping the map is "remember recent things and forget old ones", `MapState` is doing the forgetting for you and `ValueState` is not. **Size limits.** RocksDB's JNI bridge is `byte[]`-based, so the maximum supported size is 2^31 bytes per key and per value. A `ValueState` holding a very large map is a single value marching towards that ceiling. Related: states implemented with RocksDB merge operations, such as `ListState`, can silently accumulate past 2^31 bytes and then fail on the next retrieval — another reason unbounded per-key collections are a design smell. ## Where ValueState-of-map is still fine It is not always wrong. If the map is small and bounded — a handful of counters, a fixed set of flags — and you almost always read or write all of it together, the blob form is simpler and can even be cheaper: one serialised value instead of several backend entries, and one write instead of several. If you are on `HashMapStateBackend` and your state comfortably fits in heap, the performance argument largely evaporates, though the TTL-granularity argument does not. The rule of thumb: if the collection's size grows with input, or you access one entry at a time, use `MapState`. If it is a fixed-size record you touch as a unit, a `ValueState` of a small object is fine — and often a POJO is a better choice than a `Map`, because POJO types are one of the two families whose schema Flink can evolve across a savepoint. ## Migrating an existing job Switching a running job from `ValueState<Map<K,V>>` to `MapState<K,V>` is *not* a schema evolution Flink performs for you — the state name and structure change, so the old bytes are simply not read. Either give the new state a different name and accept starting empty, or take a savepoint and rewrite it offline with the State Processor API.
- Does the same argument apply on HashMapStateBackend?Much less. `HashMapStateBackend` keeps state as Java objects on the heap, so reading a `ValueState` returns the stored object with no deserialisation — a `HashMap` inside it is just a heap object you mutate. The per-access cost argument largely disappears. What survives is TTL granularity, since a `ValueState` blob still expires as a unit, and the fact that the whole map is written on every full snapshot.
- How does Flink physically distinguish two MapState entries for the same keyed key?Under `EmbeddedRocksDBStateBackend` each registered state gets its own RocksDB column family, and inside it the storage key is composed from the key group, the keyed-state key, the namespace (the window or other scope the state belongs to) and — for `MapState` — the serialised user map key. Every user map entry therefore becomes its own RocksDB record, which is exactly what makes a single `get` or `put` touch one entry instead of the whole map.
- What breaks if you swap ValueState<Map<K,V>> for MapState<K,V> and restore from a savepoint?The existing state is not migrated. The registered state's structure and serializer differ, so Flink cannot read the old bytes into the new descriptor; at best the new state starts empty, at worst the restore fails. Plan it as a data migration: read the savepoint with the State Processor API's `SavepointReader`, rewrite the entries into the new state shape, and start from the rewritten savepoint.
saying these in an interview costs you the question
- Says MapState and ValueState-of-map perform identically
- Thinks RocksDB stores Java objects rather than bytes
- Claims MapState.entries() loads the whole map into heap
- Believes TTL expires a ValueState blob entry by entry
- Says swapping the two state types migrates automatically