Skip to main content
Если вам нужна помощь, пожалуйста, сообщите о проблеме в репозитории или задайте вопрос в публичном Slack ClickHouse.
Приёмник ClickHouse Kafka Connect — это коннектор Kafka, который передаёт данные из топика Kafka в таблицу ClickHouse.

Лицензия

Коннектор Kafka Connector Sink распространяется по лицензии Apache 2.0

Требования к окружению

В окружении должен быть установлен фреймворк Kafka Connect версии 2.7 или более поздней.

Матрица совместимости версий

Основные возможности

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

Инструкции по установке

Подготовьте сведения о подключении

Чтобы подключиться к ClickHouse по HTTP(S), вам понадобится следующая информация: Сведения о подключении для вашего сервиса ClickHouse Cloud доступны в консоли ClickHouse Cloud. Выберите сервис и нажмите Connect:
Кнопка подключения сервиса ClickHouse Cloud
Выберите HTTPS. Сведения о подключении будут показаны в примере команды curl.
Сведения о подключении к ClickHouse Cloud по HTTPS
Если вы используете самоуправляемый ClickHouse, сведения о подключении задаёт ваш администратор ClickHouse.

Общие инструкции по установке

Коннектор поставляется в виде одного JAR-файла, содержащего все файлы классов, необходимые для запуска плагина. Чтобы установить плагин, выполните следующие шаги:
  • Загрузите ZIP-архив с JAR-файлом коннектора со страницы Releases репозитория Приёмник ClickHouse Kafka Connect.
  • Извлеките содержимое ZIP-файла и скопируйте его в нужное место.
  • Добавьте путь к каталогу плагина в параметр plugin.path в файле свойств Connect, чтобы Confluent Platform могла найти плагин.
  • Укажите в конфигурации имя топика, имя хоста экземпляра ClickHouse и пароль.
  • Перезапустите Confluent Platform.
  • Если вы используете Confluent Platform, войдите в интерфейс Confluent Control Center и убедитесь, что ClickHouse Sink доступен в списке доступных коннекторов.

Параметры конфигурации

Чтобы подключить ClickHouse Sink к серверу ClickHouse, необходимо указать:
  • сведения о подключении: hostname (обязательно) и порт (необязательно)
  • учетные данные пользователя: пароль (обязательно) и имя пользователя (необязательно)
  • класс коннектора: com.clickhouse.kafka.connect.ClickHouseSinkConnector (обязательно)
  • topics или topics.regex: топики Kafka, которые нужно опрашивать; имена топиков должны совпадать с именами таблиц (обязательно)
  • конвертеры key и value: указываются в зависимости от типа данных в вашем топике. Обязательно, если они еще не заданы в конфигурации воркера.
Полная таблица параметров конфигурации:

Целевые таблицы

ClickHouse Connect Sink читает сообщения из топиков Kafka и записывает их в соответствующие таблицы. ClickHouse Connect Sink записывает данные в существующие таблицы. Пожалуйста, убедитесь, что перед началом вставки данных в ClickHouse уже создана целевая таблица с подходящей схемой. Для каждого топика требуется отдельная целевая таблица в ClickHouse. Имя целевой таблицы должно совпадать с именем исходного топика.

Предобработка

Если вам нужно преобразовать исходящие сообщения перед отправкой в Приёмник ClickHouse Kafka Connect, используйте преобразования Kafka Connect.

Поддерживаемые типы данных

