skip to content

Why switch a stateful Spark Structured Streaming query to RocksDBStateStoreProvider?

level: seniorimportance: nice to knowfreq 30%

answer

  1. where the keyed map physically lives
  2. pauses, not out-of-memory errors
  3. one option takes state off the heap
  4. you pay in bytes and in disk
  5. one partition count you cannot change later

basics

~20 s

The default state store keeps every key in the executor's JVM heap, so large state means heap pressure and long GC pauses. RocksDBStateStoreProvider moves state into native memory and local disk, letting state exceed the heap at the cost of serialization and disk I/O.

solid answer

~50 s

Structured Streaming's default provider, `HDFSBackedStateStoreProvider`, holds the entire in-memory map of state for each partition on the JVM heap and writes delta and snapshot files to the checkpoint for recovery. That is fast for small state and turns into steadily worse GC behaviour as the key space grows — long pauses, then executors lost to heartbeat timeouts. Setting `spark.sql.streaming.stateStore.providerClass` to `org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider` stores state in RocksDB's native memory and local disk instead, so total state can exceed the heap and GC stops tracking it. The trade is serialization on every read and write, real disk I/O, and a separate memory budget for RocksDB that must be carved out of the executor's overhead rather than its heap. Note that state is partitioned by `spark.sql.shuffle.partitions`, which you cannot change for a stateful query once the checkpoint exists — so size that before launch.

code

properties · 3 lines
properties
spark.sql.streaming.stateStore.providerClass=org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider
spark.executor.memoryOverhead=4g
spark.sql.shuffle.partitions=200

go deeper

for a junior

Know that a stateful streaming query keeps a keyed store per partition, that it is checkpointed for recovery, and that by default it lives in the executor's Java heap.

for a middle

Explain the trade concretely: on-heap map with no serialization versus native memory plus local disk with serialization on every access, and why the second one scales past the heap.

for a senior

Show you make the call from evidence — rising numRowsTotal and GC time — and that you budget executor memory overhead when you switch. Be ready to say why shrinking the watermark or bounding a join might beat changing providers.

for a principal

Own the pre-launch decisions this forces: shuffle-partition sizing that cannot be revisited, local disk provisioning per executor, and a policy for when a pipeline's state is large enough that its history should live in a table rather than in streaming state.

