Command Palette

Search for a command to run...

Hectal
PHASE 11Advanced ~8 min· topic 2 of 5

Topic 11.2

State Stores, Changelogs, and Recovery

In one line

Stateful operations keep local state in RocksDB (or in-memory) stores, and every change is also written to a compacted changelog topic. When a task moves to another instance, its state is restored from the changelog; standby replicas keep warm copies to make failover fast. Interactive queries expose state stores for reads.

0/5 · 0%

Think of it like this

A cashier with a local till (state store) who also writes every change in a carbon-copy book kept at head office (changelog). If the cashier leaves, a new cashier rebuilds the till from the carbon copies; a trainee who's been copying along all day (standby replica) can take over instantly.

Key ideas

  1. 01

    Stores: KeyValueStore, WindowStore, SessionStore, persistent (RocksDB on local disk) or in-memory. Each store has a changelog topic named <application.id>-<store>-changelog, compacted.

  2. 02

    Recovery: when an instance starts or takes over a task, it replays the changelog from its last checkpoint into the local store before processing. With large state (tens of GB), restoration can take minutes to hours.

  3. 03

    Standby replicas (num.standby.replicas=1): other instances continuously apply the changelog to shadow copies, so failover only needs a small catch-up. Persistent volumes (StatefulSets) also keep local state across restarts.

  4. 04

    Warmup and assignment: Kafka Streams (KIP-441) assigns tasks to instances whose state is caught up and warms up others in the background, avoiding long pauses when scaling out.

  5. 05

    Interactive queries: streams.store(StoreQueryParameters.fromNameAndType(...)) reads local state; with KafkaStreams.queryMetadataForKey, an API layer can route a key lookup to the instance holding it.

  6. 06

    Rebalances interact with state: moving a stateful task means restoring state, so cooperative rebalancing, static membership and standbys matter even more here.

Code & diagrams

BalanceTopology.javajava
KTable<String, Long> balances = b.stream("ledger.transactions", Consumed.with(Serdes.String(), txnSerde))
    .groupByKey()                                              // key = accountId
    .aggregate(
        () -> 0L,
        (accountId, txn, balance) -> balance + txn.amountMinor(),   // +100_00, -30_00 ...
        Materialized.<String, Long, KeyValueStore<Bytes, byte[]>>as("balances")
            .withKeySerde(Serdes.String()).withValueSerde(Serdes.Long()));

balances.toStream().to("ledger.balances", Produced.with(Serdes.String(), Serdes.Long()));

props.put(StreamsConfig.PROCESSING_GUARANTEE_CONFIG, StreamsConfig.EXACTLY_ONCE_V2);
props.put(StreamsConfig.NUM_STANDBY_REPLICAS_CONFIG, 1);
props.put(StreamsConfig.STATE_DIR_CONFIG, "/var/lib/kafka-streams");      // persistent volume
state-recovery.mermaiddiagram
Rendering diagram…

Interview problem

The problem

Account balance from deposit and withdrawal events

Events: Deposit +100, Withdraw −30, Deposit +50. Build a stream-processing model that maintains each account's current balance. Discuss the KTable, state store, replay, recovery and changelog topic.

When it breaks

Stateful app on ephemeral pod storage with no standby replicas

What you see

Every pod restart restores the full state from the changelog: tens of minutes of downtime per deploy for large stores.

Fix & prevent

Persistent volumes, standby replicas, static membership and rolling restarts.

Explain it without notes

01

How does Kafka Streams make local RocksDB state fault-tolerant?

Practice

01

Build the balance topology, kill an instance mid-stream, and verify balances remain correct on the other instance.

Trade-offs

  • ↔

    Local state gives very fast lookups; the price is restoration time and disk management for large stores.

Done when you can

  • I can design stateful processing with stores, changelogs, standbys and interactive queries.