Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Приёмник ClickHouse Kafka Connect

Приёмник ClickHouse Kafka Connect — это коннектор Kafka, который передаёт данные из топика Kafka в таблицу ClickHouse.

Лицензия

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

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

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

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

Версия ClickHouse Kafka Connect Версия ClickHouse Kafka Connect Confluent Platform Java-клиент ClickHouse
1.4.x (последняя) 23.3+ 2.7+ 6.1+ 0.9.5
1.3.x 23.3+ 2.7+ 6.1+ от 0.8.0 до 0.9.5
1.2.x 23.3+ 2.7+ 6.1+ от 0.6.3 до 0.8.0
1.1.x 23.3+ 2.7+ 6.1+ от 0.6.0 до 0.6.1
1.0.x 23.3+ 2.7+ 6.1+ от 0.4.6 до 0.6.0

Коннектор не запускается с серверами ClickHouse версий ниже 23.3. Всегда используйте последний релиз, если нет причин закреплять более старую версию.

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

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

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

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

Чтобы подключиться к ClickHouse по HTTP(S), вам понадобится следующая информация:

Параметр(ы) Описание
HOST and PORT Обычно используется порт 8443 при использовании TLS и 8123 без TLS.
DATABASE NAME По умолчанию есть база данных default; используйте имя базы данных, к которой хотите подключиться.
USERNAME and PASSWORD По умолчанию имя пользователя — default. Используйте имя пользователя, подходящее для вашего сценария использования.

Сведения о подключении для вашего сервиса 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 и пароль.
connector.class=com.clickhouse.kafka.connect.ClickHouseSinkConnector
tasks.max=1
topics=<topic_name>
ssl=true
jdbcConnectionProperties=?sslmode=STRICT
security.protocol=SSL
hostname=<hostname>
database=<database_name>
password=<password>
ssl.truststore.location=/tmp/kafka.client.truststore.jks
port=8443
value.converter.schemas.enable=false
value.converter=org.apache.kafka.connect.json.JsonConverter
exactlyOnce=true
username=default
schemas.enable=false
  • Перезапустите Confluent Platform.
  • Если вы используете Confluent Platform, войдите в интерфейс Confluent Control Center и убедитесь, что ClickHouse Sink доступен в списке доступных коннекторов.

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

Чтобы подключить ClickHouse Sink к серверу ClickHouse, необходимо указать:

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

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

