What do CREATE STREAM AS SELECT (CSAS) and CREATE TABLE AS SELECT (CTAS) do, and what runs underneath when you execute one?
answer
- CSAS → derived STREAM
- CTAS → derived TABLE (needs GROUP BY/agg)
- Both = persistent query (long-running)
- Compiles to Kafka Streams topology
- New sink topic + own consumer group
basics
~20 sCSAS 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.
solid answer
~50 s`CREATE STREAM AS SELECT` (CSAS) and `CREATE TABLE AS SELECT` (CTAS) define **derived collections** from a query over existing streams/tables. Unlike a transient `SELECT … EMIT CHANGES`, they register a **persistent query**: a continuously-running Kafka Streams topology that ksqlDB compiles and runs on the server, surviving restarts. The query reads input topics, applies the transformation (filter, project, join, aggregate), and writes results to a **new sink Kafka topic** that backs the derived object. CSAS produces a STREAM (use for stateless transforms — filter/map/repartition, or stream-stream windowed joins). CTAS produces a TABLE and is required when the SELECT contains a **GROUP BY aggregation** or a table-table join, because the result is evolving keyed state. The persistent query keeps running until you `TERMINATE` it; `DROP` removes the object. Each persistent query has its own consumer group, state stores, and changelog/repartition topics.
go deeper
Know CSAS makes a derived stream, CTAS makes a derived table, and both keep running.
Explain persistent vs transient queries, when GROUP BY forces CTAS, and that a new sink topic is created.
Describe the compilation to a Kafka Streams topology, consumer groups, repartition/changelog topics, and lifecycle (TERMINATE/DROP).
Reason about resource cost per persistent query, schema evolution, partition-count planning, and migration/ALTER limitations at scale.
## What these statements are ksqlDB queries come in two flavors: - **Transient** queries (`SELECT … EMIT CHANGES;`) stream results back to *your client* and stop when you disconnect. - **Persistent** queries (`CREATE STREAM/TABLE … AS SELECT …`) run **on the ksqlDB server forever**, writing their output into Kafka. CSAS/CTAS are how you build persistent queries. ### CSAS — CREATE STREAM AS SELECT ```sql CREATE STREAM big_orders WITH (KAFKA_TOPIC='big_orders') AS SELECT order_id, amount FROM orders WHERE amount > 1000 EMIT CHANGES; ``` This registers a new STREAM `big_orders`, backed by a brand-new Kafka topic. A **persistent query** (named like `CSAS_BIG_ORDERS_5`) runs continuously: it consumes `orders`, filters, and produces matching rows to the new topic. Use CSAS for **stateless** transforms (filter, project, mask, rekey via `PARTITION BY`) and for **stream-stream** (windowed) joins whose output is itself a stream of events. ### CTAS — CREATE TABLE AS SELECT ```sql CREATE TABLE orders_per_user WITH (KAFKA_TOPIC='orders_per_user') AS SELECT user_id, COUNT(*) AS order_count FROM orders GROUP BY user_id EMIT CHANGES; ``` When the SELECT contains a **GROUP BY** (aggregation) or a table-table join, the result is *evolving per-key state*, so it must be a **TABLE** — you use CTAS. The persistent query maintains the running aggregate in a **state store** and emits a changelog of updates to the sink topic. ### What actually runs underneath 1. ksqlDB **compiles** the SQL into a **Kafka Streams topology** (a DAG of processors). 2. It runs that topology inside the ksqlDB server JVM as a **persistent query** with its own **consumer group**, so work is parallelized across partitions and across ksqlDB nodes in a cluster. 3. Stateful steps (aggregations, joins, windows) keep state in **RocksDB-backed state stores**, each fault-tolerant via a **changelog topic** in Kafka. 4. Operations that need to re-key data (e.g., GROUP BY on a non-key column, or a join on a non-key) insert an internal **repartition topic**. 5. The derived object's **sink topic** holds the actual output — other queries (or external consumers) can read it. ### Lifecycle and operational notes - A persistent query keeps running across server restarts; you stop it with `TERMINATE <queryId>;`. - `DROP STREAM/TABLE … [DELETE TOPIC];` removes the registered object (and optionally its backing topic). - You can declare the output topic's partitions/replicas/format in the `WITH(...)` clause; by default it inherits the source's partition count. - Schema is **inferred** from the SELECT projection. ### Edge cases - You **cannot** use CSAS when the query aggregates with GROUP BY — ksqlDB rejects it and tells you to use CTAS. - Renaming/altering a persistent query usually means terminating + recreating; in-place ALTER is limited. - Each persistent query consumes broker resources (topics, consumer-group slots, state-store disk) — they are not free.
- When are you forced to use CTAS instead of CSAS?When the SELECT contains a GROUP BY aggregation or a table-table join — the result is evolving keyed state, which must be a TABLE.
- How do you stop a persistent query created by CSAS/CTAS?Run TERMINATE <queryId>; (find the id with SHOW QUERIES), then optionally DROP the stream/table, optionally with DELETE TOPIC to also remove the backing topic.
saying these in an interview costs you the question
- Saying CSAS/CTAS only run while your client is connected (that's a transient query, not persistent)
- Claiming they query data in place without creating a new topic
- Using CSAS for a GROUP BY aggregation