skip to content

Kafka Streams

The Kafka Streams library: DSL topologies, state stores, joins, windowing and time semantics, the Processor API, and exactly-once processing. Interviewers use it to test stream-processing thinking without a separate cluster in the picture.

part ofApache Kafkaoverview, primer and where to startread it →
on this pageshow

explore

questions

61 · 11 sections

What is the difference between a KStream and a KTable in the Kafka Streams DSL, and when would you use each?

level: juniorimportance: must knowfreq 85%
basics
~20 s

A KStream is a record stream where every record is an independent event (insert/append). A KTable is a changelog stream where each record is an update keyed by its key — only the latest value per key matters.

open as a page

What is a GlobalKTable, and how does it differ from a KTable in partitioning, joins, and co-partitioning requirements?

level: seniorimportance: must knowfreq 65%
basics
~20 s

A KTable is partitioned: each task holds only its partition's keys. A GlobalKTable is fully replicated: every instance holds ALL partitions of the source topic. So GlobalKTable joins need no co-partitioning and can join on a non-key field.

open as a page

Using StreamsBuilder, how do you create stream(), table(), and globalTable() sources, and what does each produce?

level: juniorimportance: should knowfreq 60%
basics
~10 s

StreamsBuilder is the DSL entry point. builder.stream(topic) returns a KStream (event stream), builder.table(topic) returns a KTable (latest-value-per-key, materialized), and builder.globalTable(topic) returns a GlobalKTable (fully replicated). builder.build() produces the Topology.

open as a page

How does Kafka Streams split a topology into sub-topologies, and what role do repartition topics play?

level: seniorimportance: should knowfreq 55%
basics
~20 s

A topology is split into sub-topologies at points where data must be re-shuffled by key — repartition topics. Operations that change the key (selectKey, map, groupBy) force a repartition; the topology breaks there into independently-scheduled sub-topologies.

open as a page

How do you inspect a Kafka Streams topology, and why does naming DSL operators and stateful stores matter for topology compatibility?

level: principalimportance: should knowfreq 40%
basics
~10 s

Call topology.describe() to print the processor graph (sub-topologies, sources, sinks, stores). Naming operators (via Named/Materialized/Repartitioned/Joined) gives stable node and internal-topic names, so refactoring the DSL doesn't change changelog/repartition topic names and break state restore.

open as a page

In Kafka Streams, what does it mean for a DSL operation to be 'stateless', and which common operators fall into that category?

level: juniorimportance: must knowfreq 70%
basics
~10 s

A stateless operation processes each record independently, keeping no memory of past records. Examples: map, mapValues, filter, filterNot, flatMap, flatMapValues, selectKey, branch/split, merge, peek, foreach.

open as a page

What is the practical difference between map and mapValues (and filter vs transformValues) in Kafka Streams, and why would you prefer mapValues?

level: middleimportance: must knowfreq 65%
basics
~10 s

map can change both key and value; mapValues changes only the value and keeps the key. Prefer mapValues because keeping the key avoids a repartition topic when a stateful operation follows.

open as a page

Explain exactly when and how Kafka Streams creates a repartition topic from a stateless key-changing operator, and how to control or avoid it.

level: seniorimportance: must knowfreq 60%
basics
~20 s

Key-changing stateless ops (selectKey, map, flatMap) set a repartition-required flag. When a stateful op follows, Streams injects an internal repartition topic (named app-id + node + '-repartition') to re-shuffle by the new key. Use mapValues to avoid it, or repartition() to control it.

open as a page

How do you route a single KStream into multiple branches in modern Kafka Streams, and how do you recombine streams? Contrast split/branch with merge.

level: middleimportance: should knowfreq 45%
basics
~10 s

Use KStream.split() then chained .branch(predicate, Branched.as("name")) and optionally .defaultBranch(), which returns a Map<String, KStream> of routed sub-streams. merge(otherStream) does the opposite: it interleaves two same-typed streams into one. Both are stateless.

open as a page

Compare flatMap/flatMapValues with the side-effect operators peek and foreach. When is each appropriate, and what are the failure-semantics pitfalls?

level: middleimportance: should knowfreq 40%
basics
~20 s

flatMap/flatMapValues turn one input record into zero-or-many output records (flatMapValues keeps the key). peek runs a side effect per record and passes records through unchanged. foreach is terminal: it runs a side effect and ends the branch with no downstream. Side effects aren't transactional, so they can repeat on retries.