Название свойства Описание Значение по умолчанию
hostname (обязательно) Имя хоста или IP-адрес сервера Н/Д
port Порт ClickHouse: по умолчанию — 8443 (для HTTPS в облаке), а для HTTP (по умолчанию в самоуправляемом развертывании) следует использовать 8123 8443
ssl Включает SSL-подключение к ClickHouse true
jdbcConnectionProperties Свойства подключения к ClickHouse. Должны начинаться с ?, а параметры param=value должны разделяться символом & ""
username Имя пользователя базы данных ClickHouse default
password (обязательно) Пароль базы данных ClickHouse Н/Д
database Имя базы данных ClickHouse default
connector.class (Обязательно) Класс коннектора (задайте явно и оставьте значением по умолчанию) "com.clickhouse.kafka.connect.ClickHouseSinkConnector"
tasks.max Количество задач коннектора "1"
errors.retry.timeout Максимальная длительность повторных попыток Kafka Connect в миллисекундах. 0 — без повторов. -1 — бесконечные повторы. Рекомендуемое значение — более "10000" мс (10 секунд) Тайм-аут "0"
exactlyOnce Режим exactly-once включен "false"
topics (Обязательный) Топики Kafka для опроса — имена топиков должны совпадать с именами таблиц ""
key.converter (Требуется* — см. описание) Установите в соответствии с типами ваших ключей. Здесь обязательно, если вы передаёте ключи (и они не определены в конфигурации воркера). "org.apache.kafka.connect.storage.StringConverter"
value.converter (Обязательно* — см. описание) Установите в зависимости от типа данных в вашем topic. Поддерживаются форматы JSON, String, Avro или Protobuf. Здесь обязательно, если не задано в конфигурации воркера. "org.apache.kafka.connect.json.JsonConverter"
value.converter.schemas.enable Поддержка схемы конвертером значений коннектора "false"
errors.tolerance Допуск ошибок коннектора. Поддерживаются: none, all "none"
errors.deadletterqueue.topic.name Если параметр задан (при errors.tolerance=all), для неудачных батчей будет использоваться DLQ (см. Устранение неполадок) ""
errors.deadletterqueue.context.headers.enable Добавляет дополнительные заголовки для DLQ ""
clickhouseSettings Список настроек ClickHouse, разделённых запятыми (например, "insert_quorum=2, etc…") ""
topic2TableMap Список, разделённый запятыми, сопоставляющий имена topic с именами таблиц (например, "topic1=table1, topic2=table2, etc…") ""
tableRefreshInterval Время (в секундах) для обновления кэша определения таблицы 0
keeperOnCluster Позволяет задать параметр ON CLUSTER для самоуправляемых экземпляров (например, ON CLUSTER clusterNameInConfigFileDefinition) для таблицы connect_state с гарантией exactly-once (см. Распределённые DDL-запросы ""
bypassRowBinary Позволяет отключить использование RowBinary и RowBinaryWithDefaults для данных на основе схемы (Avro, Protobuf и т. д.) — следует использовать только в случаях, когда в данных могут отсутствовать столбцы, а Nullable/Default недопустимы "false"
dateTimeFormats Форматы даты-времени для разбора полей схемы DateTime64, разделённые символом ; (например, someDateField=yyyy-MM-dd HH:mm:ss.SSSSSSSSS;someOtherDateField=yyyy-MM-dd HH:mm:ss). ""
tolerateStateMismatch Позволяет коннектору отбрасывать записи с offset, "меньшим", чем текущее смещение, сохранённое в AFTER_PROCESSING (например, если отправляется смещение 5, а последним записанным смещением было 250). Используйте для восстановления ингестии после сбоя, а после завершения верните значение "false". "false"
ignorePartitionsWhenBatching Будет игнорировать партицию при сборе сообщений для вставки (но только если exactlyOnce имеет значение false). Примечание о производительности: чем больше задач у коннектора, тем меньше партиций Kafka назначается каждой задаче — в какой-то момент это может давать все меньший эффект. "false"
bufferCount (с версии v1.3.6) Количество записей, буферизуемых в памяти перед записью в ClickHouse. 0 отключает внутреннюю буферизацию. Буферизация не поддерживается при exactlyOnce=true. "0"
bufferFlushTime (с версии v1.3.6) Максимальное время в миллисекундах, в течение которого записи хранятся в буфере перед сбросом, когда exactlyOnce=false. 0 отключает сброс по времени. Значение по умолчанию — 0. Требуется только для порога по времени. Действует только при bufferCount > 0. "0"
reportInsertedOffsets (с версии v1.3.6) Включает возврат из preCommit только смещений, успешно записанных при вставке (вместо currentOffsets), когда exactlyOnce=false. Это не действует, если ignorePartitionsWhenBatching=true: в этом случае по-прежнему возвращаются currentOffsets. "false"
client_version (с версии v1.2.0) Выбирает Java-клиент ClickHouse, используемый коннектором. Допустимые значения: "V1" и "V2". "V1"

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

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

Для каждого топика требуется отдельная целевая таблица в ClickHouse. Имя целевой таблицы должно совпадать с именем исходного топика.

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

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

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

При объявленной схеме:

Тип Kafka Connect Тип ClickHouse Поддерживается Primitive
STRING String Да
STRING JSON. См. ниже (1) Да
INT8 Int8 Да
INT16 Int16 Да
INT32 Int32 Да
INT64 Int64 Да
FLOAT32 Float32 Да
FLOAT64 Float64 Да
BOOLEAN Boolean Да
ARRAY Array(T) Нет
MAP Map(Primitive, T) Нет
STRUCT Variant(T1, T2, …) Нет
STRUCT Tuple(a T1, b T2, …) Нет
STRUCT Nested(a T1, b T2, …) Нет
STRUCT JSON. См. ниже (1), (2) Нет
BYTES String Нет
org.apache.kafka.connect.data.Time Int64 / DateTime64 Нет
org.apache.kafka.connect.data.Timestamp Int32 / Date32 Нет
org.apache.kafka.connect.data.Decimal Decimal Нет
  • (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 без схемы.

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    "tasks.max": "1",
    "consumer.override.max.poll.records": "5000",
    "consumer.override.max.partition.fetch.bytes": "5242880",
    "database": "default",
    "errors.retry.timeout": "60000",
    "exactlyOnce": "false",
    "hostname": "localhost",
    "port": "8443",
    "ssl": "true",
    "jdbcConnectionProperties": "?ssl=true&sslmode=strict",
    "username": "default",
    "password": "<PASSWORD>",
    "topics": "<TOPIC_NAME>",
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    "clickhouseSettings": ""
  }
}

