skip to content

When would you choose an in-memory state store over the default RocksDB store, and what are the trade-offs and tuning levers for RocksDB?

level: seniorimportance: nice to knowfreq 35%

answer

  1. RocksDB default: on-disk, off-heap, > RAM, warm restart
  2. In-memory: heap-bound, GC, full changelog restore each start
  3. In-memory only for small bounded latency-sensitive state
  4. RocksDBConfigSetter: block cache, write buffer, bloom
  5. Bound off-heap (LRUCache + WriteBufferManager) or container OOM

basics

~20 s

RocksDB (the default) persists state to local disk, so it handles state larger than memory and restores faster from local files. In-memory stores are faster and avoid native/disk overhead but are bounded by heap and always restore fully from the changelog. Choose in-memory only for small, bounded state where heap pressure and full restores are acceptable.

solid answer

~50 s

**RocksDB** is the default persistent store: an embedded log-structured KV engine that spills to local disk under `state.dir`. It supports state much larger than RAM, keeps a local checkpoint so a warm restart replays only the changelog tail, and uses off-heap memory for its caches. **In-memory** stores (`Stores.inMemoryKeyValueStore`) keep everything on the JVM heap — lower per-op latency and no native/disk overhead, but state is capped by heap, increases GC pressure, and must be **fully restored from the changelog on every restart** (no local persistence to checkpoint against). Pick in-memory for small, bounded, latency-sensitive state where full cold restores are tolerable; pick RocksDB (the default) for large or unbounded state and faster local recovery. RocksDB tuning is done via a `RocksDBConfigSetter` (`rocksdb.config.setter`): block cache size, write buffer (memtable) size/count, and bloom filters; you can also bound total off-heap usage with a shared `WriteBufferManager`/`LRUCache`. RocksDB memory is off-heap, so it doesn't show up as JVM heap but can OOM the container if unbounded.

go deeper

for a junior

Know RocksDB is the default on-disk store and an in-memory option exists for small state.

for a middle

Contrast disk-vs-heap, capacity, and that in-memory always fully restores from the changelog.

for a senior

Tune RocksDB via RocksDBConfigSetter and bound off-heap memory; justify in-memory only for small bounded state.

for a principal

Plan fleet memory budgets across tasks×stores, container limits vs off-heap, and store-type policy as an architectural standard.

## The two store implementations Kafka Streams ships two state-store backends: - **Persistent (RocksDB)** — the **default**. RocksDB is an embedded, log-structured-merge (LSM) key-value store written in C++, accessed via JNI. It writes data to **local disk** under `state.dir`, using **off-heap** memory for its block cache and write buffers (memtables). - **In-memory** — a heap-resident map (`Stores.inMemoryKeyValueStore`, `inMemoryWindowStore`, etc.). All entries live on the **JVM heap**. ## Why RocksDB is the default 1. **State > RAM.** LSM-on-disk means the store can hold far more data than fits in memory; cold data lives on disk, hot data in the block cache. 2. **Local durability for fast restart.** RocksDB files persist on disk plus a **checkpoint file**; on a warm restart the instance reuses local state and replays only the **changelog tail**, not the whole topic. 3. **Off-heap memory.** Its caches don't bloat the JVM heap or directly drive GC. Costs: native JNI overhead per operation, disk I/O, and off-heap memory that — if unbounded — can OOM the **container** even though JVM heap looks fine. RocksDB also has its own compaction overhead. ## Why (and when) to choose in-memory In-memory stores trade durability and capacity for simplicity and latency: - **Pros:** lowest per-operation latency (pure heap access), no native library, no disk, deterministic behavior — handy in tests and for tiny lookups. - **Cons:** capped by heap; large state causes **GC pressure** and OOM risk; and crucially, there is **no local persistence**, so on *every* restart the store must be **fully restored from the changelog** — slower failover for anything but tiny state. (In-memory stores are still fault-tolerant via their changelog, just not locally persistent.) Rule of thumb: in-memory only for **small, bounded** state where full cold restores are acceptable and latency matters; otherwise keep the RocksDB default. ## Tuning RocksDB RocksDB is tuned through a **`RocksDBConfigSetter`** implementation wired via the `rocksdb.config.setter` config. Common levers: - **Block cache** (`BlockBasedTableConfig.setBlockCache`) — read cache size; bigger = fewer disk reads. - **Write buffer / memtable** (`setWriteBufferSize`, `setMaxWriteBufferNumber`) — in-memory write batching before flush to SST files. - **Bloom filters** (`setFilterPolicy`) — speed up point lookups by skipping SST files. - **Bounded total off-heap** — share a single `LRUCache` and a `WriteBufferManager` across all RocksDB instances to **cap aggregate off-heap memory** (critical in containers with memory limits; otherwise N stores × buffers can blow the cgroup limit). Don't forget to **close** native objects you create in the config setter. ## Practical gotchas - **Off-heap ≠ free.** Container OOMKills often trace to unbounded RocksDB memory, not heap. Bound it. - **state.dir on fast local disk** matters for RocksDB throughput and restart speed; ephemeral `/tmp` defeats the local-reuse benefit. - **Many stores multiply memory.** Each task × each store has its own RocksDB instance and buffers; account for the product. - Switching a store to in-memory via `Materialized.as(Stores.inMemoryKeyValueStore(name))` keeps the changelog and fault tolerance, only changing local persistence/latency characteristics.

  • Why can a Kafka Streams container get OOMKilled even though JVM heap usage looks healthy?
    RocksDB allocates its block cache and write buffers off-heap (native). Many stores/tasks multiply that, and if it's unbounded it can exceed the container's memory limit. Bound it by sharing a single LRUCache and WriteBufferManager via a RocksDBConfigSetter.
  • Are in-memory stores fault-tolerant?
    Yes — they still have a changelog and can be rebuilt. They just have no local on-disk persistence, so every restart triggers a full changelog restore rather than a tail replay from a checkpoint.

saying these in an interview costs you the question

  • Saying in-memory stores aren't fault-tolerant (they are, via changelog) — the real difference is no local persistence
  • Thinking RocksDB memory is on the JVM heap (it's off-heap/native)
  • Recommending in-memory for large state without noting heap/GC/full-restore costs
  • Ignoring that unbounded RocksDB off-heap memory OOMs containers

context