skip to content

How do you ingest a Kafka topic into ClickHouse using the Kafka table engine?

level: seniorimportance: should knowfreq 55%

answer

  1. the topic-facing table stores nothing
  2. querying it by hand steals messages
  3. a third object pumps rows into storage
  4. block size still decides part size
  5. offsets commit after the write, so expect replays

basics

~20 s

Create three objects: a Kafka engine table that consumes the topic, a MergeTree table that stores the data, and a materialized view that reads from the first and inserts into the second. The Kafka table is a consumer, not storage — selecting from it consumes messages.

solid answer

~40 s

The Kafka table engine turns a topic into a streaming source. It stores nothing: a `SELECT` against it pulls messages and advances the consumer, so it is not a table you query. The standard pattern is three objects — the `Kafka` engine table declaring `kafka_broker_list`, `kafka_topic_list`, `kafka_group_name` and `kafka_format`; a MergeTree table with the real schema and sorting key; and a `MATERIALIZED VIEW ... TO` the MergeTree that fires on each block the consumer pulls. Batch size matters as much as anywhere else in ClickHouse: `kafka_max_block_size` and the flush interval decide how many rows land in each part, so a topic with low throughput and a short flush interval still produces part pressure. Parallelism comes from `kafka_num_consumers`, capped by the topic's partition count. Delivery is at-least-once, so design the target for duplicates.

code

sql · 22 lines
sql
CREATE TABLE events_queue (raw String)
ENGINE = Kafka
SETTINGS kafka_broker_list   = 'broker:9092',
         kafka_topic_list    = 'events',
         kafka_group_name    = 'clickhouse_events',
         kafka_format        = 'JSONAsString',
         kafka_num_consumers = 4,
         kafka_handle_error_mode = 'stream';

CREATE MATERIALIZED VIEW events_mv TO events AS
SELECT
    JSONExtract(raw, 'event_time', 'DateTime') AS event_time,
    JSONExtract(raw, 'user_id', 'UInt64')      AS user_id,
    _partition                                  AS kafka_partition,
    _offset                                     AS kafka_offset
FROM events_queue
WHERE length(_error) = 0;

CREATE MATERIALIZED VIEW events_errors_mv TO events_dlq AS
SELECT _raw_message AS message, _error AS error, now() AS seen_at
FROM events_queue
WHERE length(_error) > 0;

go deeper

for a junior

Recall the three-object shape: a Kafka engine table, a MergeTree target, and a materialized view that moves rows between them. Know the Kafka table is a consumer, not storage.

for a middle

Explain the materialized view as an insert trigger, how block size and flush interval determine part size, and why consumer count is bounded by the topic's partitions.

for a senior

Demonstrate operating it: at-least-once duplicates and how the target absorbs them, dead-lettering malformed messages via the error virtual columns, and diagnosing a stalled consumer from system.kafka_consumers.

for a principal

Own the build-versus-buy call — built-in engine against an external consumer or managed connector — weighing transformation needs, schema evolution, independent scaling of ingest, and who is on call when the consumer group stalls.