Базовая конфигурация для нескольких топиков

Коннектор может читать данные из нескольких топиков

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "topics": "SAMPLE_TOPIC, ANOTHER_TOPIC, YET_ANOTHER_TOPIC",
    ...
  }
}

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

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "errors.tolerance": "all",
    "errors.deadletterqueue.topic.name": "<DLQ_TOPIC>",
    "errors.deadletterqueue.context.headers.enable": "true",
  }
}

Поддержка схемы Avro

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "io.confluent.connect.avro.AvroConverter",
    "value.converter.schema.registry.url": "<SCHEMA_REGISTRY_HOST>:<PORT>",
    "value.converter.schemas.enable": "true",
  }
}

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

Приведённое ниже сопоставление типов задаётся в io.confluent.connect.avro.AvroConverter — официальной реализации сериализатора/десериализатора Avro для Kafka Connect. Подробную информацию о логике преобразования см. в документации Kafka Connect.

✅: Поддерживается

❌: Не поддерживается

️⚠️: Поддерживается частично

Тип Avro Тип Kafka Connect Поддерживается Примечания
null N/A Не поддерживается как самостоятельный тип, но может использоваться в union
boolean BOOLEAN
int INT8/INT16/INT32 По умолчанию используется INT32. Преобразуется в INT8, если схема содержит свойство connect.type=int8 (аналогично для INT16, если connect.type=int16)
long INT64
float FLOAT32
double FLOAT64
bytes BYTES
string STRING
record STRUCT
enum STRING
array ARRAY/MAP По умолчанию используется ARRAY. Преобразуется в MAP, если поле изначально было создано через AvroData.fromConnectSchema (source)
map MAP
union STRUCT/<T> ⚠️ По умолчанию используется STRUCT. Преобразуется в тип singleton T в определении union, если flatten.singleton.unions=true (см. docs)
fixed BYTES ⚠️ Логический тип decimal для fixed не поддерживается (см. ниже)

Сопоставление типов Kafka Connect и типов ClickHouse см. в разделе Поддерживаемые типы данных.

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

Следующие схемы Avro не поддерживаются коннектором:

  • логический тип decimal для типа fixed
{"name": "decimal_18_4", "type": "fixed", "size": 8, "logicalType": "decimal", "precision": 18, "scale": 4}
  • union-типы с Nullable
{"name": "mixed_union", "type": ["null", "string", "int"], "default": null}
  • объединения в записях
{
  "name": "record_union",
  "type": [
    {
      "type": "record",
      "name": "TypeA",
      "fields": [
        {
          "name": "label",
          "type": "string"
        }
      ]
    },
    {
      "type": "record",
      "name": "TypeB",
      "fields": [
        {
          "name": "count",
          "type": "int"
        }
      ]
    }
  ]
}

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

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "io.confluent.connect.protobuf.ProtobufConverter",
    "value.converter.schema.registry.url": "<SCHEMA_REGISTRY_HOST>:<PORT>",
    "value.converter.schemas.enable": "true",
  }
}

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

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

