ClickHouse streaming has one landing pattern and four kinds of source. The pattern is a stream consumed in batches, written as parts, and fanned out through materialised views to the tables that serve queries; the sources are a Kafka topic of application events, a change-data-capture feed from a transactional database, a stream-processing layer that reshapes events before they arrive, and a direct producer that writes with async inserts. Each source changes what arrives and how; none of them changes what ClickHouse does with it.
This page sets out the landing pattern once, then walks the four sources with the topology, the delivery guarantee each gives, the failure it is prone to, and the archive post that builds it end to end. It closes with the arithmetic of consumer lag and batch size, the deduplication that turns at-least-once into effectively-once, and the monitoring that catches a stalled stream before a dashboard does. The Kafka hub covers the Kafka engine in depth; this page covers streaming as a whole.
The archive holds the Kafka-with-ClickHouse introduction, real-time stream processing with Kafka, Kafka streaming for fintech, PostgreSQL-to-ClickHouse CDC with Debezium, MySQL-to-ClickHouse with Redpanda and ksqlDB, the Kafka producer memory-leak monitor, and the payments analytics reference design.
The landing pattern every ClickHouse streaming source shares
A stream is consumed in blocks of thousands of rows, each block becomes one part in a raw landing table, and materialised views attached to that table fan every block out to the serving tables (rollups, filtered subsets, enriched copies) in the same transaction. The raw table keeps the events as they arrived; the serving tables are what dashboards read.
Batch size decides part count, part count decides merge load, and merge load decides whether the cluster keeps up, so the first number in any ClickHouse streaming design is rows per block. The how to use Kafka with ClickHouse post introduces the pattern with the Kafka engine as the consumer.
-- the landing pattern: consumer → raw table → views → serving tables
CREATE TABLE events_raw
(
ts DateTime64(3),
tenant_id UInt32,
event_type LowCardinality(String),
user_id UInt64,
amount Decimal(18, 2),
payload String,
_source LowCardinality(String), -- which stream delivered it
_offset UInt64 -- position in that stream, for dedup and audit
)
ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/events_raw', '{replica}')
PARTITION BY toYYYYMM(ts)
ORDER BY (tenant_id, event_type, ts);
CREATE MATERIALIZED VIEW events_1m_mv TO events_1m AS
SELECT toStartOfMinute(ts) AS minute, tenant_id, event_type, countState() AS events, sumState(amount) AS revenue
FROM events_raw GROUP BY minute, tenant_id, event_type;Source 1: a Kafka topic of application events, the usual ClickHouse streaming feed
The most common source: producers publish events to a topic, and ClickHouse consumes it with the Kafka table engine (a consumer inside the server, feeding a view) or with an external consumer such as Vector or a custom service that batches and inserts. The engine path is simplest and keeps the consumer group offsets in Kafka; the external path gives more control over batching and back-pressure.
Delivery is at-least-once either way: an insert that succeeded but whose offset commit failed is redelivered. The real-time stream processing with Kafka and ClickHouse post builds the engine path; the Kafka with ClickHouse for fintech post applies it where the events are payments and the SLO is seconds.
CREATE TABLE events_kafka
(
ts DateTime64(3), tenant_id UInt32, event_type String, user_id UInt64, amount Decimal(18, 2), payload String
)
ENGINE = Kafka
SETTINGS kafka_broker_list = '${KAFKA_BROKERS}',
kafka_topic_list = 'events',
kafka_group_name = 'ch_events_v1',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 4,
kafka_max_block_size = 262144,
kafka_poll_timeout_ms = 500;
CREATE MATERIALIZED VIEW events_kafka_mv TO events_raw AS
SELECT ts, tenant_id, toLowCardinality(event_type) AS event_type, user_id, amount, payload,
'kafka:events' AS _source, _offset
FROM events_kafka;
-- consumer state, per partition
SELECT database, table, consumer_id, assignments.topic, assignments.partition_id, assignments.current_offset, num_messages_read, last_poll_time, exceptions.text
FROM system.kafka_consumers ARRAY JOIN assignments, exceptions;Source 2: change data capture from PostgreSQL or MySQL into ClickHouse streaming tables
CDC turns a transactional database’s log into a stream of row changes, and ClickHouse streaming consumes that stream into ReplacingMergeTree tables keyed by the source primary key with a version column, so that inserts, updates and deletes all become versioned rows and the merge keeps the latest.
Debezium reads the PostgreSQL WAL or MySQL binlog into Kafka; a consumer (the Kafka engine, or the sink connector) writes the rows. Delivery is at-least-once and the version column makes redelivery harmless. The streaming from PostgreSQL to ClickHouse with Kafka and Debezium post builds the pipeline with every configuration file; the replication hub covers the MySQL sink-connector variant.
-- the CDC target: versioned by the source's commit position, deletes as tombstones
CREATE TABLE orders_current
(
order_id UInt64,
status LowCardinality(String),
amount Decimal(18, 2),
updated_at DateTime64(3),
_version UInt64, -- Debezium source.lsn / source.pos
_deleted UInt8
)
ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/orders_current', '{replica}', _version, _deleted)
ORDER BY order_id;
-- the view that unpacks a Debezium envelope from the Kafka engine table
CREATE MATERIALIZED VIEW orders_cdc_mv TO orders_current AS
SELECT
JSONExtractUInt(after, 'order_id') AS order_id,
JSONExtractString(after, 'status') AS status,
toDecimal64(JSONExtractFloat(after, 'amount'), 2) AS amount,
now64(3) AS updated_at,
JSONExtractUInt(source, 'lsn') AS _version,
op = 'd' AS _deleted
FROM orders_debezium_kafka;Source 3: a stream-processing layer in front
When events need joining, filtering or reshaping before they land, a stream processor (ksqlDB, Flink, Kafka Streams) does it between the source and ClickHouse, and ClickHouse receives the already-shaped stream. It costs an extra system and moves logic out of SQL that ClickHouse could often run itself in a materialised view; it earns its place when the reshaping needs state across topics, windowed joins between streams, or a schema that several consumers share. The streaming from MySQL to ClickHouse using Redpanda and ksqlDB post builds the layered pipeline and shows what ksqlDB adds over a bare CDC feed.
-- what the layer usually does upstream, expressed as the ClickHouse view that could do it instead
CREATE MATERIALIZED VIEW purchases_enriched_mv TO purchases_enriched AS
SELECT
e.ts, e.tenant_id, e.user_id, e.amount,
dictGet('product_dim', 'category', JSONExtractUInt(e.payload, 'product_id')) AS category
FROM events_raw AS e
WHERE e.event_type = 'purchase';
-- keep the external layer when the join is stream-to-stream with state, not stream-to-dimensionSource 4: direct producers with async inserts
In ClickHouse streaming terms, services that already batch can skip the broker and write straight to ClickHouse over HTTP or native protocol, with async_insert turning many small requests into server-side batches. It removes a system from the path and is the right choice for internal telemetry and moderate volumes; it gives up the broker’s replay buffer, so a ClickHouse outage is back-pressure on the producer rather than a queue that drains later.
The real-time payments analytics post shows the choice made per stream: broker for the payment events that must never be lost, direct inserts for the operational telemetry around them.
-- producer side: many small POSTs become server-side batches
SET async_insert = 1, wait_for_async_insert = 1,
async_insert_busy_timeout_ms = 500, async_insert_max_data_size = 10000000;
INSERT INTO events_raw FORMAT JSONEachRow
{"ts":"2026-09-19 12:00:00.123","tenant_id":42,"event_type":"view","user_id":7781234,"amount":0,"payload":"{}","_source":"svc:web","_offset":0}
-- what async inserts are doing
SELECT database, table, format, formatReadableSize(sum(bytes)) AS pending, count() AS entries
FROM system.asynchronous_inserts GROUP BY database, table, format;| Source | Consumer | Delivery | Replay on ClickHouse outage | Prone to |
|---|---|---|---|---|
| 1. Kafka topic | Kafka engine or external consumer | at-least-once | yes, from the committed offset | consumer lag; small blocks → many parts |
| 2. CDC (Debezium) | Kafka engine or sink connector into Replacing | at-least-once, idempotent by version | yes | schema changes at the source; FINAL cost |
| 3. Stream processor | ksqlDB / Flink output topic | as the processor guarantees | yes | logic split across two systems |
| 4. Direct producer | async inserts over HTTP / native | at-least-once if the client retries | no: back-pressure on the producer | producer memory growth on retries |

