Kafka Table Engine
ClickHouse реализует уникальную pull-модель интеграции с Kafka: сам ClickHouse становится Kafka consumer. Нет отдельного процесса-коннектора, нет sidecar-контейнера — ClickHouse читает данные из Kafka напрямую через Kafka Table Engine.
Pull-модель против push-модели
В push-модели (Kafka Connect) внешний коннектор читает из Kafka и отправляет данные в ClickHouse. В pull-модели (Kafka Engine) ClickHouse сам является consumer в consumer group — он периодически опрашивает Kafka broker и читает новые сообщения.
Трёхшаговый паттерн
Стандартный паттерн требует трёх объектов: target-таблица, Kafka engine таблица, Materialized View.
Шаг 1: Target MergeTree таблица
-- Физическое хранилище — обычная MergeTree таблица
CREATE TABLE events_local
(
event_id String,
event_type LowCardinality(String),
user_id UInt64,
ts DateTime
)
ENGINE = MergeTree
ORDER BY (event_type, user_id, ts);
Шаг 2: Kafka Engine таблица (consumer)
-- Kafka Engine таблица — cursor в Kafka topic
CREATE TABLE events_queue
(
event_id String,
event_type String,
user_id UInt64,
ts DateTime
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'events',
kafka_group_name = 'clickhouse_consumer',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 1;
Шаг 3: Materialized View — связующее звено
-- MV автоматически переносит данные из Kafka в MergeTree
CREATE MATERIALIZED VIEW events_mv TO events_local AS
SELECT event_id, event_type, user_id, ts
FROM events_queue;
После создания MV ClickHouse начинает потреблять сообщения из Kafka автоматически в фоновом режиме.
At-least-once семантика
Kafka Engine гарантирует at-least-once доставку: каждое сообщение будет обработано хотя бы один раз, но возможны дубли при определённых сценариях сбоев.
Механизм:
- ClickHouse читает batch сообщений из Kafka
- MV выполняет INSERT в MergeTree target
- Только после успешного INSERT offset коммитится в Kafka
Если ClickHouse падает после INSERT, но до коммита offset — сообщения будут прочитаны повторно при следующем старте (дубликаты). Для защиты от дублей используйте ReplacingMergeTree или дедупликацию через insert_deduplication_token.
Kafka Engine не гарантирует exactly-once семантику. При сбое между INSERT и коммитом offset возможны дубли. Для compliance-сценариев, требующих exactly-once, используйте Kafka Connect sink (урок 03).
Ключевые настройки
-- Несколько consumers для параллельного чтения из партиций
-- kafka_num_consumers должен быть <= числа партиций топика
CREATE TABLE events_queue (...)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'events',
kafka_group_name = 'clickhouse_consumer',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 3, -- 3 consumer threads для 3 партиций
kafka_max_block_size = 65536, -- максимум строк в одном batch
kafka_poll_timeout_ms = 500; -- таймаут ожидания новых сообщений
Правило масштабирования: kafka_num_consumers должен быть не больше числа партиций топика. Дополнительные consumers сверх числа партиций будут простаивать.
Kafka Engine v2: Keeper-based offset storage
В классической v1-имплементации (всё, что было до 25.x) ClickHouse полагался на сам Kafka для коммита offset. Это создаёт фундаментальную проблему two-phase commit: вначале INSERT в MergeTree, затем __consumer_offsets в Kafka — между этими двумя коммитами возможен сбой, и при перезапуске сообщения читаются повторно (at-least-once с дублями).
Начиная с ClickHouse 25.x, доступна Kafka Engine v2 — переработанная имплементация, которая хранит consumer offsets в ClickHouse Keeper. ClickHouse сам становится единственным владельцем состояния потребления, а offset коммитится атомарно вместе с данными — через тот же commit protocol, что и для replicated parts.
В 26.3 LTS Kafka Engine v2 поддерживает Keeper-based offset storage, который делает possible exactly-once consumer pattern на уровне самого Kafka Engine — без необходимости использовать внешний Kafka Connect. Включается флагом allow_experimental_kafka_offsets_storage_in_keeper = 1 и параметрами таблицы kafka_keeper_path + kafka_replica_name.
-- Kafka Engine v2 с offset в Keeper
SET allow_experimental_kafka_offsets_storage_in_keeper = 1;
CREATE TABLE events_queue_v2
(
event_id String,
event_type String,
user_id UInt64,
ts DateTime
)
ENGINE = Kafka
SETTINGS
kafka_broker_list = 'kafka:9092',
kafka_topic_list = 'events',
kafka_group_name = 'clickhouse_v2',
kafka_format = 'JSONEachRow',
kafka_num_consumers = 3,
-- v2-специфичные параметры
kafka_keeper_path = '/clickhouse/kafka/events_queue_v2',
kafka_replica_name = '{replica}';
Exactly-once consumer pattern
Combine Kafka Engine v2 с insert deduplication on the target ReplicatedMergeTree table — offset в Keeper и data block в storage коммитятся в рамках одной Keeper-транзакции, что эффективно даёт exactly-once семантику:
-- Target — ReplicatedMergeTree с insert deduplication (default)
CREATE TABLE events_local ON CLUSTER my_cluster
(
event_id String,
event_type LowCardinality(String),
user_id UInt64,
ts DateTime
) ENGINE = ReplicatedMergeTree('/clickhouse/tables/{shard}/events', '{replica}')
ORDER BY (event_type, user_id, ts);
-- MV переносит из Kafka v2 → ReplicatedMergeTree
CREATE MATERIALIZED VIEW events_mv ON CLUSTER my_cluster
TO events_local AS
SELECT event_id, event_type, user_id, ts
FROM events_queue_v2;
При сбое после INSERT, но до commit offset в Keeper, повторное чтение того же batch триггерит insert deduplication по hash блока — запись дедуплицируется, а offset двигается вперёд. Никаких дублей не появляется.
Сравнение v1 и v2
| Критерий | Kafka Engine v1 | Kafka Engine v2 |
|---|---|---|
| Хранение offsets | В Kafka (__consumer_offsets) | В ClickHouse Keeper |
| Семантика | At-least-once (возможны дубли) | Exactly-once (с insert dedup) |
| Зависимости | Только Kafka brokers | Kafka brokers + Keeper |
| Перезапуск | Возобновление с last Kafka commit | Возобновление с Keeper-state |
| Совместимость | Все версии | 25.x+ (experimental в 25.x, beta-stable в 26.x) |
| Использование с Replicated tables | Возможно, но без атомарности offset | Атомарный commit через Keeper |
v2 не отменяет v1 — обе имплементации сосуществуют. Для миграции существующего pipeline с v1 на v2: создайте новую events_queue_v2 с другим kafka_group_name (чтобы не пересекаться с v1-consumer), переведите MV на новую таблицу, удалите старую. Не пытайтесь добавить kafka_keeper_path к существующей v1-таблице — offset state в Keeper не будет инициализирован из Kafka.
Когда v2 vs Kafka Connect
Раньше для exactly-once в production применяли Kafka Connect sink с KeeperMap-движком (урок 03). С появлением v2 многие сценарии можно покрыть прямо через Kafka Engine — без отдельного Connect-кластера. Но Kafka Connect остаётся актуальным для случаев, когда требуется стандартный Kafka Connect ecosystem (трансформации SMT, dead-letter queue, schema registry на уровне connector framework).
Мониторинг
-- Состояние Kafka consumers (26.3 LTS)
SELECT *
FROM system.kafka_consumers
WHERE database = 'default';
-- Количество прочитанных сообщений
SELECT
topic,
partition,
offset,
rows_read
FROM system.kafka_consumers_info
ORDER BY topic, partition;
Для практики настройки Kafka Engine используйте лабораторную работу: labs/kafka/. Лаб включает apache/kafka:4.0.0 (KRaft) + clickhouse-server:26.3 с упражнениями по созданию Kafka Engine таблицы и Materialized View.
Ключевые выводы
- Kafka Table Engine реализует pull-модель: ClickHouse сам является Kafka consumer, не требуя отдельного коннектора.
- Трёхшаговый паттерн: target MergeTree + Kafka Engine table (cursor) + Materialized View (auto-transfer).
- At-least-once семантика: offset коммитится только после успешного INSERT в MergeTree. При сбое — возможны дубли.
kafka_num_consumersдолжен быть не больше числа партиций — иначе лишние consumers простаивают.- Для exactly-once доставки используйте Kafka Engine v2 (
kafka_keeper_path+kafka_replica_name, флагallow_experimental_kafka_offsets_storage_in_keeper) с ReplicatedMergeTree-таргетом и insert deduplication. Kafka Connect sink (урок 03) — альтернатива для сценариев, где нужен стандартный Kafka Connect ecosystem.