Если вам нужна помощь, пожалуйста, сообщите о проблеме в репозитории или задайте вопрос в публичном Slack ClickHouse.
Лицензия
Требования к окружению
Матрица совместимости версий
Основные возможности
- Сразу поддерживает семантику «ровно один раз». Основана на новой возможности ядра ClickHouse под названием KeeperMap (коннектор использует её как хранилище состояния) и позволяет обойтись минималистичной архитектурой.
- Поддержка сторонних хранилищ состояния: сейчас по умолчанию используется In-memory, но также можно использовать KeeperMap (скоро будет добавлен Redis).
- Нативная интеграция: разработана, поддерживается и сопровождается ClickHouse.
- Непрерывно тестируется с ClickHouse Cloud.
- Вставка данных как по объявленной схеме, так и без схемы.
- Поддержка всех типов данных ClickHouse.
Инструкции по установке
Подготовьте сведения о подключении
Сведения о подключении для вашего сервиса ClickHouse Cloud доступны в консоли ClickHouse Cloud.
Выберите сервис и нажмите Connect:

curl.

Общие инструкции по установке
- Загрузите ZIP-архив с JAR-файлом коннектора со страницы Releases репозитория Приёмник ClickHouse Kafka Connect.
- Извлеките содержимое ZIP-файла и скопируйте его в нужное место.
- Добавьте путь к каталогу плагина в параметр plugin.path в файле свойств Connect, чтобы Confluent Platform могла найти плагин.
- Укажите в конфигурации имя топика, имя хоста экземпляра ClickHouse и пароль.
- Перезапустите Confluent Platform.
- Если вы используете Confluent Platform, войдите в интерфейс Confluent Control Center и убедитесь, что ClickHouse Sink доступен в списке доступных коннекторов.
Параметры конфигурации
- сведения о подключении: hostname (обязательно) и порт (необязательно)
- учетные данные пользователя: пароль (обязательно) и имя пользователя (необязательно)
- класс коннектора:
com.clickhouse.kafka.connect.ClickHouseSinkConnector(обязательно) - topics или topics.regex: топики Kafka, которые нужно опрашивать; имена топиков должны совпадать с именами таблиц (обязательно)
- конвертеры key и value: указываются в зависимости от типа данных в вашем топике. Обязательно, если они еще не заданы в конфигурации воркера.
Целевые таблицы
Предобработка
Поддерживаемые типы данных
-
(1) - JSON поддерживается только при наличии в настройках ClickHouse параметра
input_format_binary_read_json_as_string=1. Это работает только для семейства форматов RowBinary, и этот параметр влияет на все столбцы в запросе вставки, поэтому все они должны иметь тип String. В этом случае connector преобразует STRUCT в строку JSON. -
(2) - Если struct содержит union, например
oneof, converter следует настроить так, чтобы он НЕ добавлял prefix/suffix к именам полей. Для этого используйте параметрgenerate.index.for.unions=falseвProtobufConverter.
Варианты конфигурации
Базовая конфигурация
localhost:8443 с включенным SSL; данные представлены в JSON без схемы.
Приведённая выше конфигурация коннектора требует включить переопределение клиентских настроек в конфигурации воркера с помощью
connector.client.config.override.policy=All. Подробнее см. в документации Kafka Connect.Базовая конфигурация для нескольких топиков
Базовая конфигурация с DLQ
Поддержка схемы Avro
Сопоставление типов Avro
io.confluent.connect.avro.AvroConverter — официальной реализации сериализатора/десериализатора Avro для Kafka Connect. Подробную информацию о логике преобразования см. в документации Kafka Connect.
✅: Поддерживается
❌: Не поддерживается
️⚠️: Поддерживается частично
Сопоставление типов Kafka Connect и типов ClickHouse см. в разделе Поддерживаемые типы данных.
Неподдерживаемые схемы Avro
- логический тип
decimalдля типа fixed
- union-типы с Nullable
- объединения в записях
Поддержка схем Protobuf
Сопоставление типов Protobuf
io.confluent.connect.protobuf.ProtobufConverter — официальной реализации сериализатора/десериализатора Protobuf для Kafka Connect. Подробную информацию о логике преобразования см. в документации Kafka Connect.
✅: Поддерживается
❌: Не поддерживается
️⚠️: Поддерживается частично
Сопоставление типов между Kafka Connect и ClickHouse см. в разделе Поддерживаемые типы данных.
Примечание о сопоставлении полей oneof со столбцами ClickHouse
oneof) в тип ClickHouse Variant. Вместо этого укажите поля oneof в схеме таблицы ClickHouse как отдельные поля Nullable.
Например:
Неподдерживаемые схемы Protobuf
- union-типы из нескольких сообщений (до версии CH 26.1)
allow_experimental_nullable_tuple_type=1 (см. эту страницу документации).
Поддержка схем JSON
Поддержка типа String
Внутренняя буферизация
poll() и сбрасывать их в ClickHouse более крупными батчами. Это может повысить пропускную способность в рабочих нагрузках, где каждый опрос возвращает множество небольших батчей по отдельным партициям.
Ключевые особенности:
bufferCountопределяет, сколько записей буферизуется перед сбросом.bufferFlushTimeзадает максимальное время ожидания (в миллисекундах) перед сбросом буферизованных записей.bufferFlushTimeдействует только приbufferCount > 0.bufferCount=0иbufferFlushTime=0оставляют буферизацию отключенной (поведение по умолчанию).- Буферизация не поддерживается, если
exactlyOnce=true.
exactlyOnce=false в конфигурации коннектора, либо отключите буферизацию, задав bufferCount=0.
Пример:
Логирование
Мониторинг
Метрики ClickHouse
Метрики Kafka Producer/Consumer
records-sent-total: Общее количество записей, отправленных в topicbytes-sent-total: Общее количество байт, отправленных в topicrecord-send-rate: Средняя скорость отправки записей в секундуbyte-rate: Среднее количество байт, отправляемых в секундуcompression-rate: Достигнутый коэффициент сжатия
records-sent-total: Общее количество записей, отправленных в партициюbytes-sent-total: Общее количество байт, отправленных в партициюrecords-lag: Текущее отставание в партицииrecords-lead: Текущее опережение в партицииreplica-fetch-lag: Информация об отставании реплик
connection-creation-total: Общее количество соединений, установленных с узлом Kafkaconnection-close-total: Общее количество закрытых соединенийrequest-total: Общее количество запросов, отправленных узлуresponse-total: Общее количество ответов, полученных от узлаrequest-rate: Средняя частота запросов в секундуresponse-rate: Средняя частота ответов в секунду
- Пропускную способность: отслеживать скорость приёма данных
- Отставание: выявлять узкие места и задержки обработки
- Сжатие: оценивать эффективность сжатия данных
- Состояние соединений: контролировать сетевую связность и стабильность
Метрики фреймворка Kafka Connect
task-count: Общее количество задач в коннектореrunning-task-count: Количество задач, выполняющихся в данный моментpaused-task-count: Количество задач, приостановленных в данный моментfailed-task-count: Количество задач, завершившихся с ошибкойdestroyed-task-count: Количество удалённых задачunassigned-task-count: Количество неназначенных задач
running, paused, failed, destroyed, unassigned
Метрики ошибок:
deadletterqueue-produce-failures: Количество неудачных записей в DLQdeadletterqueue-produce-requests: Общее число попыток записи в DLQlast-error-timestamp: Временная метка последней ошибкиrecords-skip-total: Общее количество записей, пропущенных из-за ошибокrecords-retry-total: Общее количество записей, для которых выполнялись повторные попыткиerrors-total: Общее количество возникших ошибок
offset-commit-failures: Количество неудачных фиксаций смещенийoffset-commit-avg-time-ms: Среднее время фиксации смещенийoffset-commit-max-time-ms: Максимальное время фиксации смещенийput-batch-avg-time-ms: Среднее время обработки батчаput-batch-max-time-ms: Максимальное время обработки батчаsource-record-poll-total: Общее количество записей, полученных при опросе
Рекомендации по мониторингу
- Отслеживайте отставание потребителя: Следите за
records-lagпо каждой партиции, чтобы выявлять узкие места в обработке - Контролируйте уровень ошибок: Следите за
errors-totalиrecords-skip-total, чтобы выявлять проблемы с качеством данных - Следите за состоянием задач: Отслеживайте метрики состояния задач, чтобы убедиться, что они работают корректно
- Измеряйте пропускную способность: Используйте
records-send-rateиbyte-rateдля отслеживания производительности ингестии - Следите за состоянием соединений: Проверяйте метрики соединений на уровне узла, чтобы выявлять проблемы в сети
- Отслеживайте эффективность сжатия: Используйте
compression-rateдля оптимизации передачи данных
Ограничения
- Удаление данных не поддерживается.
- Размер батча наследуется из свойств Kafka Consumer.
- При использовании KeeperMap для exactly-once, если смещение было изменено или сброшено назад, необходимо удалить содержимое KeeperMap для соответствующего топика. (Подробнее см. в руководстве по устранению неполадок ниже)
Настройка производительности и оптимизация пропускной способности
Когда нужна настройка производительности?
- Высоконагруженные рабочие нагрузки: когда обрабатываются миллионы событий в секунду из топиков Kafka
- Отставание потребителя: когда ваш коннектор не успевает за скоростью поступления данных, и отставание растет
- Ограниченные ресурсы: когда нужно оптимизировать использование CPU, памяти или сети
- Несколько топиков: когда данные одновременно потребляются из нескольких топиков с большим объемом трафика
- Маленький размер сообщений: когда приходится работать со множеством небольших сообщений, для которых полезен батчинг на стороне сервера
- Вы обрабатываете небольшие или умеренные объемы данных (< 10,000 сообщений/секунду)
- Отставание потребителя остается стабильным и приемлемым для вашего сценария использования
- Настройки коннектора по умолчанию уже соответствуют вашим требованиям к пропускной способности
- Ваш кластер ClickHouse без труда справляется с входящей нагрузкой
Как устроен поток данных
- Фреймворк Kafka Connect в фоновом режиме забирает сообщения из топиков Kafka
- Коннектор считывает сообщения из внутреннего буфера фреймворка
- Коннектор объединяет сообщения в батчи в зависимости от размера выборки
- ClickHouse получает батч вставки по HTTP/S
- ClickHouse обрабатывает вставку (синхронно или асинхронно)
Настройка размера батча в Kafka Connect
Настройки fetch
fetch.min.bytes: Минимальный объём данных, после достижения которого фреймворк передаёт данные коннектору (по умолчанию: 1 байт)fetch.max.bytes: Максимальный объём данных, получаемый за один запрос (по умолчанию: 52428800 / 50 МБ)fetch.max.wait.ms: Максимальное время ожидания перед возвратом данных, если значениеfetch.min.bytesне достигнуто (по умолчанию: 500 мс)
В Confluent Cloud для изменения этих настроек нужно открыть обращение в поддержку через Confluent Cloud.
Настройки опроса
max.poll.records: Максимальное количество записей, возвращаемых за один опрос (по умолчанию: 500)max.partition.fetch.bytes: Максимальный объём данных на партицию (по умолчанию: 1048576 / 1 MB)
В Confluent Cloud для изменения этих настроек нужно открыть обращение в службу поддержки через Confluent Cloud.
Рекомендуемые настройки для высокой пропускной способности
Чтобы использовать указанные выше свойства, нужно разрешить переопределение клиентских настроек в конфигурации воркера с помощью
connector.client.config.override.policy=All. Подробнее см. в документации Kafka Connect.- Более крупные батчи = Лучшая производительность ингестии в ClickHouse, меньше частей, ниже накладные расходы
- Более крупные батчи = Более высокое использование памяти, возможное увеличение сквозной задержки
- Слишком крупные батчи = Риск тайм-аутов, ошибок OutOfMemory или превышения
max.poll.interval.ms
Асинхронные вставки
Когда использовать асинхронную вставку
- Много небольших батчей: Ваш коннектор часто отправляет небольшие батчи (< 1000 строк в батче)
- Высокий параллелизм: Несколько задач коннектора записывают данные в одну и ту же таблицу
- Распределённое развертывание: На разных хостах запущено много экземпляров коннектора
- Накладные расходы на создание частей: Вы сталкиваетесь с ошибками “too many parts”
- Смешанная рабочая нагрузка: Вы совмещаете ингестию в реальном времени с нагрузкой от запросов
- Вы уже отправляете большие батчи (> 10 000 строк в батче) с контролируемой частотой
- Вам нужна немедленная доступность данных для запросов (запросы должны видеть данные мгновенно)
- Семантика «ровно один раз» с
wait_for_async_insert=0не соответствует вашим требованиям - В вашем случае вместо этого можно выиграть от улучшения батчинга на стороне клиента
Как работают асинхронные вставки
- Получает запрос на вставку от коннектора
- Записывает данные во внутренний буфер в памяти (вместо немедленной записи на диск)
- Возвращает коннектору успешный результат (если
wait_for_async_insert=0) - Записывает буфер на диск, когда выполняется одно из следующих условий:
- Буфер достигает
async_insert_max_data_size(по умолчанию: 100 MB) - С момента первой вставки прошло
async_insert_busy_timeout_msмиллисекунд (по умолчанию: 1000 ms) - Достигнуто максимальное количество накопленных запросов (
async_insert_max_query_number, по умолчанию: 100)
- Буфер достигает
Включение асинхронной вставки
clickhouseSettings:
async_insert=1: Включает асинхронную вставкуwait_for_async_insert=1(рекомендуется): Коннектор ожидает, пока данные будут записаны в хранилище ClickHouse, прежде чем подтверждать получение. Обеспечивает гарантии доставки.wait_for_async_insert=0: Коннектор подтверждает получение сразу после буферизации. Обеспечивает более высокую производительность, но данные могут быть потеряны при сбое сервера до сброса на диск.
Тонкая настройка поведения асинхронной вставки
async_insert_max_data_size(по умолчанию: 104857600 / 100 MB): Максимальный размер буфера до сбросаasync_insert_busy_timeout_ms(по умолчанию: 1000): Максимальное время (мс) до сбросаasync_insert_stale_timeout_ms(по умолчанию: 0): Время (мс) с момента последней вставки до сбросаasync_insert_max_query_number(по умолчанию: 100): Максимальное количество запросов до сброса
- Преимущества: Меньше частей, выше производительность слияния, ниже нагрузка на CPU, выше пропускная способность при высоком параллелизме
- Что учитывать: Данные становятся доступны для запросов не сразу, немного увеличивается общая задержка
- Риски: Потеря данных при сбое сервера, если
wait_for_async_insert=0, возможна повышенная нагрузка на оперативную память при больших буферах
Асинхронные вставки с семантикой «ровно один раз»
exactlyOnce=true с асинхронными вставками:
wait_for_async_insert=1 вместе с exactly-once, чтобы смещения фиксировались только после сохранения данных.
Подробнее об async inserts см. в документации ClickHouse об async inserts.
Параллелизм коннектора
Количество задач на коннектор
- Максимально эффективное количество задач = количество партиций топика
- Каждая задача поддерживает собственное подключение к ClickHouse
- Больше задач = выше накладные расходы и выше вероятность конкуренции за ресурсы
tasks.max, равного количеству партиций топика, затем корректируйте это значение на основе метрик CPU и пропускной способности.
Игнорирование партиций при формировании батчей
exactlyOnce=false. Этот параметр может повысить пропускную способность за счёт формирования более крупных батчей, но при этом не сохраняются гарантии порядка в пределах каждой партиции.
Несколько топиков с высокой пропускной способностью
topic2TableMap для сопоставления топиков с таблицами и сталкиваетесь с узким местом на этапе вставки, из-за чего возникает отставание потребителя, рассмотрите возможность создать отдельный коннектор для каждого топика.
Основная причина в том, что сейчас батчи вставляются в каждую таблицу последовательно.
Рекомендация: Для нескольких топиков с большим объёмом данных разверните отдельный экземпляр коннектора для каждого топика, чтобы максимально повысить параллельную пропускную способность вставки.
Что учитывать при выборе движка таблицы ClickHouse
MergeTree: Лучший вариант для большинства сценариев, обеспечивает баланс между производительностью запросов и вставкиReplicatedMergeTree: Требуется для высокой доступности, добавляет накладные расходы на репликацию*MergeTreeс правильно заданнымORDER BY: Оптимизируйте под свои шаблоны запросов
Пул соединений и тайм-ауты
socket_timeout(по умолчанию: 30000 мс): Максимальное время ожидания при операциях чтенияconnection_timeout(по умолчанию: 10000 мс): Максимальное время на установление соединения
Мониторинг и устранение проблем с производительностью
- Отставание потребителя: используйте инструменты мониторинга Kafka, чтобы отслеживать отставание по каждой партиции
- Метрики коннектора: отслеживайте
receivedRecords,recordProcessingTime,taskProcessingTimeчерез JMX (см. Мониторинг) - Метрики ClickHouse:
system.asynchronous_inserts: отслеживайте использование буфера асинхронной вставкиsystem.parts: отслеживайте число частей, чтобы выявлять проблемы со слияниемsystem.merges: отслеживайте активные слиянияsystem.events: отслеживайтеInsertedRows,InsertedBytes,FailedInsertQuery
Краткая сводка рекомендаций
- Начните с настроек по умолчанию, затем измерьте производительность и настраивайте систему по фактическим результатам
- Предпочитайте более крупные батчи: по возможности ориентируйтесь на 10 000–100 000 строк на одну вставку
- Используйте асинхронную вставку, если отправляете много небольших батчей или работаете в условиях высокого параллелизма
- Всегда используйте
wait_for_async_insert=1с семантикой «ровно один раз» - Масштабируйте горизонтально: увеличивайте
tasks.maxвплоть до числа партиций - Используйте отдельный коннектор для каждого топика с большим объемом данных для максимальной пропускной способности
- Постоянно отслеживайте метрики: контролируйте отставание потребителя, количество частей и активность слияния
- Тщательно тестируйте: всегда проверяйте изменения конфигурации под реалистичной нагрузкой перед развертыванием в production
Пример: конфигурация для высокой пропускной способности
Приведённая выше конфигурация коннектора требует включить переопределение клиентских настроек в конфигурации воркера через
connector.client.config.override.policy=All. Подробнее см. в документации Kafka Connect.- Обрабатывает до 10 000 записей за один цикл опроса
- Формирует батчи по нескольким партициям для более крупных операций вставки
- Использует асинхронную вставку с буфером 16 МБ
- Запускает 8 параллельных задач (число должно соответствовать количеству ваших партиций)
- Оптимизирована под пропускную способность, а не под строгий порядок
Устранение неполадок
”Несоответствие состояния для топика [someTopic] и партиции [0]”
Эта корректировка может повлиять на гарантии exactly-once.
”При каких ошибках коннектор будет выполнять повторную попытку?”
ClickHouseException- Это общее исключение, которое может сгенерировать ClickHouse. Обычно оно возникает, когда сервер перегружен, и следующие коды ошибок считаются особенно характерными для временных сбоев:- 3 - UNEXPECTED_END_OF_FILE
- 107 - FILE_DOESNT_EXIST
- 159 - TIMEOUT_EXCEEDED
- 164 - READONLY
- 202 - TOO_MANY_SIMULTANEOUS_QUERIES
- 203 - NO_FREE_CONNECTION
- 209 - SOCKET_TIMEOUT
- 210 - NETWORK_ERROR
- 241 - MEMORY_LIMIT_EXCEEDED
- 242 - TABLE_IS_READ_ONLY
- 252 - TOO_MANY_PARTS
- 285 - TOO_FEW_LIVE_REPLICAS
- 319 - UNKNOWN_STATUS_OF_INSERT
- 425 - SYSTEM_ERROR
- 999 - KEEPER_EXCEPTION
SocketTimeoutException- Это исключение возникает при тайм-ауте сокета.UnknownHostException- Это исключение возникает, когда не удаётся определить хост.IOException- Это исключение возникает при проблемах с сетью.
”Все мои данные пустые/состоят из нулей”
_ в качестве разделителя). Поля в таблице тогда будут иметь формат “field1_field2_field3” (то есть “before_id”, “after_id” и т. д.).
”Я хочу использовать ключи Kafka в ClickHouse”
KeyToValue, чтобы переместить ключ в поле value (под новым именем _key):