ksqlDB Streams and Tables

ksqlDB Streams, Tables and Continuous Queries

ksqlDB is a server built on Kafka Streams 129 . CREATE STREAM declares a schema over an existing topic; CREATE TABLE ... AS SELECT starts a persistent query, a Kafka Streams application that runs inside the server and maintains a table in a topic and a state store. Push queries (EMIT CHANGES) stream every change to the client; pull queries read a table's current value, like a database lookup.

listings/l06111_ksqldb.sh: a stream, a persistent query, and a pull queryShell
# Start ksqlDB, declare a stream over the order events, and keep a table of counts per type
docker compose -f compose/ksqldb.yaml up -d 2>/dev/null
until curl -sf localhost:33088/info >/dev/null; do sleep 2; done
docker exec -i l3-ksqldb ksql http://localhost:8088 <<'SQL' 2>/dev/null | grep 'v8\|reated'
SET 'auto.offset.reset' = 'earliest';
CREATE STREAM order_events (order_key STRING KEY, event_id BIGINT, ts STRING, type STRING,
  order_id BIGINT, total DOUBLE) WITH (KAFKA_TOPIC = 'booknest.order-events',
  VALUE_FORMAT = 'JSON');
CREATE TABLE events_by_type AS
  SELECT type, COUNT(*) AS events FROM order_events GROUP BY type EMIT CHANGES;
SQL
sleep 40                                             # let the persistent query catch up
H='Content-Type: application/vnd.ksqlapi.delimited.v1'                # a pull query over REST
curl -s -X POST localhost:33088/query-stream -H "$H" \
  -d '{"sql": "SELECT * FROM events_by_type;"}' | tail -n +2 | jq -r @tsv     # skip the header
Output
CLI v8.3.2, Server v8.3.2 located at http://localhost:8088
 Stream created
 Created query with ID CTAS_EVENTS_BY_TYPE_1
order_cancelled   5890
order_delivered   93009
order_placed   100000
order_paid   94109
order_shipped   93803
order_returned   3929

One row per event type. order_paid shows three more than Consuming for Analytics counted because the test events of Dead Letter Queues share the topic; the one that is not JSON went to ksqlDB's processing log instead. Each persistent query is a Kafka Streams application inside the server, with its own internal topics and state stores.