Ниже приведено сопоставление типов, определенное в io.confluent.connect.protobuf.ProtobufConverter — официальной реализации сериализатора/десериализатора Protobuf для Kafka Connect. Подробную информацию о логике преобразования см. в документации Kafka Connect.

✅: Поддерживается

❌: Не поддерживается

️⚠️: Поддерживается частично

Тип Protobuf Тип Kafka Connect Тип ClickHouse Поддерживается Примечания
double FLOAT64 Float64
float FLOAT32 Float32
int32 INT8/INT16/INT32 Int32 По умолчанию используется INT32. Преобразуется в INT8, если в схеме задан параметр connect.type=int8 (аналогично для INT16 при connect.type=int16)
sint32 INT8/INT16/INT32 Int32 По умолчанию используется INT32. Преобразуется в INT8, если в схеме задан параметр connect.type=int8 (аналогично для INT16 при connect.type=int16)
sfixed32 INT8/INT16/INT32 Int32 По умолчанию используется INT32. Преобразуется в INT8, если в схеме задан параметр connect.type=int8 (аналогично для INT16 при connect.type=int16)
uint32 INT64 UInt32
fixed32 INT64 UInt32
int64 INT64 Int64
uint64 INT64 UInt64
sint64 INT64 Int64
fixed64 INT64 UInt64
sfixed64 INT64 Int64
bool BOOLEAN Bool
string STRING String
bytes BYTES String
enum INT32/STRING Int32 По умолчанию используется STRING. Преобразуется в INT32, если int.for.enums=true (см. документацию schema registry)
message STRUCT Tuple / JSON ⚠️ См. раздел ниже о неподдерживаемых схемах
repeated T (where T is not a map entry) ARRAY Array(T)
map<K, V> MAP Map(K, V)
oneof STRUCT Tuple / Variant ⚠️ См. раздел ниже о преобразовании oneof в схему ClickHouse
google.protobuf.DoubleValue FLOAT64 Nullable(Float64)
google.protobuf.FloatValue FLOAT32 Nullable(Float32)
google.protobuf.Int64Value INT64 Nullable(Int64)
google.protobuf.UInt64Value INT64 Nullable(UInt64)
google.protobuf.UInt32Value INT64 Nullable(UInt32)
google.protobuf.Int32Value INT32 Nullable(Int32)
google.protobuf.BoolValue BOOLEAN Nullable(Bool)
google.protobuf.StringValue STRING Nullable(String)
google.protobuf.BytesValue BYTES Nullable(String)
google.protobuf.Timestamp org.apache.kafka.connect.data.Timestamp DateTime64(3)
google.type.Date org.apache.kafka.connect.data.Date Date
google.type.TimeOfDay org.apache.kafka.connect.data.Time Int32 / Int64
google.protobuf.Duration STRUCT Tuple(seconds Int64, nano Nullable(Int32))
google.protobuf.Any N/A N/A
google.protobuf.Empty N/A N/A

Сопоставление типов между Kafka Connect и ClickHouse см. в разделе Поддерживаемые типы данных.

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

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

Например:

syntax = "proto3";

package com.clickhouse.kafka.connect.proto.test;

message StringIntUnion {
  oneof mixed {
    string mixed_string = 2;
    int32 mixed_int = 3;
  }
}

соответствует следующему определению таблицы ClickHouse:

CREATE TABLE IF NOT EXISTS `StringIntUnion`
(
    mixed_string Nullable(String),
    mixed_int Nullable(Int32)
) ENGINE = ...;

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

Следующие схемы Protobuf не поддерживаются коннектором:

  • union-типы из нескольких сообщений (до версии CH 26.1)
syntax = "proto3";

package com.clickhouse.kafka.connect.proto.test;