open as a page

What is a state store in Kafka Streams, and why does a stream-processing application need one?

level: juniorimportance: must knowfreq 70%
basics
~20 s

A state store is a local key-value store on each Kafka Streams instance that holds data needed for stateful operations like aggregations, joins, and windowing — so the app can remember results across messages instead of treating each event in isolation.

open as a page

How does Kafka Streams make a local state store fault-tolerant? Explain changelog topics.

level: middleimportance: must knowfreq 75%
basics
~20 s

Every update to a state store is also written to a special Kafka topic called a changelog. The changelog is the durable backup: if an instance dies, another instance replays the changelog to rebuild the store. Changelogs are log-compacted so they stay bounded.

open as a page

How do you create and configure a state store in Kafka Streams? Contrast Materialized (DSL) with Stores/StoreBuilder (Processor API).

level: middleimportance: should knowfreq 55%
basics
~20 s

In the DSL you configure stores with Materialized — passing it to operators like count() or aggregate() to set the store name, serdes, and persistence. In the lower-level Processor API you build a store explicitly with the Stores factory and a StoreBuilder, then register it on the topology.

open as a page

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.

level: seniorimportance: should knowfreq 45%
basics
~20 s

Each 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.

open as a page

Walk through what happens to state stores during a rebalance, and explain how standby replicas and state.dir affect restoration time.

level: seniorimportance: should knowfreq 50%
basics
~20 s

When tasks move between instances during a rebalance, the new owner must rebuild each store by replaying its changelog topic before it can process records. This restore can be slow. Standby replicas keep warm copies on other instances to avoid or shorten it, and a persistent state.dir lets an instance reuse local state instead of restoring from scratch.

open as a page

What kinds of joins does Kafka Streams support, and how do they differ at a high level?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Kafka Streams supports stream-stream joins (over a time window), stream-table joins (enrich a record with the latest table value), and table-table joins (combine two changelog tables). GlobalKTable joins let a stream join a fully-replicated table without matching partitioning.

open as a page

What is the co-partitioning requirement for Kafka Streams joins, and how do you satisfy or avoid it?

level: middleimportance: must knowfreq 65%
basics
~20 s

Co-partitioning means both join inputs must use the same key, the same number of partitions, and the same partitioning strategy, so matching keys land on the same task. You satisfy it by re-keying and repartitioning; GlobalKTable joins avoid it entirely.

open as a page

Explain KStream-KTable join semantics: what triggers output, what inner vs left produce, and why a GlobalKTable join differs.

level: middleimportance: must knowfreq 60%
basics
~20 s

A stream-table join enriches each stream record with the table's current value for that key. Only stream records trigger output (table updates don't). Inner emits only on a match; left emits the stream record paired with null when no match. GlobalKTable joins skip co-partitioning and join on any value-derived key.

open as a page

How does a foreign-key KTable-KTable join (KIP-213) work, and how is it different from an equi-key table join?

level: seniorimportance: should knowfreq 35%
basics
~20 s

A foreign-key join lets a left KTable join a right KTable on a key extracted from the left value, not the left's own key. Kafka Streams repartitions the left table by the foreign key internally, so the two tables don't need to be co-partitioned, and updates on either side correctly propagate.

open as a page

How do JoinWindows work for stream-stream joins, including grace period and inner vs outer/left emission timing?

level: seniorimportance: should knowfreq 50%
basics
~20 s

A JoinWindows defines a time difference: record A joins record B if their timestamps are within the window. The grace period allows late records before the window closes. Inner results emit immediately on a match; left/outer null-paired results emit only after the window plus grace expires.

open as a page

What are the four window types in Kafka Streams, and how do they differ?

level: juniorimportance: must knowfreq 78%
basics
~10 s

Tumbling (fixed, non-overlapping), hopping (fixed size, overlapping by advance interval), sliding (window around record pairs within a time difference), and session (activity-gap-based, dynamic size). The first three are time-based; sessions are data-driven.

open as a page

What is the grace period in Kafka Streams windowing (ofSizeAndGrace), and what happens to records that arrive after it?

level: middleimportance: must knowfreq 72%
basics
~20 s

The grace period is extra time after a window's end during which late (out-of-order) records are still accepted and update the window's result. Once stream time passes window-end + grace, the window is closed and later records are dropped.

open as a page

How do JoinWindows work in a KStream-KStream join, and what does the window size mean?

