Перейти к содержанию
Learning Platform
Глоссарий Troubleshooting
Урок 12.02 · 35 мин
Продвинутый
KafkaingestionMVat-least-once

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.

Kafka Engine: трёхшаговый паттерн потребления
Kafka topic (events)Kafka topic: источник данных. Kafka engine table является consumer в consumer group. kafka_group_name определяет имя consumer group — ClickHouse хранит offsets в Kafka (не в Keeper). При перезапуске потребление продолжается с последнего закоммиченного offset.
pull: Kafka Engine читает сообщения
Kafka Engine table (events_queue)Kafka Engine table (events_queue): виртуальная таблица-очередь. SELECT из неё читает сообщения из Kafka и коммитит offset. Данные нигде не хранятся локально — Kafka Engine это только cursor. Прямые SELECT из Kafka table не рекомендуются в production — используйте MV.
MATERIALIZED VIEW автоматически
MergeTree target (events_local)MergeTree target table (events_local): физическое хранилище данных. MV выбирает данные из Kafka table и INSERT-ит в target. Offset в Kafka коммитится только после успешного INSERT в MergeTree — это обеспечивает at-least-once семантику.

Шаг 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 доставку: каждое сообщение будет обработано хотя бы один раз, но возможны дубли при определённых сценариях сбоев.

Механизм:

  1. ClickHouse читает batch сообщений из Kafka
  2. MV выполняет INSERT в MergeTree target
  3. Только после успешного INSERT offset коммитится в Kafka

Если ClickHouse падает после INSERT, но до коммита offset — сообщения будут прочитаны повторно при следующем старте (дубликаты). Для защиты от дублей используйте ReplacingMergeTree или дедупликацию через insert_deduplication_token.

WARNING

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.

INFO

В 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 v1Kafka Engine v2
Хранение offsetsВ Kafka (__consumer_offsets)В ClickHouse Keeper
СемантикаAt-least-once (возможны дубли)Exactly-once (с insert dedup)
ЗависимостиТолько Kafka brokersKafka brokers + Keeper
ПерезапускВозобновление с last Kafka commitВозобновление с Keeper-state
СовместимостьВсе версии25.x+ (experimental в 25.x, beta-stable в 26.x)
Использование с Replicated tablesВозможно, но без атомарности offsetАтомарный commit через Keeper
WARNING

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;
TIP

Для практики настройки Kafka Engine используйте лабораторную работу: labs/kafka/. Лаб включает apache/kafka:4.0.0 (KRaft) + clickhouse-server:26.3 с упражнениями по созданию Kafka Engine таблицы и Materialized View.


Ключевые выводы

  1. Kafka Table Engine реализует pull-модель: ClickHouse сам является Kafka consumer, не требуя отдельного коннектора.
  2. Трёхшаговый паттерн: target MergeTree + Kafka Engine table (cursor) + Materialized View (auto-transfer).
  3. At-least-once семантика: offset коммитится только после успешного INSERT в MergeTree. При сбое — возможны дубли.
  4. kafka_num_consumers должен быть не больше числа партиций — иначе лишние consumers простаивают.
  5. Для 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.
Streaming: event time, watermarks, windows Apache Kafka: основы и архитектура

Закончили урок?

Отметьте его как пройденный, чтобы отслеживать свой прогресс

Войдите чтобы оценить урок

Прогресс модуля
0 из 11