message TwoRecords {
  oneof payload {
    TypeA type_a = 2;
    TypeB type_b = 3;
  }

  // translates to Nullable(Tuple(label String)) in ClickHouse, which is unsupported
  message TypeA {
    string label = 1;
  }

  // translates to Nullable(Tuple(count Int32)) in ClickHouse, which is unsupported
  message TypeB {
    int32 count = 1;
  }
}

Начиная с версии CH 26.1 эта схема поддерживается при allow_experimental_nullable_tuple_type=1 (см. эту страницу документации).

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

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
  }
}

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

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

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "value.converter": "org.apache.kafka.connect.storage.StringConverter",
    "customInsertFormat": "true",
    "insertFormat": "CSV"
  }
}

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

Внутренняя буферизация позволяет задаче sink-коннектора накапливать записи из нескольких вызовов poll() и сбрасывать их в ClickHouse более крупными батчами. Это может повысить пропускную способность в рабочих нагрузках, где каждый опрос возвращает множество небольших батчей по отдельным партициям.

Ключевые особенности:

  • bufferCount определяет, сколько записей буферизуется перед сбросом.
  • bufferFlushTime задает максимальное время ожидания (в миллисекундах) перед сбросом буферизованных записей.
  • bufferFlushTime действует только при bufferCount > 0.
  • bufferCount=0 и bufferFlushTime=0 оставляют буферизацию отключенной (поведение по умолчанию).
  • Буферизация не поддерживается, если exactlyOnce=true.

Почему буферизация несовместима с режимом exactly-once: Буферизация изменяет границы батчей, из-за чего нарушаются дедупликация блоков ClickHouse и работа автомата состояний offset'ов коннектора. Чтобы решить эту проблему, либо отключите режим exactly-once, задав exactlyOnce=false в конфигурации коннектора, либо отключите буферизацию, задав bufferCount=0.

Пример:

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "exactlyOnce": "false",
    "bufferCount": "5000",
    "bufferFlushTime": "2000"
  }
}

Логирование

Логирование автоматически обеспечивается платформой Kafka Connect. Пункт назначения и формат журналов можно настроить через файл конфигурации Kafka Connect.

Если вы используете Confluent Platform, журналы можно просмотреть, выполнив CLI-команду:

confluent local services connect log

Подробнее см. в официальном руководстве.

Мониторинг

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

Метрики ClickHouse

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

com.clickhouse:type=ClickHouseKafkaConnector,name=SinkTask{id}
Имя метрики Тип Описание
receivedRecords long Общее количество полученных записей.
recordProcessingTime long Общее время в наносекундах, затраченное на группировку и преобразование записей в единую структуру.
taskProcessingTime long Общее время в наносекундах, затраченное на обработку и вставку данных в ClickHouse.

Метрики 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 мс)
Настройки опроса

Коннектор опрашивает буфер фреймворка на наличие сообщений:

  • max.poll.records: Максимальное количество записей, возвращаемых за один опрос (по умолчанию: 500)
  • max.partition.fetch.bytes: Максимальный объём данных на партицию (по умолчанию: 1048576 / 1 MB)

Для оптимальной производительности в ClickHouse используйте более крупные батчи:

# Увеличить количество записей за один опрос
consumer.override.max.poll.records=5000

# Увеличить размер считываемых данных для партиции (5 МБ)
consumer.override.max.partition.fetch.bytes=5242880

# Опционально: увеличить минимальный объём считываемых данных для накопления перед обработкой (1 МБ)
consumer.override.fetch.min.bytes=1048576

# Опционально: уменьшить время ожидания, если задержка критична
consumer.override.fetch.max.wait.ms=300

Важно: настройки 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:

{
  "name": "clickhouse-connect",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    ...
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1"
  }
}

Ключевые настройки:

  • async_insert=1: Включает асинхронную вставку
  • wait_for_async_insert=1 (рекомендуется): Коннектор ожидает, пока данные будут записаны в хранилище ClickHouse, прежде чем подтверждать получение. Обеспечивает гарантии доставки.
  • wait_for_async_insert=0: Коннектор подтверждает получение сразу после буферизации. Обеспечивает более высокую производительность, но данные могут быть потеряны при сбое сервера до сброса на диск.
