В наших примерах мы используем дистрибутив Confluent для Kafka Connect.
Ниже описана простая установка, при которой сообщения считываются из одного топика Kafka, а строки вставляются в таблицу ClickHouse. Мы рекомендуем Confluent Cloud, который предлагает достаточно щедрый бесплатный уровень для тех, у кого нет собственной среды Kafka.
Обратите внимание, что для коннектора JDBC требуется схема (с коннектором JDBC нельзя использовать обычные JSON или CSV). Хотя схема может быть закодирована в каждом сообщении, настоятельно рекомендуется использовать Schema Registry от Confluent, чтобы избежать связанных с этим накладных расходов. Предоставленный скрипт вставки автоматически определяет схему по сообщениям и регистрирует её в реестре — поэтому этот скрипт можно повторно использовать и для других датасетов. Предполагается, что ключи Kafka имеют тип String. Более подробную информацию о схемах Kafka можно найти здесь.
Лицензия
Коннектор JDBC распространяется по лицензии Confluent Community
Порядок действий
Подготовьте сведения о подключении
Чтобы подключиться к ClickHouse по HTTP(S), вам понадобится следующая информация:
| Параметр(ы) | Описание |
|---|---|
HOST and PORT |
Обычно используется порт 8443 при использовании TLS и 8123 без TLS. |
DATABASE NAME |
По умолчанию есть база данных default; используйте имя базы данных, к которой хотите подключиться. |
USERNAME and PASSWORD |
По умолчанию имя пользователя — default. Используйте имя пользователя, подходящее для вашего сценария использования. |
Сведения о подключении для вашего сервиса ClickHouse Cloud доступны в консоли ClickHouse Cloud. Выберите сервис и нажмите Connect:

Выберите HTTPS. Сведения о подключении будут показаны в примере команды curl.