## The three-object pattern ClickHouse consumes Kafka with a pipeline of three schema objects, and getting the roles right is most of the interview answer. ```sql -- 1. the consumer: stores nothing CREATE TABLE events_queue ( event_time DateTime, user_id UInt64, action String ) ENGINE = Kafka SETTINGS kafka_broker_list = 'broker:9092', kafka_topic_list = 'events', kafka_group_name = 'clickhouse_events', kafka_format = 'JSONEachRow', kafka_num_consumers = 4; -- 2. the storage CREATE TABLE events ( event_time DateTime, user_id UInt64, action LowCardinality(String) ) ENGINE = MergeTree PARTITION BY toYYYYMM(event_time) ORDER BY (user_id, event_time); -- 3. the pump CREATE MATERIALIZED VIEW events_mv TO events AS SELECT event_time, user_id, action FROM events_queue; ``` **The Kafka engine table holds no data.** It is a cursor over a topic. Running `SELECT * FROM events_queue` by hand pulls messages and moves the consumer forward, which means an innocent debugging query silently steals rows from the pipeline. This is the single most common mistake and the thing to say out loud in an interview. **The materialized view is an insert trigger,** not a cached result. Each block the consumer reads is passed through the view's `SELECT` and inserted into the target table. That is where you do light transformation: type casting, flattening, dropping fields, deriving columns. ## Batching, and why it is still the whole game Every block the view pushes becomes a data part, so the Kafka path is subject to exactly the same part-explosion pressure as any other write path. Two settings govern the block size: `kafka_max_block_size` (rows collected before a flush) and the Kafka flush interval (time before an incomplete block is flushed anyway). A low-volume topic with an aggressive flush interval will happily produce a part every few hundred milliseconds. Size these so each part carries a meaningful number of rows. ## Parallelism `kafka_num_consumers` sets how many consumers the table runs; it cannot usefully exceed the topic's partition count, since each partition is consumed by at most one consumer in the group. `kafka_thread_per_consumer` gives each consumer its own thread so they do not share one polling thread. Across a replicated ClickHouse cluster you typically create the Kafka table on several nodes with the *same* consumer group name and let the group split partitions between them, with each node writing into its local shard or replica. ## Delivery semantics The pipeline is **at-least-once**. Offsets advance after a block has been written; a failure between the write and the commit, or a consumer rebalance mid-flight, can replay messages that were already stored. Plan for it rather than hoping: - store into a `ReplacingMergeTree` keyed on a business key so duplicates collapse on merge (with `FINAL` or aggregation at read time until they do), - or deduplicate downstream in the query, - or make the payload naturally idempotent. Do not claim exactly-once. An interviewer who works with this stack will notice. ## Malformed messages One unparseable message can otherwise stall the whole consumer. `kafka_skip_broken_messages` allows a bounded number of parse failures per block to be skipped. `kafka_handle_error_mode = 'stream'` instead exposes failures through the `_error` and `_raw_message` virtual columns, so a second materialized view can route them into a dead-letter table while the main view keeps only clean rows. Virtual columns also give you provenance — `_topic`, `_partition`, `_offset`, `_key`, `_timestamp` — and storing `_partition` and `_offset` in the target table makes duplicate hunting and replay auditing straightforward. ## Operating it `system.kafka_consumers` shows consumer state, assignments and last exceptions — the first place to look when rows stop arriving. Server logs carry rebalance and parse errors. `DETACH TABLE` / `ATTACH TABLE` on the materialized view is the way to pause and resume ingestion (for example while doing maintenance on the target) without dropping the consumer group. ## Scope note: Kinesis There is no Kinesis table engine in the open-source ClickHouse server the way there is for Kafka. Ingesting from Kinesis means either ClickPipes on ClickHouse Cloud, which offers managed connectors for streaming sources, or an external process that reads the stream and inserts in batches. Saying "same as Kafka but with a Kinesis engine" is a factual error. ## When not to use the engine at all The built-in consumer is attractive because it needs no extra service, but it couples ingestion lifecycle to database DDL, offers limited transformation, and gives you at-least-once with manual error routing. Many teams prefer an external consumer or a managed connector when they need richer transformation, schema-registry integration, or independent scaling of the ingest tier. The engine shines for straightforward, high-volume, low-transformation feeds.

  • Why should you never run an ad-hoc SELECT against a Kafka engine table?
    Because reading from it consumes messages and advances the consumer group's offsets. The rows you see in your debugging query are rows the materialized view will never insert into the target table, so a casual `SELECT * ... LIMIT 10` silently drops data. Inspect the target MergeTree table instead, or use a dedicated consumer group for exploration.
  • How do you handle a malformed message that would otherwise stall the consumer?
    Either allow a bounded number of parse failures per block with `kafka_skip_broken_messages`, or set `kafka_handle_error_mode = 'stream'` and read the `_error` and `_raw_message` virtual columns. The second is preferable in production: one materialized view routes clean rows to the target table and a second routes failures to a dead-letter table, so nothing is lost silently.
  • What delivery guarantee does this pipeline give, and how do you cope with it?
    At-least-once. Offsets advance after the block is written, so a crash or rebalance in between replays messages. Cope by storing into a ReplacingMergeTree keyed on a business identifier so duplicates collapse on merge, by deduplicating at query time, or by making the payload idempotent. Persisting the `_partition` and `_offset` virtual columns makes duplicates auditable.
  • How many consumers should the Kafka engine table run?
    At most the topic's partition count, since each partition is served by one consumer in the group. Set `kafka_num_consumers` accordingly and enable `kafka_thread_per_consumer` so they do not share a single polling thread. Across a cluster, create the table on several nodes with the same consumer group name and let the group divide partitions between them.

saying these in an interview costs you the question

  • Thinks the Kafka engine table stores the consumed messages
  • Runs SELECT on the Kafka table to check ingestion is working
  • Claims the pipeline delivers exactly-once semantics
  • Sets kafka_num_consumers far above the topic's partition count
  • Says ClickHouse has a Kinesis table engine like the Kafka one

context