Перейти к основному содержанию
Движок таблицы Kafka можно использовать для чтения данных из и записи данных в Apache Kafka и другие брокеры с поддержкой Kafka API (например, Redpanda, Amazon MSK).

Kafka в ClickHouse

Если вы используете ClickHouse Cloud, мы рекомендуем вместо этого ClickPipes. ClickPipes изначально поддерживает подключения к частным сетям, независимое масштабирование ресурсов ингестии и кластера, а также комплексный мониторинг при стриминге данных из Kafka в ClickHouse.
Для использования движка таблицы Kafka вам следует в общих чертах понимать, как работают materialized views ClickHouse.

Обзор

Сначала мы рассмотрим наиболее распространённый сценарий использования: применение движка таблицы Kafka для вставки данных из Kafka в ClickHouse. Движок таблицы Kafka позволяет ClickHouse напрямую читать из топика Kafka. Хотя это удобно для просмотра сообщений в топике, по своей конструкции этот движок допускает только однократное чтение: когда к таблице выполняется запрос, он считывает данные из очереди и увеличивает смещение консьюмера перед возвратом результатов вызывающей стороне. На практике данные нельзя прочитать повторно без сброса этих смещений. Чтобы сохранить данные, прочитанные через движок таблицы, нужен способ захватить их и вставить в другую таблицу. Эту возможность нативно предоставляют materialized view на основе триггеров. Materialized view инициирует чтение из движка таблицы, получая батчи документов. Предложение TO определяет пункт назначения данных — обычно это таблица семейства MergeTree. Этот процесс показан ниже:

Шаги

1. Подготовка
Если у вас уже есть данные в целевом топике, вы можете адаптировать приведённый ниже пример для использования со своим набором данных. В качестве альтернативы здесь доступен пример набора данных GitHub. Этот набор данных используется в примерах ниже; в нём используется сокращённая схема и подмножество строк (в частности, мы ограничились событиями GitHub, относящимися к репозиторию ClickHouse) по сравнению с полным набором данных, доступным здесь, для краткости. Тем не менее его достаточно, чтобы работало большинство запросов, опубликованных вместе с набором данных.
2. Настройте ClickHouse
Этот шаг обязателен, если вы подключаетесь к защищённому кластеру Kafka. Эти параметры нельзя передать через команды SQL DDL, поэтому их нужно настроить в файле config.xml ClickHouse. Мы предполагаем, что вы подключаетесь к экземпляру, защищённому с помощью SASL. Это самый простой вариант при работе с Confluent Cloud.
Либо поместите приведенный выше фрагмент в новый файл в каталоге conf.d/, либо добавьте его в существующие файлы конфигурации. О настройках, доступных для конфигурирования, см. здесь. Мы также создадим базу данных KafkaEngine, которую будем использовать в этом руководстве:
После создания базы данных переключитесь на неё:
3. Создайте целевую таблицу
Подготовьте целевую таблицу. В примере ниже для краткости используется сокращённая схема GitHub. Обратите внимание: хотя здесь используется движок таблицы MergeTree, этот пример можно легко адаптировать для любого движка из семейства MergeTree.
4. Создайте и заполните топик
Далее создадим топик. Для этого можно использовать несколько инструментов. Если Kafka запущена локально на нашей машине или внутри контейнера Docker, хорошо подойдет RPK. Мы можем создать топик с именем github и 5 партициями, выполнив следующую команду:
Если мы используем Kafka в Confluent Cloud, возможно, нам будет удобнее воспользоваться Confluent CLI:
Теперь нам нужно наполнить этот топик данными, и для этого мы используем kcat. Если вы запускаете Kafka локально с отключённой аутентификацией, можно выполнить команду, похожую на следующую:
Или следующее, если в нашем кластере Kafka для аутентификации используется SASL:
Датасет содержит 200 000 строк, поэтому его приём займёт всего несколько секунд. Если вы хотите работать с более крупным датасетом, ознакомьтесь с разделом о больших датасетах в репозитории ClickHouse/kafka-samples на GitHub.
5. Создайте таблицу с движком таблицы Kafka
В приведённом ниже примере создаётся таблица с движком Kafka и той же схемой, что и у таблицы MergeTree. Это не является обязательным требованием, поскольку в целевой таблице можно использовать alias или эфемерные столбцы. Однако настройки важны — обратите внимание, что JSONEachRow используется как тип данных для чтения JSON из топика Kafka. Значения github и clickhouse обозначают имя топика и имя группы потребителей соответственно. На самом деле топики могут быть списком значений.
Ниже мы рассмотрим настройки движка и вопросы настройки производительности. На этом этапе простой запрос select к таблице github_queue должен прочитать несколько строк. Обратите внимание, что это сдвинет смещения consumer вперёд, из-за чего эти строки нельзя будет прочитать повторно без сброса. Также обратите внимание на limit и обязательный parameter stream_like_engine_allow_direct_select.
6. Создайте materialized view
Materialized view свяжет две ранее созданные таблицы, считывая данные из движка таблицы Kafka и вставляя их в целевую таблицу MergeTree. Можно выполнить ряд преобразований данных. Мы ограничимся простым чтением и вставкой. Использование * предполагает, что имена столбцов совпадают (с учетом регистра).
В момент создания materialized view подключается к движку Kafka и начинает чтение, вставляя строки в целевую таблицу. Этот процесс будет продолжаться бесконечно: новые сообщения, добавляемые в Kafka, будут считываться. При необходимости вы можете повторно запустить скрипт вставки, чтобы добавить в Kafka дополнительные сообщения.
7. Убедитесь, что строки были вставлены
Убедитесь, что данные есть в целевой таблице:
Вы должны увидеть 200 000 строк:

