skip to content

Streams DSL and Topology

Building a topology from KStream, KTable and GlobalKTable, and reading how it splits into sub-topologies. Interviewers ask because describing your topology reveals where repartition topics and state actually appear.

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

questions

5

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%

answer

  1. KStream = events (append)
  2. KTable = upserts (latest per key)
  3. null in KTable = tombstone/delete
  4. stream-table duality
  5. KTable materialized + compacted changelog

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.

solid answer

~50 s

A KStream models an unbounded sequence of independent events: every record is appended, and two records with the same key are two distinct facts (e.g. 'page-view', 'click'). A KTable models a changelog: each record is an UPSERT for its key, so the table holds only the latest value per key, and a null value is a tombstone (delete). Use a KStream for event/transaction streams where every occurrence matters (logs, clicks, payments). Use a KTable for mutable entity state where you only care about the current value (user profile, account balance, inventory level). Under the hood a KTable is materialized into a local state store backed by a compacted changelog topic, which is what enables lookups and joins against current state. The DSL lets you convert between them: stream().toTable() / aggregate() turns a stream into a table, and table.toStream() emits the changelog as events.

go deeper

for a junior

Recall the one-liner: KStream = events, KTable = latest value per key; null deletes in a KTable.

for a middle

Explain stream-table duality and the conversions (toTable/aggregate/toStream) and when each fits.

for a senior

Discuss materialization into a RocksDB state store + compacted changelog, tombstone semantics, and cache/commit-interval compaction of updates.

for a principal

Reason about modeling choices for a whole pipeline: which abstraction minimizes repartitions and state, and the correctness implications of tombstones and update compaction downstream.

## First principles Kafka Streams gives you three core abstractions built on top of Kafka topics. Two of them — **KStream** and **KTable** — express the famous **stream-table duality**: a stream of events and a table of state are two views of the same data. ### KStream — a record (event) stream A `KStream<K, V>` represents an **unbounded, append-only sequence of records**. Each record is an **independent, immutable fact**. If two records arrive with the same key, they are two separate events, NOT an update of one another. Example: a topic of page views ``` ("alice", view-home) , ("alice", view-cart) ``` These are two distinct page views by Alice — you keep both. ### KTable — a changelog (update) stream A `KTable<K, V>` represents an **evolving table** where each record is an **UPSERT** (insert-or-update) for its key. Only the **latest value per key** is meaningful. A record with a **null value is a tombstone** that deletes the key. Example: a topic of account balances ``` ("alice", 100) , ("alice", 80) ``` The table now holds `alice -> 80`; the 100 is overwritten. A KTable is **materialized** into a local **state store** (RocksDB by default), backed by a **compacted changelog topic** for fault tolerance. ### Stream-table duality - A **stream is a changelog of a table**: replay the changelog and you rebuild the table. - A **table is a snapshot of a stream**: take the latest value per key and you get the table. The DSL exposes conversions: - `KStream.toTable()` or aggregations (`groupByKey().reduce()/aggregate()/count()`) → KTable. - `KTable.toStream()` → emits the table's changelog as a KStream of events. ### When to use which - **KStream**: clickstreams, transactions, IoT readings, audit logs — anything where every occurrence is a fact you must process. - **KTable**: current entity state — user profiles, configuration, inventory, balances — where you only care about the latest value and want efficient key lookups/joins. ### Edge cases - **Null semantics differ**: in a KStream a null value is just a record (and may be filtered); in a KTable a null value **deletes** the key (tombstone). - KTable updates are **not** guaranteed to emit one downstream record per input — record caching and `commit.interval.ms` can **compact** intermediate updates so you only see the last value within a window of time (this is the suppression/cache behavior, not a bug). - A KTable needs a **keyed** source; if your data isn't keyed by the entity, you must `groupBy`/`selectKey` and aggregate, which triggers a repartition.

  • What does a null value mean in a KStream versus a KTable?
    In a KStream a null value is just another record (you may filter it). In a KTable a null value is a tombstone that deletes the key from the table and its state store/changelog.
  • How do you convert a KStream into a KTable?
    Either KStream.toTable() (treat the stream directly as a changelog) or a stateful aggregation: stream.groupByKey().reduce()/aggregate()/count(), which produces a KTable. groupBy with a new key triggers a repartition first.

saying these in an interview costs you the question

  • Saying a KTable keeps all records like a KStream — it only keeps the latest value per key.
  • Claiming KStream supports key-based lookups/joins-against-state — that's the KTable's materialized store.
  • Forgetting that a null value deletes a key in a KTable (treating it as a normal record).

context

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