При объявленной схеме:
  • (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.
Без объявленной схемы: Запись преобразуется в JSON и отправляется в ClickHouse как значение в формате JSONEachRow.

Варианты конфигурации

Ниже приведены несколько распространённых вариантов конфигурации, которые помогут вам быстро начать работу.

Базовая конфигурация

Простейшая конфигурация для начала работы — предполагается, что Kafka Connect запущен в распределенном режиме, а сервер ClickHouse работает на 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

Следующие схемы Avro не поддерживаются коннектором:
  • логический тип decimal для типа fixed
  • union-типы с Nullable
  • объединения в записях

Поддержка схем Protobuf

Обратите внимание: если вы столкнулись с проблемами из-за отсутствующих классов, имейте в виду, что не во всех средах есть конвертер Protobuf, и вам может потребоваться альтернативная версия JAR-файла, включающая зависимости.

Сопоставление типов Protobuf

Ниже приведено сопоставление типов, определенное в io.confluent.connect.protobuf.ProtobufConverter — официальной реализации сериализатора/десериализатора Protobuf для Kafka Connect. Подробную информацию о логике преобразования см. в документации Kafka Connect. ✅: Поддерживается ❌: Не поддерживается ️⚠️: Поддерживается частично Сопоставление типов между Kafka Connect и ClickHouse см. в разделе Поддерживаемые типы данных.

Примечание о сопоставлении полей oneof со столбцами ClickHouse

Коннектор не поддерживает преобразование union Protobuf (oneof) в тип ClickHouse Variant. Вместо этого укажите поля oneof в схеме таблицы ClickHouse как отдельные поля Nullable. Например:
соответствует следующему определению таблицы ClickHouse:

Неподдерживаемые схемы Protobuf

Следующие схемы Protobuf не поддерживаются коннектором:
  • union-типы из нескольких сообщений (до версии CH 26.1)
Начиная с версии CH 26.1 эта схема поддерживается при allow_experimental_nullable_tuple_type=1 (см. эту страницу документации).

Поддержка схем JSON

Поддержка типа String

Коннектор поддерживает String Converter для различных форматов ClickHouse: JSON, CSV и TSV.

Внутренняя буферизация

Внутренняя буферизация позволяет задаче sink-коннектора накапливать записи из нескольких вызовов poll() и сбрасывать их в ClickHouse более крупными батчами. Это может повысить пропускную способность в рабочих нагрузках, где каждый опрос возвращает множество небольших батчей по отдельным партициям. Ключевые особенности:
  • bufferCount определяет, сколько записей буферизуется перед сбросом.
  • bufferFlushTime задает максимальное время ожидания (в миллисекундах) перед сбросом буферизованных записей.
  • bufferFlushTime действует только при bufferCount > 0.
  • bufferCount=0 и bufferFlushTime=0 оставляют буферизацию отключенной (поведение по умолчанию).
  • Буферизация не поддерживается, если exactlyOnce=true.
Почему буферизация несовместима с режимом exactly-once: Буферизация изменяет границы батчей, из-за чего нарушаются дедупликация блоков ClickHouse и работа автомата состояний offset’ов коннектора. Чтобы решить эту проблему, либо отключите режим exactly-once, задав exactlyOnce=false в конфигурации коннектора, либо отключите буферизацию, задав bufferCount=0. Пример:

Логирование

Логирование автоматически обеспечивается платформой Kafka Connect. Пункт назначения и формат журналов можно настроить через файл конфигурации Kafka Connect. Если вы используете Confluent Platform, журналы можно просмотреть, выполнив CLI-команду:
Подробнее см. в официальном руководстве.

Мониторинг

ClickHouse Kafka Connect публикует метрики времени выполнения через Java Management Extensions (JMX). В Kafka Connector JMX включен по умолчанию.

Метрики ClickHouse

Коннектор публикует пользовательские метрики под следующим именем MBean:

Метрики Kafka Producer/Consumer

Коннектор предоставляет стандартные метрики продюсера и консьюмера Kafka, которые помогают оценивать поток данных, пропускную способность и производительность. Метрики на уровне topic:
  • records-sent-total: Общее количество записей, отправленных в topic
  • bytes-sent-total: Общее количество байт, отправленных в topic
  • record-send-rate: Средняя скорость отправки записей в секунду
  • byte-rate: Среднее количество байт, отправляемых в секунду
  • compression-rate: Достигнутый коэффициент сжатия
Метрики на уровне партиции:
  • records-sent-total: Общее количество записей, отправленных в партицию
  • bytes-sent-total: Общее количество байт, отправленных в партицию
  • records-lag: Текущее отставание в партиции
  • records-lead: Текущее опережение в партиции
  • replica-fetch-lag: Информация об отставании реплик
Метрики соединений на уровне узла:
  • connection-creation-total: Общее количество соединений, установленных с узлом Kafka
  • connection-close-total: Общее количество закрытых соединений
  • request-total: Общее количество запросов, отправленных узлу
  • response-total: Общее количество ответов, полученных от узла
  • request-rate: Средняя частота запросов в секунду
  • response-rate: Средняя частота ответов в секунду
Эти метрики помогают отслеживать:
  • Пропускную способность: отслеживать скорость приёма данных
  • Отставание: выявлять узкие места и задержки обработки
  • Сжатие: оценивать эффективность сжатия данных
  • Состояние соединений: контролировать сетевую связность и стабильность

Метрики фреймворка Kafka Connect

Коннектор интегрируется с фреймворком 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: Количество неудачных записей в DLQ
  • deadletterqueue-produce-requests: Общее число попыток записи в DLQ
  • last-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: Общее количество записей, полученных при опросе

Рекомендации по мониторингу

  1. Отслеживайте отставание потребителя: Следите за records-lag по каждой партиции, чтобы выявлять узкие места в обработке
  2. Контролируйте уровень ошибок: Следите за errors-total и records-skip-total, чтобы выявлять проблемы с качеством данных
  3. Следите за состоянием задач: Отслеживайте метрики состояния задач, чтобы убедиться, что они работают корректно
  4. Измеряйте пропускную способность: Используйте records-send-rate и byte-rate для отслеживания производительности ингестии
  5. Следите за состоянием соединений: Проверяйте метрики соединений на уровне узла, чтобы выявлять проблемы в сети
  6. Отслеживайте эффективность сжатия: Используйте compression-rate для оптимизации передачи данных
Подробные определения JMX-метрик и сведения об интеграции с Prometheus см. в файле конфигурации jmx-export-connector.yml.

Ограничения

  • Удаление данных не поддерживается.
  • Размер батча наследуется из свойств Kafka Consumer.
  • При использовании KeeperMap для exactly-once, если смещение было изменено или сброшено назад, необходимо удалить содержимое KeeperMap для соответствующего топика. (Подробнее см. в руководстве по устранению неполадок ниже)

Настройка производительности и оптимизация пропускной способности

В этом разделе описаны стратегии настройки производительности для Приёмника ClickHouse Kafka Connect. Настройка производительности особенно важна при работе со сценариями с высокой пропускной способностью, а также когда нужно оптимизировать использование ресурсов и минимизировать задержку.

Когда нужна настройка производительности?

Настройка производительности обычно требуется в следующих случаях:
  • Высоконагруженные рабочие нагрузки: когда обрабатываются миллионы событий в секунду из топиков Kafka
  • Отставание потребителя: когда ваш коннектор не успевает за скоростью поступления данных, и отставание растет
  • Ограниченные ресурсы: когда нужно оптимизировать использование CPU, памяти или сети
  • Несколько топиков: когда данные одновременно потребляются из нескольких топиков с большим объемом трафика
  • Маленький размер сообщений: когда приходится работать со множеством небольших сообщений, для которых полезен батчинг на стороне сервера
Настройка производительности обычно НЕ нужна, когда:
  • Вы обрабатываете небольшие или умеренные объемы данных (< 10,000 сообщений/секунду)
  • Отставание потребителя остается стабильным и приемлемым для вашего сценария использования
  • Настройки коннектора по умолчанию уже соответствуют вашим требованиям к пропускной способности
  • Ваш кластер ClickHouse без труда справляется с входящей нагрузкой

Как устроен поток данных

Прежде чем приступать к настройке, важно понимать, как данные проходят через коннектор:
  1. Фреймворк Kafka Connect в фоновом режиме забирает сообщения из топиков Kafka
  2. Коннектор считывает сообщения из внутреннего буфера фреймворка
  3. Коннектор объединяет сообщения в батчи в зависимости от размера выборки
  4. ClickHouse получает батч вставки по HTTP/S
  5. ClickHouse обрабатывает вставку (синхронно или асинхронно)
Производительность можно оптимизировать на каждом из этих этапов.

Настройка размера батча в Kafka Connect

Первый уровень оптимизации — управлять тем, какой объём данных коннектор получает из Kafka за один батч.
Настройки fetch
Kafka Connect (фреймворк) в фоновом режиме получает сообщения из топиков Kafka независимо от коннектора:
  • 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.
Для оптимальной производительности в ClickHouse используйте более крупные батчи:
Чтобы использовать указанные выше свойства, нужно разрешить переопределение клиентских настроек в конфигурации воркера с помощью connector.client.config.override.policy=All. Подробнее см. в документации Kafka Connect.
Важно: настройки fetch в Kafka Connect относятся к сжатым данным, тогда как ClickHouse получает несжатые данные. Подбирайте эти настройки с учётом вашего коэффициента сжатия. Компромиссы:
  • Более крупные батчи = Лучшая производительность ингестии в ClickHouse, меньше частей, ниже накладные расходы
  • Более крупные батчи = Более высокое использование памяти, возможное увеличение сквозной задержки
  • Слишком крупные батчи = Риск тайм-аутов, ошибок OutOfMemory или превышения max.poll.interval.ms
Подробнее: документация Confluent | документация Kafka

Асинхронные вставки

Асинхронные вставки особенно полезны, если коннектор отправляет сравнительно небольшие батчи или если вы хотите дополнительно оптимизировать ингестию, переложив батчинг на ClickHouse.
Когда использовать асинхронную вставку
Рассмотрите возможность включения асинхронной вставки, если:
  • Много небольших батчей: Ваш коннектор часто отправляет небольшие батчи (< 1000 строк в батче)
  • Высокий параллелизм: Несколько задач коннектора записывают данные в одну и ту же таблицу
  • Распределённое развертывание: На разных хостах запущено много экземпляров коннектора
  • Накладные расходы на создание частей: Вы сталкиваетесь с ошибками “too many parts”
  • Смешанная рабочая нагрузка: Вы совмещаете ингестию в реальном времени с нагрузкой от запросов
НЕ используйте асинхронную вставку, если:
  • Вы уже отправляете большие батчи (> 10 000 строк в батче) с контролируемой частотой
  • Вам нужна немедленная доступность данных для запросов (запросы должны видеть данные мгновенно)
  • Семантика «ровно один раз» с wait_for_async_insert=0 не соответствует вашим требованиям
  • В вашем случае вместо этого можно выиграть от улучшения батчинга на стороне клиента
Как работают асинхронные вставки
Когда асинхронные вставки включены, ClickHouse:
  1. Получает запрос на вставку от коннектора
  2. Записывает данные во внутренний буфер в памяти (вместо немедленной записи на диск)
  3. Возвращает коннектору успешный результат (если wait_for_async_insert=0)
  4. Записывает буфер на диск, когда выполняется одно из следующих условий:
    • Буфер достигает async_insert_max_data_size (по умолчанию: 100 MB)
    • С момента первой вставки прошло async_insert_busy_timeout_ms миллисекунд (по умолчанию: 1000 ms)
    • Достигнуто максимальное количество накопленных запросов (async_insert_max_query_number, по умолчанию: 100)
Это значительно уменьшает количество создаваемых частей и повышает общую пропускную способность.
Включение асинхронной вставки
Добавьте настройки async insert в параметр конфигурации 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

Выберите подходящий движок таблицы ClickHouse для своей задачи:
  • MergeTree: Лучший вариант для большинства сценариев, обеспечивает баланс между производительностью запросов и вставки
  • ReplicatedMergeTree: Требуется для высокой доступности, добавляет накладные расходы на репликацию
  • *MergeTree с правильно заданным ORDER BY: Оптимизируйте под свои шаблоны запросов
Настройки, которые следует учитывать:
Для настроек вставки на уровне коннектора:

Пул соединений и тайм-ауты

Коннектор использует HTTP-соединения с ClickHouse. Настройте тайм-ауты для сетей с высокой задержкой:
  • socket_timeout (по умолчанию: 30000 мс): Максимальное время ожидания при операциях чтения
  • connection_timeout (по умолчанию: 10000 мс): Максимальное время на установление соединения
Увеличьте эти значения, если при работе с большими батчами возникают ошибки тайм-аута.

Мониторинг и устранение проблем с производительностью

Отслеживайте следующие ключевые метрики:
  1. Отставание потребителя: используйте инструменты мониторинга Kafka, чтобы отслеживать отставание по каждой партиции
  2. Метрики коннектора: отслеживайте receivedRecords, recordProcessingTime, taskProcessingTime через JMX (см. Мониторинг)
  3. Метрики ClickHouse:
    • system.asynchronous_inserts: отслеживайте использование буфера асинхронной вставки
    • system.parts: отслеживайте число частей, чтобы выявлять проблемы со слиянием
    • system.merges: отслеживайте активные слияния
    • system.events: отслеживайте InsertedRows, InsertedBytes, FailedInsertQuery
Распространенные проблемы с производительностью:

Краткая сводка рекомендаций

  1. Начните с настроек по умолчанию, затем измерьте производительность и настраивайте систему по фактическим результатам
  2. Предпочитайте более крупные батчи: по возможности ориентируйтесь на 10 000–100 000 строк на одну вставку
  3. Используйте асинхронную вставку, если отправляете много небольших батчей или работаете в условиях высокого параллелизма
  4. Всегда используйте wait_for_async_insert=1 с семантикой «ровно один раз»
  5. Масштабируйте горизонтально: увеличивайте tasks.max вплоть до числа партиций
  6. Используйте отдельный коннектор для каждого топика с большим объемом данных для максимальной пропускной способности
  7. Постоянно отслеживайте метрики: контролируйте отставание потребителя, количество частей и активность слияния
  8. Тщательно тестируйте: всегда проверяйте изменения конфигурации под реалистичной нагрузкой перед развертыванием в production

Пример: конфигурация для высокой пропускной способности

Ниже приведён полный пример, оптимизированный для высокой пропускной способности:
Приведённая выше конфигурация коннектора требует включить переопределение клиентских настроек в конфигурации воркера через connector.client.config.override.policy=All. Подробнее см. в документации Kafka Connect.
Эта конфигурация:
  • Обрабатывает до 10 000 записей за один цикл опроса
  • Формирует батчи по нескольким партициям для более крупных операций вставки
  • Использует асинхронную вставку с буфером 16 МБ
  • Запускает 8 параллельных задач (число должно соответствовать количеству ваших партиций)
  • Оптимизирована под пропускную способность, а не под строгий порядок

Устранение неполадок

”Несоответствие состояния для топика [someTopic] и партиции [0]

Это происходит, когда смещение, сохранённое в KeeperMap, отличается от смещения, сохранённого в Kafka — обычно после удаления топика или ручной корректировки смещения. Чтобы это исправить, нужно удалить старые значения, сохранённые для указанной комбинации топика и партиции:
Эта корректировка может повлиять на гарантии 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 - Это исключение возникает при проблемах с сетью.

”Все мои данные пустые/состоят из нулей”

Скорее всего, поля в ваших данных не совпадают с полями в таблице — особенно часто это встречается при использовании CDC (и формата Debezium). Одно из распространённых решений — добавить преобразование flatten в конфигурацию коннектора:
Это преобразует ваши данные из вложенного JSON в плоский JSON (с использованием _ в качестве разделителя). Поля в таблице тогда будут иметь формат “field1_field2_field3” (то есть “before_id”, “after_id” и т. д.).

”Я хочу использовать ключи Kafka в ClickHouse”

По умолчанию ключи Kafka не сохраняются в поле value, но вы можете использовать преобразование KeyToValue, чтобы переместить ключ в поле value (под новым именем _key):
Последнее изменение 23 июля 2026 г.