Типовые операции

Остановка и возобновление потребления сообщений
Чтобы остановить потребление сообщений, можно отсоединить таблицу с движком Kafka:
Это не повлияет на смещения группы потребителей. Чтобы возобновить чтение и продолжить с предыдущего смещения, повторно подключите таблицу.
Добавление метаданных Kafka
Может быть полезно отслеживать метаданные исходных сообщений Kafka после их приёма в ClickHouse. Например, может понадобиться узнать, какую часть конкретного топика или партиции мы уже обработали. Для этого движок таблицы Kafka предоставляет несколько виртуальных столбцов. Их можно сохранить в виде столбцов в целевой таблице, изменив схему и оператор SELECT в materialized view. Сначала выполните описанную выше операцию остановки, а затем добавьте столбцы в целевую таблицу.
Ниже мы добавляем информационные столбцы, чтобы указать исходный топик и партицию, из которой пришла строка.
Далее нужно убедиться, что виртуальные столбцы сопоставлены должным образом. Виртуальные столбцы имеют префикс _. Полный список виртуальных столбцов можно найти здесь. Чтобы обновить таблицу, добавив виртуальные столбцы, потребуется удалить materialized view, повторно выполнить ATTACH для таблицы с движком Kafka и заново создать materialized view.
У недавно прочитанных строк должны быть метаданные.
Результат будет выглядеть так:
Изменение настроек движка Kafka
Мы рекомендуем удалить таблицу с движком Kafka и заново создать её с новыми настройками. В ходе этого процесса materialized view изменять не нужно — потребление сообщений возобновится, как только таблица с движком Kafka будет создана заново.
Отладка проблем
Ошибки, например связанные с аутентификацией, не отображаются в ответах на DDL-запросы к движку Kafka. Для диагностики проблем рекомендуем использовать основной файл журнала ClickHouse clickhouse-server.err.log. Дополнительное журналирование трассировки для базовой клиентской библиотеки Kafka librdkafka можно включить через конфигурацию.
Обработка некорректных сообщений
Kafka часто используют как «свалку» для данных. Из-за этого в топиках оказываются сообщения в разных форматах и с несогласованными именами полей. По возможности избегайте этого и используйте возможности Kafka, такие как Kafka Streams или ksqlDB, чтобы сообщения были корректно сформированы и согласованы еще до записи в Kafka. Если это невозможно, в ClickHouse есть несколько возможностей, которые могут помочь.
  • Обрабатывайте поле сообщения как строки. При необходимости в операторе materialized view можно использовать функции для очистки и приведения типов. Это не решение для production, но может помочь при разовой ингестии.
  • Если вы читаете JSON из топика в формате JSONEachRow, используйте настройку input_format_skip_unknown_fields. При записи данных ClickHouse по умолчанию генерирует исключение, если входные данные содержат столбцы, которых нет в целевой таблице. Однако если эта опция включена, такие лишние столбцы будут игнорироваться. Опять же, это не решение уровня production и оно может запутать других.
  • Обратите внимание на настройку kafka_skip_broken_messages. Она требует, чтобы пользователь задал допустимый уровень ошибок для некорректных сообщений на block с учетом kafka_max_block_size. Если этот порог превышен (в абсолютном количестве сообщений), снова будет применяться стандартное поведение с исключением, а остальные сообщения будут пропущены.