Если вы используете самоуправляемый ClickHouse, сведения о подключении задаёт ваш администратор ClickHouse.
Установите Kafka Connect и коннектор
Мы исходим из того, что вы скачали пакет Confluent и установили его локально. Следуйте инструкциям по установке коннектора, приведённым здесь.
Если вы используете метод установки через confluent-hub, ваши локальные файлы конфигурации будут обновлены.
Для отправки данных из Kafka в ClickHouse мы используем компонент Sink коннектора.
Подготовьте конфигурацию
Следуйте этим инструкциям по настройке Connect в соответствии с типом вашей установки, учитывая различия между автономным и распределённым кластером. Если вы используете Confluent Cloud, вам подходит распределённая конфигурация.
Следующие параметры важны при использовании коннектора JDBC с ClickHouse. Полный список параметров приведён здесь:
_connection.url_- должен иметь видjdbc:clickhouse://<clickhouse host>:<clickhouse http port>/<target database>connection.user- пользователь с правами на запись в целевую базу данныхtable.name.format- таблица ClickHouse, в которую выполняется вставка данных. Она должна существовать.batch.size- количество строк, отправляемых в одном батче. Убедитесь, что здесь задано достаточно большое значение. Согласно рекомендациям ClickHouse, значение 1000 следует считать минимумом.tasks.max- коннектор JDBC Sink поддерживает запуск одной или нескольких задач. Это можно использовать для повышения производительности. Наряду с размером батча это ваш основной способ повысить производительность.value.converter.schemas.enable- установите false, если используете Schema Registry, и true, если встраиваете схемы в сообщения.value.converter- задайте в соответствии с типом данных, например для JSON:io.confluent.connect.json.JsonSchemaConverter.key.converter- установитеorg.apache.kafka.connect.storage.StringConverter. Мы используем ключи String.pk.mode- для ClickHouse неактуален. Установите none.auto.create- не поддерживается и должно быть false.auto.evolve- для этого параметра мы рекомендуем false, хотя в будущем он может поддерживаться.insert.mode- установите значение "insert". Другие режимы в настоящее время не поддерживаются.key.converter- задайте в соответствии с типами ваших ключей.value.converter- задайте в зависимости от типа данных в вашем топике. Эти данные должны иметь поддерживаемую схему — в форматах JSON, Avro или Protobuf.
Если вы используете наш пример набора данных для тестирования, убедитесь, что заданы следующие параметры:
value.converter.schemas.enable- установите false, так как мы используем Schema Registry. Установите true, если вы встраиваете схему в каждое сообщение.key.converter- установите "org.apache.kafka.connect.storage.StringConverter". Мы используем ключи String.value.converter- установите "io.confluent.connect.json.JsonSchemaConverter".value.converter.schema.registry.url- укажите URL сервера схем вместе с учётными данными для него через параметрvalue.converter.schema.registry.basic.auth.user.info.
Примеры файлов конфигурации для примера данных GitHub можно найти здесь, если Connect запущен в автономном режиме, а Kafka размещён в Confluent Cloud.
Создайте таблицу ClickHouse
Убедитесь, что таблица создана, предварительно удалив её, если она уже существует после предыдущих примеров. Ниже показан пример, совместимый с уменьшенным набором данных GitHub. Обратите внимание на отсутствие типов Array и Map, которые в настоящее время не поддерживаются:
CREATE TABLE github
(
file_time DateTime,
event_type Enum('CommitCommentEvent' = 1, 'CreateEvent' = 2, 'DeleteEvent' = 3, 'ForkEvent' = 4, 'GollumEvent' = 5, 'IssueCommentEvent' = 6, 'IssuesEvent' = 7, 'MemberEvent' = 8, 'PublicEvent' = 9, 'PullRequestEvent' = 10, 'PullRequestReviewCommentEvent' = 11, 'PushEvent' = 12, 'ReleaseEvent' = 13, 'SponsorshipEvent' = 14, 'WatchEvent' = 15, 'GistEvent' = 16, 'FollowEvent' = 17, 'DownloadEvent' = 18, 'PullRequestReviewEvent' = 19, 'ForkApplyEvent' = 20, 'Event' = 21, 'TeamAddEvent' = 22),
actor_login LowCardinality(String),
repo_name LowCardinality(String),
created_at DateTime,
updated_at DateTime,
action Enum('none' = 0, 'created' = 1, 'added' = 2, 'edited' = 3, 'deleted' = 4, 'opened' = 5, 'closed' = 6, 'reopened' = 7, 'assigned' = 8, 'unassigned' = 9, 'labeled' = 10, 'unlabeled' = 11, 'review_requested' = 12, 'review_request_removed' = 13, 'synchronize' = 14, 'started' = 15, 'published' = 16, 'update' = 17, 'create' = 18, 'fork' = 19, 'merged' = 20),
comment_id UInt64,
path String,
ref LowCardinality(String),
ref_type Enum('none' = 0, 'branch' = 1, 'tag' = 2, 'repository' = 3, 'unknown' = 4),
creator_user_login LowCardinality(String),
number UInt32,
title String,
state Enum('none' = 0, 'open' = 1, 'closed' = 2),
assignee LowCardinality(String),
closed_at DateTime,
merged_at DateTime,
merge_commit_sha String,
merged_by LowCardinality(String),
review_comments UInt32,
member_login LowCardinality(String)
) ENGINE = MergeTree ORDER BY (event_type, repo_name, created_at)Запустите Kafka Connect
Запустите Kafka Connect в автономном или распределённом режиме.
./bin/connect-standalone connect.properties.ini github-jdbc-sink.properties.iniДобавьте данные в Kafka
Отправьте сообщения в Kafka с помощью предоставленных скрипта и конфигурации. Вам потребуется изменить github.config, чтобы добавить в него учётные данные Kafka. В настоящее время скрипт настроен для использования с Confluent Cloud.
python producer.py -c github.configЭтот скрипт можно использовать для вставки любого ndjson-файла в топик Kafka. Он попытается автоматически определить схему. Предоставленная примерная конфигурация вставит только 10k сообщений — измените здесь, если это необходимо. Эта конфигурация также удаляет из набора данных все несовместимые поля Array при вставке в Kafka.
Это необходимо, чтобы коннектор JDBC мог преобразовывать сообщения в операторы INSERT. Если вы используете собственные данные, убедитесь, что либо добавляете схему в каждое сообщение (установив _value.converter.schemas.enable _в true), либо ваш клиент публикует сообщения со ссылкой на схему в registry.
Kafka Connect должен начать потреблять сообщения и вставлять строки в ClickHouse. Обратите внимание, что предупреждения вида "[JDBC Compliant Mode] Transaction isn't supported." ожидаемы, и их можно игнорировать.
Простое чтение из целевой таблицы "Github" должно подтвердить вставку данных.
SELECT count() FROM default.github;| count\(\) |
| :--- |
| 10000 |Рекомендуемые материалы для дальнейшего чтения
- Параметры конфигурации Kafka Sink Connector
- Подробный разбор Kafka Connect — JDBC Source Connector
- Подробный разбор Kafka Connect JDBC Sink: работа с первичными ключами
- Kafka Connect в действии: JDBC Sink — для тех, кто предпочитает смотреть, а не читать.
- Подробный разбор Kafka Connect — конвертеры и сериализация