Kafka в ClickHouse
Если вы используете ClickHouse Cloud, мы рекомендуем вместо этого ClickPipes. ClickPipes изначально поддерживает подключения к частным сетям, независимое масштабирование ресурсов ингестии и кластера, а также комплексный мониторинг при стриминге данных из Kafka в ClickHouse.
Обзор
Шаги
1. Подготовка
2. Настройте ClickHouse
config.xml ClickHouse. Мы предполагаем, что вы подключаетесь к экземпляру, защищённому с помощью SASL. Это самый простой вариант при работе с Confluent Cloud.
KafkaEngine, которую будем использовать в этом руководстве:
3. Создайте целевую таблицу
4. Создайте и заполните топик
github и 5 партициями, выполнив следующую команду:
5. Создайте таблицу с движком таблицы Kafka
JSONEachRow используется как тип данных для чтения JSON из топика Kafka. Значения github и clickhouse обозначают имя топика и имя группы потребителей соответственно. На самом деле топики могут быть списком значений.
github_queue должен прочитать несколько строк. Обратите внимание, что это сдвинет смещения consumer вперёд, из-за чего эти строки нельзя будет прочитать повторно без сброса. Также обратите внимание на limit и обязательный parameter stream_like_engine_allow_direct_select.
6. Создайте materialized view
7. Убедитесь, что строки были вставлены
Типовые операции
Остановка и возобновление потребления сообщений
Добавление метаданных 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. Вставка строк напрямую
2. Использование materialized views
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.