Семантика доставки и проблемы с дубликатами
Kafka движок таблицы имеет семантику at-least-once. Дубликаты возможны в ряде известных, хотя и редких, случаев. Например, сообщения могут быть прочитаны из Kafka и успешно вставлены в ClickHouse. Но до того, как новое смещение будет зафиксировано, соединение с Kafka может быть потеряно. В такой ситуации требуется повторная попытка обработки блока. Блок может быть дедуплицирован при использовании distributed таблицы или ReplicatedMergeTree в качестве целевой таблицы. Хотя это снижает вероятность появления дублирующихся строк, такой подход опирается на идентичность блоков. Такие события, как ребалансировка Kafka, могут нарушить это допущение, что в редких случаях приводит к дубликатам.
Вставки на основе кворума
В некоторых случаях, когда в ClickHouse нужны более высокие гарантии доставки, могут потребоваться вставки на основе кворума. Это нельзя настроить для materialized view или целевой таблицы. Однако это можно задать для пользовательских профилей, например.

Из ClickHouse в Kafka

Хотя это и более редкий сценарий использования, данные ClickHouse также можно сохранять в Kafka. Например, мы вручную вставим строки в таблицу с движком таблицы Kafka. Эти данные будут прочитаны тем же Kafka-движком, а его materialized view поместит их в таблицу MergeTree. Наконец, мы покажем, как использовать materialized views при вставке в Kafka, чтобы наполнять таблицы на основе существующих исходных таблиц.

Шаги

Нашу первоначальную цель лучше всего иллюстрирует следующая схема: Мы предполагаем, что вы уже создали таблицы и представления, описанные в разделе Kafka to ClickHouse, и что топик был полностью прочитан.
1. Вставка строк напрямую
Сначала проверьте количество строк в целевой таблице.
У вас должно быть 200 000 строк:
Теперь выполните вставку строк из целевой таблицы GitHub обратно в движок таблицы Kafka github_queue. Обратите внимание: мы используем формат JSONEachRow и ограничиваем выборку 100 строками с помощью LIMIT.
Снова пересчитайте строки в GitHub, чтобы убедиться, что их количество увеличилось на 100. Как показано на схеме выше, строки были вставлены в Kafka через движок таблицы Kafka, затем тот же движок снова их считал и вставил в целевую таблицу GitHub с помощью нашего materialized view!
Вы должны увидеть ещё 100 строк:
2. Использование materialized views
Мы можем использовать materialized views, чтобы отправлять сообщения в движок Kafka (и в топик), когда документы вставляются в таблицу. Когда строки вставляются в таблицу GitHub, срабатывает materialized view, в результате чего строки снова вставляются в движок Kafka и в новый топик. Это лучше всего показано на иллюстрации: Создайте новый Kafka топик github_out или его эквивалент. Убедитесь, что движок таблицы Kafka github_out_queue указывает на этот топик.
Теперь создайте новое materialized view github_out_mv, которое будет указывать на таблицу GitHub и при срабатывании выполнять вставку строк в указанный выше движок. В результате новые записи из таблицы GitHub будут отправляться в наш новый топик Kafka.
Если выполнить вставку в исходный топик github, созданный в разделе Kafka to ClickHouse, документы автоматически появятся в топике “github_clickhouse”. Убедитесь в этом с помощью встроенных инструментов Kafka. Например, ниже мы выполняем вставку 100 строк в топик github с помощью kcat для топика, размещённого в Confluent Cloud:
Чтение из топика github_out должно подтвердить, что сообщения доставлены.
Хотя это и довольно сложный пример, он показывает возможности materialized views в сочетании с движком Kafka.

Кластеры и производительность

Работа с кластерами ClickHouse

