Kafka Topics from Flink

Reading and Writing Kafka Topics from Flink

Flink 129 's Kafka 129 connector lives outside the core, versioned for a Flink release: here flink-sql-connector-kafka 5.0.0-2.2, built for Flink 2.2 and loaded unchanged by 2.3.0 from /opt/flink/lib. In SQL, a Kafka topic becomes a table: the kafka connector reads or appends records in a format (JSON, Avro 129 with a schema registry, and others), and upsert-kafka treats a compacted topic as a table keyed by its primary key. The source tracks offsets in Flink's checkpoints, and the sink can write in Kafka transactions ('sink.delivery-guarantee' = 'exactly-once'). BookNest's job, run through the SQL client:

compose/flink/booknest.sql: daily revenue per channel, from Kafka to KafkaSQL
-- booknest.sql: Flink SQL over booknest.orders: revenue per channel and day, back into Kafka
-- (the SQL client splits statements at a semicolon that ends a line, so no comments after one)
SET 'table.dml-sync' = 'true';
CREATE TABLE orders (
  order_id BIGINT, order_ts STRING, channel STRING, status STRING,
  items ARRAY<ROW<book_id INT, qty INT, unit_price DOUBLE>>,
  ts AS TO_TIMESTAMP(order_ts, 'yyyy-MM-dd''T''HH:mm:ss''Z'''),   -- event time
  WATERMARK FOR ts AS ts - INTERVAL '1' HOUR
) WITH ('connector' = 'kafka', 'topic' = 'booknest.orders', 'format' = 'json',
  'properties.bootstrap.servers' = 'l3-kafka:9092',
  'scan.startup.mode' = 'earliest-offset', 'scan.bounded.mode' = 'latest-offset');
CREATE TABLE daily_revenue (day_start TIMESTAMP(3), channel STRING, lines BIGINT,
  revenue DOUBLE) WITH ('connector' = 'kafka', 'topic' = 'booknest.daily-revenue',
  'properties.bootstrap.servers' = 'l3-kafka:9092', 'format' = 'json');
CREATE VIEW order_lines AS
  SELECT o.ts, o.channel, i.qty * i.unit_price AS revenue
  FROM orders AS o CROSS JOIN UNNEST(o.items) AS i (book_id, qty, unit_price)
  WHERE o.status <> 'cancelled';
INSERT INTO daily_revenue
  SELECT window_start, channel, COUNT(*), ROUND(SUM(revenue), 2)
  FROM TABLE(TUMBLE(TABLE order_lines, DESCRIPTOR(ts), INTERVAL '1' DAY))
  GROUP BY window_start, window_end, channel;

ts is a computed column parsed from order_ts, and the WATERMARK clause makes it the event-time attribute, allowing records up to an hour late. scan.bounded.mode stops the source at today's end offsets, so this streaming job finishes once it has read the 100,000 orders; drop it and the job runs forever, emitting each day as the watermark passes it.