Explain the Kafka Streams record cache and how commit.interval.ms and cache.max.bytes.buffering affect what downstream sees and when the changelog is written.
answer
- Record cache = dedup + batch in front of store
- Flush on size budget OR commit.interval.ms
- cache.max.bytes.buffering → statestore.cache.max.bytes
- commit.interval default 30s (100ms EOS)
- Cache affects emission/changelog timing, not store reads
basics
~20 sEach state store has an in-memory record cache that deduplicates and batches updates. Downstream operators and the changelog only see the latest value per key when the cache is flushed — which happens when the cache fills up or at every commit (commit.interval.ms). Disabling the cache makes every update flow through immediately.
solid answer
~50 sKafka Streams puts a **record cache** in front of each state store. Updates to a key first land in the cache; the cache **dedupes** (collapses multiple updates to the same key into the latest) and **batches** writes. Cached values are flushed downstream and to the **changelog** when (a) the cache exceeds its size budget — historically `cache.max.bytes.buffering`, replaced by `statestore.cache.max.bytes` — or (b) a **commit** occurs every `commit.interval.ms` (default 30s, or 100ms under EOS). Because of this, downstream operators and the changelog see *fewer, later* updates: only the latest value per key per flush, not every intermediate change. This is a feature for aggregations — it suppresses noisy intermediate emissions and reduces changelog volume — but it adds latency and means emission is not per-record. Setting cache size to 0 or `Materialized.withCachingDisabled()` forces every update through immediately (more, earlier records, higher changelog/output volume). The cache affects *forwarding/changelog timing*, not store correctness — the store itself always reflects the latest write.
go deeper
Know there's a cache that batches updates so downstream sees fewer, later records.
Explain the two flush triggers (size budget and commit.interval.ms) and that disabling caching emits every update.
Reason about latency vs changelog-volume trade-offs, EOS commit interval, and store-correctness-vs-emission-timing separation.
Set cache/commit policy across a fleet for latency SLAs and Kafka write budget, and distinguish suppress() from caching for windowed final results.
## What the record cache is In front of every state store, Kafka Streams maintains an in-memory **record cache** (sometimes called the dedup cache). When a stateful operator updates a key, the new value goes into the cache first. The cache serves two purposes: 1. **Deduplication / update suppression** — if the same key is updated many times between flushes, only the **latest** value is forwarded downstream and written to the changelog. Intermediate values are collapsed. 2. **Batching / write efficiency** — flushing many entries at once is cheaper than per-record I/O to RocksDB and Kafka. ## When the cache flushes The cache is flushed (entries forwarded downstream and written to the changelog) on **two triggers**: - **Size pressure:** when total cache usage exceeds the budget. The classic config is **`cache.max.bytes.buffering`** (default 10 MB, shared across all stores in an instance, divided among stream threads). Newer versions rename this to **`statestore.cache.max.bytes`**. - **Commit:** at every **commit**, governed by **`commit.interval.ms`** (default **30000** ms at-least-once; **100** ms under exactly-once for lower latency). A commit flushes caches, flushes producers, and commits offsets. So the **effective emission/changelog latency** is bounded by whichever trigger comes first — a busy key may flush on size pressure well before the commit, while a quiet key waits up to `commit.interval.ms`. ## Observable effects - **Downstream sees the latest value per key per flush, not every change.** A `count()` that increments a key 50 times between flushes emits *one* update (the final count for that flush window), not 50. This is why people see "missing intermediate results" — they're suppressed, not lost; the store is correct. - **Changelog volume drops** because superseded intermediate values are never written. Larger cache + longer commit interval ⇒ fewer changelog records ⇒ less Kafka write load, at the cost of latency. - **Latency rises** with cache size and commit interval — results appear later. ## Disabling the cache - `Materialized.withCachingDisabled()` (DSL) or `StoreBuilder.withCachingDisabled()` (PAPI), or globally setting the cache size to **0**, makes **every update flow through immediately**: more records downstream, every intermediate value written to the changelog, lower latency, higher volume. Useful when you genuinely need per-update emission (e.g., a change-data feed) or for deterministic testing. ## Important boundaries - The cache governs **forwarding and changelog timing**, *not* store correctness. A read of the store (including Interactive Queries) always reflects the most recent write, cached or not. - The cache is **per instance, shared across that instance's stores/threads** under the global byte budget — it is not per-key or per-store unbounded. - `commit.interval.ms` also drives offset commits and producer flush; lowering it for fresher output increases commit overhead and (under EOS) transaction frequency. - Don't confuse this **record cache** with the **RocksDB block cache** (a separate native LRU for RocksDB reads, tuned via a `RocksDBConfigSetter`). ## Tuning intuition - Want low-latency, every-update output → small/zero cache, short commit interval. - Want efficient aggregation with suppressed noise and minimal changelog/output → larger cache, longer commit interval (and consider `suppress()` for windowed final-result semantics).
- A developer complains their aggregation 'isn't emitting every update.' What's the cause and fix?The record cache is collapsing intermediate updates and only flushing the latest per key on size/commit. Fix by reducing/zeroing the cache or calling Materialized.withCachingDisabled() if per-update emission is required — accepting higher output and changelog volume.
- Does the cache change what an Interactive Query reads from the store?No. Interactive Queries and any store read always return the most recent value, cached or flushed. The cache only delays forwarding to downstream operators and the changelog, not store correctness.
- Why is commit.interval.ms much lower (100ms) under exactly-once?EOS commits a Kafka transaction each interval; a long interval would make end-to-end latency very high and hold transactions open. Streams defaults EOS to 100ms to keep latency reasonable, trading more frequent commits.
saying these in an interview costs you the question
- Claiming the cache changes the value stored or read from the store (it doesn't — only emission timing)
- Saying intermediate updates are lost rather than suppressed/collapsed
- Confusing the Streams record cache with RocksDB's block cache
- Forgetting that commit.interval.ms also forces a cache flush, not just size pressure