Благодаря группам потребителей Kafka несколько экземпляров ClickHouse могут читать из одного и того же топика. Каждому потребителю назначается одна партиция топика в соотношении 1:1. При масштабировании чтения ClickHouse с использованием движка таблицы Kafka учитывайте, что общее число потребителей в кластере не может превышать число партиций в топике. Поэтому заранее убедитесь, что для топика настроено достаточное число партиций. Несколько экземпляров ClickHouse можно настроить на чтение из топика с одним и тем же идентификатором группы потребителей, указанным при создании движка таблицы Kafka. Таким образом, каждый экземпляр будет читать из одной или нескольких партиций, записывая сегменты в свою локальную целевую таблицу. Целевые таблицы, в свою очередь, можно настроить на использование ReplicatedMergeTree, чтобы справляться с дублированием данных. Такой подход позволяет масштабировать чтение из Kafka вместе с кластером ClickHouse при условии, что в Kafka достаточно партиций.

Настройка производительности

При увеличении пропускной способности таблицы с движком Kafka учитывайте следующее:
  • Производительность зависит от размера сообщений, формата и типов целевых таблиц. Показатель в 100 тыс. строк/с для одного движка таблицы считается достижимым. По умолчанию сообщения читаются блоками; это регулируется параметром kafka_max_block_size. По умолчанию он равен max_insert_block_size, то есть 1,048,576. Если сообщения не являются исключительно большими, это значение почти всегда стоит увеличивать. Значения в диапазоне от 500 тыс. до 1 млн вполне обычны. Протестируйте и оцените, как это влияет на пропускную способность.
  • Число потребителей для движка таблицы можно увеличить с помощью kafka_num_consumers. Однако по умолчанию вставки будут выполняться последовательно в одном потоке, если kafka_thread_per_consumer не изменить относительно значения по умолчанию 1. Установите это значение в 1, чтобы сбросы на диск выполнялись параллельно. Обратите внимание: создание таблицы с движком Kafka с N потребителями (и kafka_thread_per_consumer=1) логически эквивалентно созданию N движков Kafka, каждый со своей materialized view и kafka_thread_per_consumer=0.
  • Увеличение числа потребителей не даётся бесплатно. Каждый потребитель поддерживает собственные буферы и потоки, увеличивая накладные расходы на сервер. Учитывайте эти накладные расходы и по возможности сначала линейно масштабируйте нагрузку по кластеру.
  • Если пропускная способность потока сообщений Kafka непостоянна, а задержки допустимы, рассмотрите увеличение stream_flush_interval_ms, чтобы на диск сбрасывались более крупные блоки.
  • background_message_broker_schedule_pool_size задаёт количество потоков, выполняющих фоновые задачи. Эти потоки используются для стриминга Kafka. Этот параметр применяется при запуске сервера ClickHouse и не может быть изменён в пользовательском сеансе; по умолчанию его значение равно 16. Если вы видите тайм-ауты в журнале, возможно, стоит увеличить это значение.
  • Для взаимодействия с Kafka используется библиотека librdkafka, которая также создаёт потоки. Поэтому большое число таблиц Kafka или потребителей может приводить к большому количеству переключений контекста. Либо распределите эту нагрузку по кластеру, по возможности реплицируя только целевые таблицы, либо рассмотрите использование движка таблицы для чтения из нескольких топиков — поддерживается список значений. Из одной таблицы можно читать через несколько materialized view, каждая из которых фильтрует данные из определённого топика.
Любые изменения настроек следует тестировать. Мы рекомендуем отслеживать отставание потребителей Kafka, чтобы убедиться, что масштабирование выбрано правильно.

Дополнительные настройки

Помимо настроек, рассмотренных выше, интерес могут представлять следующие:
  • Kafka_max_wait_ms — Время ожидания в миллисекундах при чтении сообщений из Kafka перед повторной попыткой. Задаётся на уровне профиля пользователя; значение по умолчанию — 5000.
Все настройки базовой библиотеки librdkafka также можно указывать в файлах конфигурации ClickHouse внутри элемента kafka — имена настроек должны быть XML-элементами, в которых точки заменены на подчёркивания, например.
Это настройки для опытных пользователей, и мы рекомендуем обратиться к документации Kafka за более подробными разъяснениями.
Последнее изменение 12 июня 2026 г.