Тонкая настройка поведения асинхронной вставки

Вы можете тонко настроить поведение сброса при асинхронной вставке:

"clickhouseSettings": "async_insert=1,wait_for_async_insert=1,async_insert_max_data_size=104857600,async_insert_busy_timeout_ms=1000"

Распространённые параметры настройки:

  • 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 с асинхронными вставками:

{
  "config": {
    "exactlyOnce": "true",
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1"
  }
}

Важно: Всегда используйте wait_for_async_insert=1 вместе с exactly-once, чтобы смещения фиксировались только после сохранения данных.

Подробнее об async inserts см. в документации ClickHouse об async inserts.

Параллелизм коннектора

Увеличьте степень параллелизма, чтобы повысить пропускную способность:

Количество задач на коннектор
"tasks.max": "4"

Каждая задача обрабатывает подмножество партиций топика. Больше задач = больше параллелизма, но:

  • Максимально эффективное количество задач = количество партиций топика
  • Каждая задача поддерживает собственное подключение к ClickHouse
  • Больше задач = выше накладные расходы и выше вероятность конкуренции за ресурсы

Рекомендация: Начните с tasks.max, равного количеству партиций топика, затем корректируйте это значение на основе метрик CPU и пропускной способности.

Игнорирование партиций при формировании батчей

По умолчанию коннектор формирует батчи сообщений по каждой партиции. Чтобы повысить пропускную способность, можно формировать батчи без учёта партиций:

"ignorePartitionsWhenBatching": "true"

** Предупреждение**: Используйте только при exactlyOnce=false. Этот параметр может повысить пропускную способность за счёт формирования более крупных батчей, но при этом не сохраняются гарантии порядка в пределах каждой партиции.

Несколько топиков с высокой пропускной способностью

Если ваш коннектор настроен на подписку на несколько топиков, вы используете topic2TableMap для сопоставления топиков с таблицами и сталкиваетесь с узким местом на этапе вставки, из-за чего возникает отставание потребителя, рассмотрите возможность создать отдельный коннектор для каждого топика.

Основная причина в том, что сейчас батчи вставляются в каждую таблицу последовательно.

Рекомендация: Для нескольких топиков с большим объёмом данных разверните отдельный экземпляр коннектора для каждого топика, чтобы максимально повысить параллельную пропускную способность вставки.

Что учитывать при выборе движка таблицы ClickHouse

Выберите подходящий движок таблицы ClickHouse для своей задачи:

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

Настройки, которые следует учитывать:

CREATE TABLE my_table (...)
ENGINE = MergeTree()
ORDER BY (timestamp, id)
SETTINGS 
    -- Увеличить максимальное количество потоков вставки для параллельной записи частей
    max_insert_threads = 4,
    -- Разрешить вставки с кворумом для надёжности (ReplicatedMergeTree)
    insert_quorum = 2

Для настроек вставки на уровне коннектора:

"clickhouseSettings": "insert_quorum=2,insert_quorum_timeout=60000"

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

Коннектор использует HTTP-соединения с ClickHouse. Настройте тайм-ауты для сетей с высокой задержкой:

"jdbcConnectionProperties": "?socket_timeout=300000,connection_timeout=30000"
  • socket_timeout (по умолчанию: 30000 мс): Максимальное время ожидания при операциях чтения
  • connection_timeout (по умолчанию: 10000 мс): Максимальное время на установление соединения

Увеличьте эти значения, если при работе с большими батчами возникают ошибки тайм-аута.

Сетевой буфер клиента

При использовании клиента V2 ("client_version": "V2") данные копируются между сокетом и памятью приложения через буфер, размер которого задаётся параметром client_network_buffer_size (в байтах). По умолчанию его размер составляет 300000 байт. Увеличение размера буфера может повысить пропускную способность в соединениях с высокой задержкой или высокой пропускной способностью, но увеличивает использование памяти на каждую задачу подключения.

