Kafka в ClickHouse
Если вы используете ClickHouse Cloud, мы рекомендуем вместо этого ClickPipes. ClickPipes изначально поддерживает подключения к частным сетям, независимое масштабирование ресурсов ингестии и кластера, а также комплексный мониторинг при стриминге данных из Kafka в ClickHouse.
Обзор
Шаги
1
Подготовьте
Если у вас есть данные, загруженные в целевой топик, вы можете адаптировать приведённый ниже пример для использования со своим набором данных. Либо можно воспользоваться примером набора данных GitHub, доступным здесь. Этот набор данных используется в примерах ниже и для краткости включает сокращённую схему и подмножество строк (в частности, мы ограничиваемся событиями GitHub, относящимися к репозиторию ClickHouse) по сравнению с полным набором данных, доступным здесь. Тем не менее его достаточно, чтобы работало большинство запросов, опубликованных вместе с набором данных.
2
Настройте ClickHouse
Этот шаг обязателен, если вы подключаетесь к защищённому кластеру Kafka. Эти настройки нельзя передать через команды SQL DDL, их необходимо задать в Либо поместите приведённый выше фрагмент в новый файл в каталоге conf.d/, либо объедините его с существующими файлами конфигурации. Сведения о доступных для настройки параметрах см. здесь.Мы также создадим базу данных с именем После создания базы данных вам потребуется переключиться на неё:
config.xml ClickHouse. Предполагается, что вы подключаетесь к экземпляру, защищённому с помощью SASL. Это самый простой способ при работе с Confluent Cloud.KafkaEngine, которую будем использовать в этом руководстве:3
Создайте целевую таблицу
Подготовьте целевую таблицу. В примере ниже для краткости используется сокращённая схема GitHub. Обратите внимание: хотя здесь используется движок таблицы MergeTree, этот пример можно легко адаптировать для любого движка из семейства MergeTree.
4
Создайте топик и заполните его
Далее мы создадим топик. Для этого можно использовать несколько инструментов. Если Kafka запущен локально на вашей машине или в контейнере Docker, хорошо подойдёт RPK. Мы можем создать топик с именем Если мы запускаем Kafka в Confluent Cloud, то можем предпочесть использовать Confluent CLI:Теперь нужно заполнить этот topic данными; для этого мы будем использовать kcat. Если Kafka запущен локально и аутентификация отключена, можно выполнить команду, аналогичную следующей:Или следующий вариант, если в нашем кластере Kafka для аутентификации используется SASL:Датасет содержит 200 000 строк, поэтому его приём займёт всего несколько секунд. Если вы хотите работать с более крупным датасетом, ознакомьтесь с разделом о больших датасетах в GitHub-репозитории ClickHouse/kafka-samples.
github с 5 партициями, выполнив следующую команду:5
Создайте движок таблицы Kafka
В примере ниже создаётся движок таблицы с той же схемой, что и у таблицы MergeTree. Это не является строгим требованием, так как в целевой таблице могут быть псевдонимы или эфемерные столбцы. Однако настройки важны — обратите внимание на использование Ниже мы рассмотрим настройки движка и настройку производительности. На этом этапе простой запрос select к таблице
JSONEachRow в качестве типа данных для чтения JSON из топика Kafka. Значения github и clickhouse представляют собой имена топика и группы потребителей соответственно. На самом деле топики могут быть списком значений.github_queue должен прочитать несколько строк. Обратите внимание, что это сдвинет смещения потребителя вперёд, из-за чего эти строки нельзя будет прочитать повторно без сброса. Обратите внимание на ограничение и обязательный параметр stream_like_engine_allow_direct_select.6
Создайте materialized view
materialized view свяжет две ранее созданные таблицы, считывая данные из таблицы Kafka и вставляя их в целевую таблицу MergeTree. Мы можем выполнить ряд преобразований данных. Здесь мы ограничимся простым чтением и вставкой. Использование * предполагает, что имена столбцов совпадают (с учётом регистра).В момент создания materialized view подключается к движку Kafka и начинает читать данные, вставляя строки в целевую таблицу. Этот процесс будет продолжаться бесконечно, при этом новые сообщения, вставляемые в Kafka, будут потребляться. При необходимости вы можете повторно запустить скрипт вставки, чтобы добавить в Kafka дополнительные сообщения.
7
Убедитесь, что строки вставлены
Убедитесь, что в целевой таблице есть данные:Вы должны увидеть 200 000 строк:
Типовые операции
Остановка и возобновление потребления сообщений
Добавление метаданных Kafka
_.
Полный список виртуальных столбцов можно найти здесь.
Чтобы обновить таблицу, добавив виртуальные столбцы, потребуется удалить materialized view, повторно выполнить ATTACH для таблицы с движком Kafka и заново создать materialized view.
Изменение настроек движка Kafka
Отладка проблем
clickhouse-server.err.log. Дополнительное журналирование трассировки для базовой клиентской библиотеки Kafka librdkafka можно включить через конфигурацию.
Обработка некорректных сообщений
- Обрабатывайте поле сообщения как строки. При необходимости в операторе materialized view можно использовать функции для очистки и приведения типов. Это не решение для production, но может помочь при разовой ингестии.
- Если вы читаете JSON из топика в формате JSONEachRow, используйте настройку
input_format_skip_unknown_fields. При записи данных ClickHouse по умолчанию генерирует исключение, если входные данные содержат столбцы, которых нет в целевой таблице. Однако если эта опция включена, такие лишние столбцы будут игнорироваться. Опять же, это не решение уровня production и оно может запутать других. - Обратите внимание на настройку
kafka_skip_broken_messages. Она требует, чтобы пользователь задал допустимый уровень ошибок для некорректных сообщений на block с учетомkafka_max_block_size. Если этот порог превышен (в абсолютном количестве сообщений), снова будет применяться стандартное поведение с исключением, а остальные сообщения будут пропущены.
Семантика доставки и проблемы с дубликатами
Вставки на основе кворума
Из ClickHouse в Kafka
Шаги
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 Теперь создайте новое materialized view Если выполнить вставку в исходный топик github, созданный в разделе Kafka в ClickHouse, документы автоматически появятся в топике “github_clickhouse”. Убедитесь в этом с помощью встроенных инструментов Kafka. Например, ниже мы выполняем вставку 100 строк в топик github с помощью kcat для топика, размещённого в Confluent Cloud:Чтение из топика Несмотря на то что пример достаточно сложный, он наглядно демонстрирует возможности materialized views при использовании совместно с движком Kafka.
github_out или его эквивалент. Убедитесь, что движок таблицы Kafka github_out_queue указывает на этот топик.github_out_mv, которое будет указывать на таблицу GitHub и при срабатывании выполнять вставку строк в указанный выше движок. В результате новые записи из таблицы GitHub будут отправляться в наш новый топик Kafka.github_out должно подтвердить доставку сообщений.Кластеры и производительность
Работа с кластерами ClickHouse
Настройка производительности
- Производительность зависит от размера сообщений, формата и типов целевых таблиц. Показатель в 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_max_wait_ms — Время ожидания в миллисекундах при чтении сообщений из Kafka перед повторной попыткой. Задаётся на уровне профиля пользователя; значение по умолчанию — 5000.