## Where streaming state lives Every stateful operator — windowed aggregation, `dropDuplicates`, stream–stream join, `flatMapGroupsWithState` — keeps a keyed map that must survive across micro-batches and across failures. Spark abstracts that behind a *state store provider*, one instance per shuffle partition per operator. The provider owns both the hot path (get and put on every record) and durability (writing versioned data into the query's checkpoint so a replayed batch can rewind state to the right version). Open-source Spark ships two providers, and the choice between them is one of the few tuning decisions that changes a streaming job's failure mode rather than just its speed. ## The default: HDFSBackedStateStoreProvider The default keeps the whole partition's state as an in-memory map **on the JVM heap**, and writes a delta file per batch into `state/` under the checkpoint, periodically compacting deltas into a snapshot. The name refers to where durability goes, not where the data lives at runtime. For small state — a few hundred thousand keys per partition — this is excellent: every access is a plain hash-map lookup with no serialization. The problem is that streaming state is long-lived by construction. A deduplication set, a session store, or an aggregation over a large key space grows until the heap is dominated by objects that are, by definition, never garbage. Old-generation occupancy climbs, full GCs get longer, and the symptom that reaches you is not "out of memory" but *stalls*: batch duration spikes, then an executor misses heartbeats and is lost, then its state has to be reloaded from the checkpoint on another executor, which makes the next batch slower still. ## The alternative: RocksDBStateStoreProvider Setting `spark.sql.streaming.stateStore.providerClass` = `org.apache.spark.sql.execution.streaming.state.RocksDBStateStoreProvider` switches state into an embedded RocksDB instance per partition. RocksDB is a native LSM-tree store: hot data sits in a native-memory block cache and memtables, colder data spills to files on the executor's local disk. Because none of it is on the JVM heap, the garbage collector no longer walks it, and total state size is bounded by disk rather than by heap. The costs are real and worth stating in an interview, because a candidate who only names the benefit has not operated it: - **Serialization on every access.** Keys and values cross the JVM/native boundary as bytes, so per-record cost rises even when everything is cached. - **Disk I/O and compaction.** LSM compaction consumes CPU and disk bandwidth in the background, and the executor needs local disk sized for the state plus compaction headroom. - **A separate memory budget.** RocksDB's memory is off-heap, so it must be accounted for outside the executor heap — typically by raising `spark.executor.memoryOverhead`. Forgetting this converts a GC problem into the cluster manager killing the container for exceeding its memory limit, which looks like a regression. - **Tuning surface.** The `spark.sql.streaming.stateStore.rocksdb.*` family exposes block cache and write buffer sizing; changelog checkpointing is available to reduce the cost of uploading state files each batch by writing a change log instead of full snapshots. ## The decision Stay on the default when state is small and bounded — a short watermark, a modest key space, per-batch state in the tens of megabytes. Move to RocksDB when state is large or unbounded in practice: long dedup windows, session state over many users, stream–stream joins with wide time bounds, or any query where you can watch old-generation heap climb across hours and batch duration climb with it. The diagnostic that justifies the switch is `stateOperators[].numRowsTotal` in `lastProgress` trending up alongside GC time, not a hunch. Also note that managed Spark distributions do not necessarily default the same way as open-source Spark, so check what your platform actually runs before concluding which provider a job is using. ## The constraint that catches people Whichever provider you use, state is sharded by `spark.sql.shuffle.partitions`, and that number is baked into the checkpoint's state layout. For a **stateful** query you cannot change it across restarts — the state stored for partition 37 of 200 has no meaning under a 400-partition layout. This is one of the few Spark settings you must get approximately right *before* the query first runs, because fixing it later means starting from a fresh checkpoint and rebuilding or abandoning history. Size it for the state you expect at peak, not for the throughput of the first week. ## Reducing state instead of relocating it Before changing providers, ask whether the state should be that large at all. A watermark threshold sized from intuition rather than measured lag is the most common cause of bloat; so is `dropDuplicates` without a watermark, which retains seen keys forever, and a stream–stream join with no time-range condition. Shrinking the threshold, adding `dropDuplicatesWithinWatermark`, or bounding a join often removes the problem more cheaply than moving gigabytes of state to disk. RocksDB is the right answer when the state is genuinely required.

  • What symptom points at heap-resident state rather than a slow sink?
    Batch duration climbing steadily over hours with input rate flat, growing `stateOperators[].numRowsTotal` in `lastProgress`, and executor GC time rising as a share of task time — often ending with executors lost to heartbeat timeouts rather than an OutOfMemoryError. A slow sink instead shows up as time concentrated in the write stage with state size flat.
  • Why must you raise executor memory overhead when enabling RocksDB state?
    RocksDB's block cache, memtables and index blocks are native memory, outside the JVM heap. The cluster manager enforces a limit on the container's total memory, so unaccounted native usage gets the executor killed for exceeding it. Budget the RocksDB memory in `spark.executor.memoryOverhead` — otherwise switching providers trades GC pauses for container kills, which looks like a regression.
  • Why can't you change spark.sql.shuffle.partitions for a stateful query after it starts?
    State is sharded by that number and stored per partition inside the checkpoint. Partition 37 of 200 holds a specific hash range; under 400 partitions that range no longer corresponds to anything, so the stored state cannot be reinterpreted. Spark therefore refuses the change. Size the setting before the first run, since fixing it later means a fresh checkpoint and rebuilt history.
  • What should you try before switching providers?
    Shrink the state. Size the watermark threshold from measured arrival lag rather than intuition, use `dropDuplicatesWithinWatermark` instead of an unbounded `dropDuplicates`, and put a time-range condition on stream–stream joins so unmatched rows are not buffered indefinitely. RocksDB is the right answer once the state is genuinely needed, not as a first response to state you did not intend to keep.

saying these in an interview costs you the question

  • Thinks the default provider keeps state on disk because of its name
  • Enables RocksDB without budgeting extra executor memory overhead
  • Plans to raise shuffle partitions later to shrink per-partition state
  • Assumes RocksDB is faster per record than an on-heap map
  • Treats growing state as a sink problem rather than a watermark problem

context