level: seniorimportance: should knowfreq 58%
basics
~20 s

JoinWindows define how close in event time two records from the two streams must be to join. With JoinWindows.ofTimeDifferenceAndGrace(d), records join if their timestamps differ by at most d (symmetric: a within [b-d, b+d]). Each stream is buffered in a window store for that span.

open as a page

How do session windows work, including session merging and the inactivity gap?

level: seniorimportance: should knowfreq 50%
basics
~20 s

A session window groups records for a key that arrive within an inactivity gap of each other. The window grows with each new record and closes when no record arrives for longer than the gap. Out-of-order records can merge two adjacent sessions into one.

open as a page

How are windowed aggregation results keyed and stored — explain Windowed keys, windowed serdes, and the segmented store layout?

level: principalimportance: should knowfreq 40%
basics
~20 s

A windowed aggregation produces a KTable keyed by Windowed<K> (the original key plus the window's start/end). It is stored in a segmented WindowStore (RocksDB split into time segments) and serialized with a WindowedSerdes that encodes key + window timestamp. Old segments are dropped wholesale when they fall outside retention.

open as a page

What are the three notions of time in Kafka Streams (event time, processing time, ingestion time), and how do they differ?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Event time = when the event actually happened (set by the producer). Ingestion time = when the broker appended the record. Processing time = when Kafka Streams processes it. They differ because of network and processing delays.

open as a page

What is the TimestampExtractor interface in Kafka Streams, and what do FailOnInvalidTimestamp, LogAndSkipOnInvalidTimestamp, and WallclockTimestampExtractor do?

level: middleimportance: must knowfreq 60%
basics
~10 s

TimestampExtractor tells Streams which timestamp to use per record. FailOnInvalidTimestamp (default) reads the record timestamp and throws on a negative one; LogAndSkipOnInvalidTimestamp logs and skips bad ones; WallclockTimestampExtractor ignores the record and uses System.currentTimeMillis().

open as a page

What is 'stream time' in Kafka Streams, how does it advance, and why does it never go backwards?

level: seniorimportance: must knowfreq 50%
basics
~20 s

Stream time is the app's notion of the current event-time clock per task: the maximum record timestamp seen so far. It only moves forward — a later record with an older timestamp does not pull it back. It drives window closing and lateness decisions.

open as a page

How does Kafka Streams handle out-of-order and late records in windowed operations, and what role does the grace period play?

level: seniorimportance: should knowfreq 45%
basics
~20 s

Out-of-order records (timestamp below stream time but within the grace period) are still added to their window and update the result. Once stream time passes window end plus the grace period, the window closes and later-arriving records for it are dropped as 'late'.

open as a page

What problem does max.task.idle.ms solve in Kafka Streams, and what are the trade-offs of tuning it?

level: principalimportance: should knowfreq 35%
basics
~10 s

When a task reads several partitions, max.task.idle.ms tells Streams how long to pause/wait for an empty-but-active partition to deliver data before processing ahead with other partitions. It trades join/merge time-ordering correctness against latency.

open as a page

In Kafka Streams, what is the difference between groupByKey() and groupBy(), and why does one of them trigger a repartition?

level: juniorimportance: must knowfreq 70%
basics
~20 s

groupByKey() groups by the existing record key with no repartition. groupBy() picks a NEW key, so Streams must repartition (re-shuffle records to topic partitions) so all records with the same new key land on the same task.

open as a page

Compare count(), reduce(), and aggregate() on a KGroupedStream. When must you use aggregate() with an Initializer and Aggregator?

level: middleimportance: must knowfreq 65%
basics
~20 s

count() returns how many records per key. reduce() combines two values of the SAME type into one (e.g. sum of longs). aggregate() is the general form: an Initializer creates a starting value of ANY type, and an Aggregator folds each record into it — use it when the result type differs from the input.

open as a page

What problem does suppress(Suppressed.untilWindowCloses(...)) solve in windowed aggregations, and what are its operational requirements and caveats?

level: seniorimportance: must knowfreq 55%
basics
~20 s

By default a windowed aggregation emits an update on every record, so downstream sees many intermediate counts per window. suppress(Suppressed.untilWindowCloses(...)) buffers updates and emits only the FINAL result for each window after the window plus grace period has closed. It needs a buffer config (e.g. unbounded or a bytes/records limit) and only works on windowed KTables.

open as a page

How do windowed aggregations work in Kafka Streams? Cover the window types, the resulting key type, and how late records and grace periods are handled.

level: seniorimportance: must knowfreq 60%
basics
~20 s

You call windowedBy(...) on a KGroupedStream before count/reduce/aggregate, so records are bucketed by time. The result is a KTable keyed by Windowed<K> (key + window bounds). Window types include tumbling, hopping, sliding, and session windows. Late records that arrive within the grace period still update their window; beyond grace they are dropped.

open as a page

An aggregation produces a KTable. Explain what the resulting KTable's update stream emits, and why a downstream consumer may see many intermediate values per key.

level: middleimportance: should knowfreq 50%
basics
~20 s

A KTable is a changelog: each key maps to its latest aggregate. After every input record, the aggregate for that key changes, so the KTable emits a new (key, newAggregate) update. Downstream therefore sees the running aggregate update again and again, not just the final value.

open as a page

What is the Processor API in Kafka Streams, and how does it differ from the high-level DSL?

level: juniorimportance: must knowfreq 60%
basics
~20 s

The Processor API (PAPI) is the low-level Kafka Streams API. You write a Processor that handles each record one at a time, can attach state stores, and forward results downstream. The DSL is a higher-level, declarative layer (map, filter, join) built on top of it.

open as a page

What does ProcessorContext.forward() do, and how do you control where a forwarded record goes?

level: middleimportance: must knowfreq 55%
basics
~20 s

context.forward(record) sends a record from your processor to its downstream child nodes. You can call it zero, one, or many times per input record. Passing forward(record, "childName") routes it to only one named child instead of all of them.

open as a page

Explain ProcessorContext.schedule() and the difference between PunctuationType.STREAM_TIME and WALL_CLOCK_TIME punctuators.

level: seniorimportance: must knowfreq 50%
basics
~20 s

schedule() registers a punctuator — a callback that fires periodically. STREAM_TIME advances by the timestamps of records flowing through, so it only fires when data arrives and progresses. WALL_CLOCK_TIME advances by the system clock, so it fires on a real-time interval even if no records arrive.

open as a page

How do you attach and access a state store from a custom Processor, including in the DSL via process()?

level: seniorimportance: must knowfreq 45%
basics
~20 s

Register the store with a StoreBuilder (addStateStore in PAPI, or builder.addStateStore in the DSL), connect it to the processor by name, then look it up in init() via context.getStateStore("name"). In the DSL, pass the store names as the second argument to process(supplier, "storeName").

open as a page

Walk through the lifecycle of a Processor: ProcessorSupplier.get(), init(), process(), and close(). Why is a supplier used instead of a single instance?

level: middleimportance: should knowfreq 35%
basics
~20 s

A ProcessorSupplier.get() returns a new Processor for each task/thread, so instances aren't shared across threads. init() runs once per task to grab the context and stores, process() runs per record, and close() runs once when the task shuts down to release resources.

open as a page

What are Interactive Queries in Kafka Streams, and how do you read the value for a key from a local state store?

level: juniorimportance: must knowfreq 70%
basics
~20 s

Interactive Queries let your app read directly from the state stores Kafka Streams already maintains, instead of querying an external database. You call KafkaStreams.store(...) to get a read-only store and look up a key with get().

open as a page

In a multi-instance Kafka Streams app, how do you find which instance can answer an Interactive Query for a given key?

level: middleimportance: must knowfreq 60%
basics
~20 s

A local store only holds the partitions assigned to its instance. To find the owner of a key, call streams.queryMetadataForKey(storeName, key, serializer). It returns the host:port (from application.server) of the instance that owns that key's partition, so you can route the query there.

open as a page

How would you build the RPC layer that turns local Interactive Queries into a cluster-wide queryable service?

level: seniorimportance: should knowfreq 45%
basics
~20 s

Run an HTTP (or gRPC) server on each instance with a query endpoint. On a request, use queryMetadataForKey to find the owning host. If it's local, read the local store and return the value; if remote, forward the request to that host's endpoint and relay the result.

open as a page

How do standby replicas (num.standby.replicas) change Interactive Query availability and routing, and what are the consistency tradeoffs?

level: seniorimportance: should knowfreq 35%
basics
~20 s

Setting num.standby.replicas > 0 makes other instances keep warm copies of a store's state. If the active instance fails, queryMetadataForKey still lists standby hosts, so you can serve reads from a standby for higher availability — but a standby may be slightly behind, so its data can be stale.

open as a page

What is the IQv2 API (KIP-796), and how do you interactively query windowed/session stores versus key-value stores?

level: principalimportance: nice to knowfreq 25%
basics
~20 s

IQv2 is a newer, type-safe query API (StateQueryRequest + Query types like KeyQuery/RangeQuery) returning a StateQueryResult, designed to be extensible and partition-aware. Windowed stores use a different queryable type (windowStore/sessionStore) and you query by key plus a time range, not just a key.

open as a page

What does setting processing.guarantee=exactly_once_v2 in a Kafka Streams application actually guarantee, and how do you enable it?

level: juniorimportance: must knowfreq 70%
basics
~20 s

It guarantees each input record affects the application's state and output exactly once, even if the app crashes and restarts — no duplicates and no lost updates. You enable it by setting the config processing.guarantee to exactly_once_v2.

open as a page

How does Kafka Streams keep consumer offsets, changelog writes, and output topic writes atomically consistent under EOS?

level: middleimportance: must knowfreq 60%
basics
~20 s

Streams uses a single Kafka transaction per commit. The output writes, changelog (state) writes, and the consumed input offsets are all added to that transaction and committed together via sendOffsetsToTransaction + commitTransaction, so they advance as one atomic unit.

open as a page

A team enabled exactly_once_v2 in a Streams app but a downstream service still sees duplicate/phantom records. What's the likely cause and fix?

level: middleimportance: must knowfreq 50%
basics
~10 s

The downstream consumer is almost certainly reading with isolation.level=read_uncommitted (the default), so it sees records from aborted transactions and uncommitted writes. Fix: set isolation.level=read_committed on the downstream consumer so it only reads committed output.

open as a page

How does commit.interval.ms behave under EOS, and how do you tune it against the latency/throughput tradeoff?

level: seniorimportance: should knowfreq 40%
basics
~20 s

Under EOS, commit.interval.ms defaults to 100 ms (vs 30000 ms at-least-once) because each commit ends a transaction, and downstream read_committed consumers only see output after a commit. Smaller interval = lower end-to-end latency but more commit overhead; larger interval = higher throughput but more latency and bigger replay on failure.

open as a page

What changed between the original exactly_once and exactly_once_v2 (KIP-447), and why was v2 introduced?

level: seniorimportance: should knowfreq 45%
basics
~20 s

The original exactly_once used one transactional producer per task (input partition), which didn't scale — many producers and slow rebalances. exactly_once_v2 (KIP-447) uses one producer per stream thread and fences zombies via consumer-group metadata, so it scales far better. v2 needs brokers >= 2.5.

open as a page

In Kafka Streams, what is a stream task and how does the framework decide how many tasks an application has?

level: juniorimportance: must knowfreq 70%
basics
~20 s

A task is the smallest unit of parallelism in Kafka Streams. The number of tasks equals the number of partitions of the busiest (most-partitioned) input topic in the topology. Each task processes a fixed set of partitions.

open as a page

What does num.stream.threads control, and how do threads relate to tasks and instances when scaling a Kafka Streams app?

level: middleimportance: must knowfreq 65%
basics
~20 s

num.stream.threads sets how many StreamThreads run inside one application instance. Tasks are distributed across all threads of all instances. You scale up (more threads per instance) or out (more instances), but never beyond the total number of tasks.

open as a page

What does num.standby.replicas do in Kafka Streams, and how does it affect failover and availability?

level: seniorimportance: must knowfreq 55%
basics
~20 s

num.standby.replicas keeps shadow copies of a stateful task's state store on other instances, continuously fed from the changelog topic. On failover, an up-to-date standby takes over almost immediately instead of rebuilding state from scratch, cutting recovery time.

open as a page

What are internal repartition and changelog topics in Kafka Streams, and how do they relate to scaling?

level: middleimportance: should knowfreq 40%
basics
~20 s

Streams auto-creates two kinds of internal topics: repartition topics (to re-shuffle data by a new key before stateful ops) and changelog topics (a durable backup of each state store). Their partition counts mirror the input, which also bounds parallelism.

open as a page

How does cooperative rebalancing work in Kafka Streams, and why was it introduced over the older eager protocol?

level: seniorimportance: should knowfreq 45%
basics
~20 s

Cooperative rebalancing (the default since 2.4) lets instances keep tasks they already own during a rebalance instead of giving everything up. Only the tasks that must move are revoked, avoiding a global stop-the-world pause and unnecessary state restoration.

open as a page