Все параметры клиента настраиваются через jdbcConnectionProperties, а не через clickhouseSettings. Например:

"jdbcConnectionProperties": "?client_network_buffer_size=1000000"

Если необходимо задать несколько свойств клиента, разделите их символом &:

"jdbcConnectionProperties": "?ssl=true&client_network_buffer_size=1000000"

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

Отслеживайте следующие ключевые метрики:

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

Распространенные проблемы с производительностью:

Симптом Возможная причина Решение
Высокое отставание потребителя Слишком маленькие батчи Увеличьте max.poll.records, включите асинхронную вставку
Ошибки "Too many parts" Частые мелкие вставки Включите асинхронную вставку, увеличьте размер батча
Ошибки тайм-аута Слишком большой размер батча, медленная сеть Уменьшите размер батча, увеличьте socket_timeout, проверьте сеть
Высокая загрузка CPU Слишком много мелких частей Включите асинхронную вставку, увеличьте параметры слияния
Ошибки OutOfMemory Слишком большой размер батча Уменьшите max.poll.records, max.partition.fetch.bytes
Неравномерная загрузка задач Неравномерное распределение партиций Выполните перебалансировку или скорректируйте tasks.max

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

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

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

Ниже приведён полный пример, оптимизированный для высокой пропускной способности:

{
  "name": "clickhouse-high-throughput",
  "config": {
    "connector.class": "com.clickhouse.kafka.connect.ClickHouseSinkConnector",
    "tasks.max": "8",
    
    "topics": "high_volume_topic",
    "hostname": "my-clickhouse-host.cloud",
    "port": "8443",
    "database": "default",
    "username": "default",
    "password": "<PASSWORD>",
    "ssl": "true",
    
    "value.converter": "org.apache.kafka.connect.json.JsonConverter",
    "value.converter.schemas.enable": "false",
    
    "exactlyOnce": "false",
    "ignorePartitionsWhenBatching": "true",
    
    "consumer.override.max.poll.records": "10000",
    "consumer.override.max.partition.fetch.bytes": "5242880",
    "consumer.override.fetch.min.bytes": "1048576",
    "consumer.override.fetch.max.wait.ms": "500",
    
    "clickhouseSettings": "async_insert=1,wait_for_async_insert=1,async_insert_max_data_size=16777216,async_insert_busy_timeout_ms=1000,socket_timeout=300000"
  }
}

Эта конфигурация:

  • Обрабатывает до 10 000 записей за один цикл опроса
  • Формирует батчи по нескольким партициям для более крупных операций вставки
  • Использует асинхронную вставку с буфером 16 МБ
  • Запускает 8 параллельных задач (число должно соответствовать количеству ваших партиций)
  • Оптимизирована под пропускную способность, а не под строгий порядок

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

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

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

-- Сначала определите базу данных, используемую для хранения данных.
SELECT * FROM [database].connect_state

-- Определите ключ, соответствующий топику и партиции.
ALTER TABLE [database].connect_state DELETE WHERE key = [keyname]

"При каких ошибках коннектор будет выполнять повторную попытку?"

Сейчас основное внимание уделяется выявлению временных ошибок, при которых можно выполнить повторную попытку, в том числе:

  • 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 в конфигурацию коннектора:

transforms=flatten
transforms.flatten.type=org.apache.kafka.connect.transforms.Flatten$Value
transforms.flatten.delimiter=_

Это преобразует ваши данные из вложенного JSON в плоский JSON (с использованием _ в качестве разделителя). Поля в таблице тогда будут иметь формат "field1_field2_field3" (то есть "before_id", "after_id" и т. д.).

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

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

transforms=keyToValue
transforms.keyToValue.type=com.clickhouse.kafka.connect.transforms.KeyToValue
transforms.keyToValue.field=_key
Navigation