skip to content

ksqlDB Streaming SQL

ksqlDB's SQL over Kafka: streams versus tables, push versus pull queries, and materialized views on the Kafka Streams runtime. Interviewers ask when SQL is enough and when you need a real Streams application.

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

questions

6

In ksqlDB, what is the difference between a STREAM and a TABLE, and when would you use each?

level: juniorimportance: must knowfreq 78%

answer

  1. Stream = append-only facts
  2. Table = latest value per key (upsert)
  3. Null value = tombstone delete (table)
  4. Aggregation yields a TABLE
  5. Stream-table duality

basics

~20 s

A STREAM is an unbounded, append-only sequence of independent events (an immutable log). A TABLE represents the latest value per key (a mutable, upserted view). Use a STREAM for facts/events, a TABLE for current state.

solid answer

~50 s

ksqlDB models two core abstractions over Kafka topics. A STREAM is an append-only, unbounded sequence of immutable events — every record is an independent fact, and the same key can appear many times (e.g., a page-view log). A TABLE is a changelog interpretation: each record is an upsert keyed by the message key, so the table holds the latest value per key, and a record with a null value is a tombstone (delete). The same topic can be read as either, depending on the semantics you want: stream = 'a thing happened', table = 'the current state of a thing'. Aggregations (COUNT, SUM) produce TABLEs because they hold evolving per-key state. Choosing wrong changes join and aggregation behavior — e.g., a stream-stream join is windowed, a stream-table join is an enrichment lookup against current state.

go deeper

for a junior

Know stream = events/append-only, table = current state per key, and give one example of each.

for a middle

Explain upsert/tombstone semantics, PRIMARY KEY vs KEY, and that aggregations produce tables.

for a senior

Articulate the stream-table duality and how the choice changes join/aggregation behavior and topic retention/compaction.

for a principal

Reason about modeling a domain as facts vs. state, compaction strategy, and how duality underpins event-sourcing/CQRS designs.

## First principles: the stream-table duality Kafka stores data in **topics**, which are partitioned, ordered, append-only logs of records. Each record has a key, a value, a timestamp, and an offset. ksqlDB puts a SQL layer on top of these logs, and it interprets a topic in one of two ways. ### STREAM A **STREAM** treats the topic as an *append-only sequence of immutable events*. Every record is an independent fact. If key `user-42` appears five times, that is five distinct events — nothing is overwritten. Think of a transaction log, click events, or sensor readings. Semantically: 'this happened, then this happened, then this happened.' ```sql CREATE STREAM clicks (user_id VARCHAR, url VARCHAR) WITH (KAFKA_TOPIC='clicks', VALUE_FORMAT='JSON'); ``` ### TABLE A **TABLE** treats the topic as a *changelog*: each record is an **upsert** keyed by the record key. The table holds only the **latest value per key**. A record with the same key replaces the previous value; a record with a **null value** is a **tombstone** that deletes the key. Think of it as the materialized 'current state'. Semantically: 'the current value of X is …'. ```sql CREATE TABLE users (user_id VARCHAR PRIMARY KEY, country VARCHAR) WITH (KAFKA_TOPIC='users', VALUE_FORMAT='JSON'); ``` Note the **PRIMARY KEY** on a table vs. the (optional) **KEY** on a stream. ### The duality This is the famous **stream-table duality**: a table is the result of *replaying* a stream of upserts; a stream is the sequence of *changes* applied to a table. ksqlDB lets you move between them. An aggregation over a stream produces a table (evolving per-key state), and reading a table's underlying changelog gives you a stream again. ### Why the choice matters - **Aggregations** (`COUNT`, `SUM`, `LATEST_BY_OFFSET`) over a stream always yield a **TABLE**, because the running result per group is mutable state. - **Joins** behave differently: stream-stream joins must be **windowed** (you join events near each other in time); a stream-table join is a non-windowed **enrichment lookup** against the table's current state; table-table joins maintain a continuously-updated joined view. - **Retention/compaction**: tables are typically backed by **log-compacted** topics so the latest value per key survives; streams use time/size retention. ### Edge cases - A null-valued record is just data in a stream but a **delete** in a table. - Out-of-order records: a table upsert uses the record as-is; ksqlDB does not reorder by timestamp for plain table updates. - A stream without a declared key still works; a table **requires** a primary key to define the upsert semantics.

  • What does a record with a null value mean in a TABLE versus a STREAM?
    In a TABLE it is a tombstone — it deletes that key from the materialized state. In a STREAM it is just another event with a null payload; nothing is deleted.
  • Why does a COUNT aggregation produce a TABLE and not a STREAM?
    The running count per group is evolving, keyed state — the latest count per key — which is exactly the upsert/changelog semantics of a TABLE.

saying these in an interview costs you the question

  • Saying a TABLE keeps all historical records (it keeps only the latest per key)
  • Claiming a topic can only be one or the other (the same topic can be read as stream or table)
  • Thinking a null value deletes from a STREAM

context

open as a page

What do CREATE STREAM AS SELECT (CSAS) and CREATE TABLE AS SELECT (CTAS) do, and what runs underneath when you execute one?

level: middleimportance: must knowfreq 70%

basics

~20 s

CSAS and CTAS create a new, continuously-updated derived stream or table from a SELECT over existing ones. Each launches a persistent query — a long-running Kafka Streams job — that writes results into a new backing Kafka topic.

open as a page

Explain the difference between a push query (EMIT CHANGES) and a pull query in ksqlDB.

level: middleimportance: must knowfreq 72%

basics

~20 s

A push query (EMIT CHANGES) subscribes to a continuous stream of updates and never finishes until cancelled — it pushes new results as they arrive. A pull query is a point-in-time lookup against a materialized table's current state and returns immediately.

open as a page

What are the join types and co-partitioning requirements in ksqlDB, and how do stream-stream, stream-table, and table-table joins differ?

level: seniorimportance: should knowfreq 50%

basics

~20 s

Stream-stream joins are windowed (you join events within a time window) and produce a stream. Stream-table joins are non-windowed enrichment lookups against the table's current state. Table-table joins maintain a continuously-updated joined table. All require co-partitioned inputs (same key, same partition count).

open as a page

How does ksqlDB materialize state for aggregations and pull queries, and what role do RocksDB state stores and changelog topics play in fault tolerance?

level: seniorimportance: should knowfreq 48%

basics

~20 s

Aggregations build a materialized view stored in a local RocksDB state store on disk, partitioned across ksqlDB nodes. Each store is backed by a compacted Kafka changelog topic, so on failure or restart the state is rebuilt by replaying that changelog.

open as a page

How do windowed aggregations work in ksqlDB, and what are the differences between tumbling, hopping, and session windows?

level: seniorimportance: should knowfreq 55%

basics

~10 s

Windowed aggregations group events into time buckets before aggregating. Tumbling windows are fixed-size, non-overlapping; hopping windows are fixed-size but overlap by a smaller advance; session windows are activity-based, closing after a gap of inactivity.

open as a page