ClickHouse streaming lag arithmetic: blocks, parts and the merge budget
Consumer lag in ClickHouse streaming is a function of three numbers: the rate the stream produces, the block size the consumer collects before inserting, and the rate merges can absorb the parts those inserts create. A consumer collecting 10,000-row blocks from a 200,000-row-per-second stream inserts twenty parts a second per consumer, which no merge pool sustains; the same stream at 250,000-row blocks is under one part per second and settles.
The kafka_max_block_size and kafka_poll_timeout_ms pair (or the equivalent in an external consumer) is where the trade is made, and the check is parts per partition against the merge budget from the MergeTree hub.
-- inserts per minute and rows per insert, the two sides of the trade
SELECT toStartOfMinute(event_time) AS m, count() AS inserts, round(avg(written_rows)) AS rows_per_insert, sum(written_rows) AS rows
FROM system.query_log
WHERE type = 'QueryFinish' AND query_kind = 'Insert' AND tables = ['default.events_raw'] AND event_time > now() - INTERVAL 1 HOUR
GROUP BY m ORDER BY m;
-- and whether merges are keeping up
SELECT partition, count() AS active_parts FROM system.parts WHERE table = 'events_raw' AND active GROUP BY partition ORDER BY partition DESC LIMIT 3;From at-least-once to effectively-once
Every source above can deliver a block twice, and ClickHouse streaming handles it at two layers. Replicated tables deduplicate identical inserted blocks by checksum for a window of recent blocks (replicated_deduplication_window), which absorbs the common case of a retried insert of the same batch.
For rows that arrive again in a different batch, the table design carries the dedup: ReplacingMergeTree keyed by the event’s own identifier, or the stream offset stored in the row so that a read can filter duplicates. Exactly-once in the strict sense is a producer-side property; the design goal on the ClickHouse side is that a redelivery never changes a query result.
-- layer 1: block-level, automatic on Replicated tables
SELECT name, value FROM system.merge_tree_settings WHERE name IN ('replicated_deduplication_window', 'replicated_deduplication_window_seconds');
-- for non-replicated tables: non_replicated_deduplication_window
-- layer 2: row-level, by design
CREATE TABLE events_dedup
(
event_id UUID, ts DateTime64(3), tenant_id UInt32, event_type LowCardinality(String), user_id UInt64, amount Decimal(18, 2)
)
ENGINE = ReplicatedReplacingMergeTree('/clickhouse/tables/{shard}/events_dedup', '{replica}')
ORDER BY (tenant_id, event_type, ts, event_id);
-- reads over a window: SELECT ... FROM events_dedup FINAL WHERE ... or count(DISTINCT event_id)Monitoring a ClickHouse streaming pipeline
Four numbers, in priority order: freshness (age of the newest row in the landing table), consumer lag (offset behind the topic head, from the broker or system.kafka_consumers), parts per partition on the landing table, and view failures in system.query_views_log. Freshness is the one that catches everything, because every other failure eventually shows as stale data. The producer side has its own failure mode: a client that buffers on retries can grow without bound, and the Python script to monitor Kafka producer memory leaks post is the archive’s watchdog for it. The monitoring hub covers the cluster-level numbers underneath.
-- the streaming dashboard, one row
SELECT
(SELECT dateDiff('second', max(ts), now()) FROM events_raw WHERE ts > now() - INTERVAL 1 DAY) AS freshness_s,
(SELECT max(assignments.current_offset) FROM system.kafka_consumers ARRAY JOIN assignments) AS consumer_offset,
(SELECT max(cnt) FROM (SELECT count() AS cnt FROM system.parts WHERE table = 'events_raw' AND active GROUP BY partition)) AS max_parts_per_partition,
(SELECT countIf(status != 'QueryFinish') FROM system.query_views_log WHERE event_time > now() - INTERVAL 1 HOUR) AS view_failures_1h;Schema evolution on a live ClickHouse streaming pipeline
ClickHouse streaming sources change shape while they run: a producer adds a field, a CDC source gains a column, a processor renames one. The landing table absorbs additive change cheaply (ADD COLUMN with a default is instant, and input_format_skip_unknown_fields lets the consumer keep going when a field it does not know appears); renames and type changes are coordinated: the view is altered first to map old and new names, then the source is changed, then the old mapping is removed.
A raw String payload column beside the typed columns is the safety net that keeps unknown fields queryable until the schema catches up.
-- additive change without stopping the consumer
ALTER TABLE events_raw ON CLUSTER 'ch_prod' ADD COLUMN campaign LowCardinality(String) DEFAULT '';
ALTER TABLE events_kafka MODIFY SETTING input_format_skip_unknown_fields = 1; -- or on the session of the external consumer
-- then the view: DROP + CREATE with the new column mapped; the Kafka engine table buffers offsets meanwhileChoosing between the four sources
The choice is made per stream, not per platform, and three questions settle it. Must the events survive a ClickHouse outage without loss? Then a broker (sources 1 to 3), because it holds the replay buffer; direct inserts are for streams where a gap is acceptable or the producer has its own buffer.
Do the events originate in a transactional database? Then CDC (source 2), because re-implementing the change feed in application code drifts from the database. Does the shaping need state across streams? Then a processor (source 3); otherwise the materialised view does the shaping and the processor is a system too many.
A payments estate typically ends up with all four: broker for the payment events, CDC for the account state, a processor where two streams must be joined by time, and direct inserts for the telemetry that watches the rest.
Back-pressure: what happens when ClickHouse is the slow side
Every source has a behaviour when inserts slow down, and it is decided before it happens. The Kafka engine simply stops polling while an insert is in flight, so lag grows on the broker and nothing is lost; an external consumer does the same if it commits offsets only after a successful insert, and loses data if it commits first. A CDC feed behaves as its broker does.
Direct producers hit TOO_MANY_PARTS or a timeout and must retry with a bounded buffer, which is where the memory-leak watchdog earns its keep. The design rule is that back-pressure lands on a queue with a known depth and a known alert, never on a producer’s heap.
-- the two signals that ClickHouse is the slow side
SELECT value FROM system.metrics WHERE metric = 'DelayedInserts'; -- inserts being throttled by parts_to_delay_insert
SELECT count() FROM system.errors WHERE name = 'TOO_MANY_PARTS' AND last_error_time > now() - INTERVAL 1 HOUR;Version notes
The Kafka engine has been stable since 19.x; system.kafka_consumers is 23.x and later. async_insert has been stable since 22.x. replicated_deduplication_window defaults to 1,000 blocks; confirm the value on the running version. The Debezium connector versions in the archive posts predate current releases; the Kafka engine documentation is the source for engine settings, and the connector’s own documentation for the envelope format. Confirm both on the running versions.
Reading the archive
The archive spans 22.x to 26.8, and the connector and engine versions in the older posts have moved on; the patterns hold, the settings names are checked against the version notes above.
Start with how to use Kafka with ClickHouse for the ClickHouse streaming landing pattern, then real-time stream processing and the fintech post for source 1, the Debezium post for source 2, the Redpanda and ksqlDB post for source 3, and the payments design for the per-stream choice between sources. The producer memory-leak monitor is the one operational post and belongs on every producer host.
ChistaDATA designs the source-to-serving path as a ClickHouse consulting engagement and keeps the four monitoring numbers above as standing alerts under managed services. Every consumer setting and every view on this page is tested on staging at production rate before it goes live, with the landing table backed up and a tested restore in place.