Коннектор HTTP Sink не зависит от типа данных, поэтому ему не требуется схема Kafka; кроме того, он поддерживает специфичные для ClickHouse типы данных, такие как Map и Array. За эту дополнительную гибкость приходится платить небольшим усложнением конфигурации.
Ниже мы опишем простую установку: получение сообщений из одного топика Kafka и вставку строк в таблицу ClickHouse.
Шаги быстрого старта
Соберите сведения о подключении
Чтобы подключиться к 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 и коннектор HTTP Sink
У вас есть два варианта:
-
Самоуправляемый: Загрузите пакет Confluent и установите его локально. Следуйте инструкциям по установке коннектора, описанным здесь. Если вы используете метод установки
confluent-hub, ваши локальные файлы конфигурации будут обновлены. -
Confluent Cloud: Для тех, кто использует Confluent Cloud для хостинга Kafka, доступна полностью управляемая версия HTTP Sink. Для этого требуется, чтобы ваша среда ClickHouse была доступна из Confluent Cloud.
Создайте целевую таблицу в ClickHouse
Перед проверкой подключения давайте сначала создадим тестовую таблицу в ClickHouse Cloud; эта таблица будет получать данные из Kafka:
CREATE TABLE default.my_table
(
`side` String,
`quantity` Int32,
`symbol` String,
`price` Int32,
`account` String,
`userid` String
)
ORDER BY tuple()Настройка HTTP Sink
Создайте Kafka топик и экземпляр коннектора HTTP Sink:

Настройте коннектор HTTP Sink:
- Укажите имя созданного топика
- Аутентификация
HTTP Url- URL ClickHouse Cloud с указанным запросомINSERT<protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow. Примечание: запрос должен быть URL-кодирован.Endpoint Authentication type- BASICAuth username- имя пользователя ClickHouseAuth password- пароль ClickHouse

- Конфигурация
Input Kafka record value formatЗависит от исходных данных, но в большинстве случаев это JSON или Avro. В следующих настройках мы исходим изJSON.- В разделе
advanced configurations:HTTP Request Method- установите значение POSTRequest Body Format- jsonBatch batch size- согласно рекомендациям ClickHouse, установите значение не менее 1000.Batch json as array- trueRetry on HTTP codes- 400-500, но при необходимости скорректируйте; например, это может измениться, если перед ClickHouse используется HTTP proxy.Maximum Reties- значение по умолчанию (10) подходит, но при необходимости вы можете увеличить его для более надёжных повторных попыток.

Проверка подключения
Создайте сообщение в топике, настроенном для вашего HTTP Sink

и убедитесь, что созданное сообщение было записано в ваш экземпляр ClickHouse.
Устранение неполадок
HTTP Sink не выполняет батчинг сообщений
Коннектор HTTP Sink не выполняет батчинг запросов для сообщений с разными значениями заголовков Kafka.
- Убедитесь, что у ваших записей Kafka одинаковый ключ.
- При добавлении параметров в URL HTTP API для каждой записи может формироваться уникальный URL. По этой причине батчинг отключается при использовании дополнительных параметров URL.
400 Некорректный запрос
CANNOT_PARSE_QUOTED_STRING
Если при вставке объекта JSON в столбец String в HTTP Sink возникает следующая ошибка:
Code: 26. DB::ParsingException: Cannot parse JSON string: expected opening quote: (while reading the value of key key_name): While executing JSONEachRowRowInputFormat: (at row 1). (CANNOT_PARSE_QUOTED_STRING)Установите параметр input_format_json_read_objects_as_strings=1 в URL в виде закодированной строки SETTINGS%20input_format_json_read_objects_as_strings%3D1
Загрузите набор данных GitHub (необязательно)
Обратите внимание: в этом примере сохраняются поля типа Array из набора данных GitHub. Мы предполагаем, что в примерах у вас есть пустой топик github, и для вставки сообщений в Kafka используем kcat.
Подготовьте конфигурацию
Следуйте этим инструкциям по настройке Connect в соответствии с типом вашей установки, учитывая различия между автономным и распределённым кластером. Если вы используете Confluent Cloud, вам нужен распределённый вариант.
Самый важный параметр — http.api.url. HTTP-интерфейс ClickHouse требует передавать оператор INSERT как параметр URL. Он должен включать формат (JSONEachRow в данном случае) и целевую базу данных. Формат должен соответствовать данным Kafka, которые будут преобразованы в строку в HTTP-полезной нагрузке. Эти параметры должны быть URL-кодированы. Ниже показан пример такого формата для набора данных Github (при условии, что вы запускаете ClickHouse локально):
<protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow
http://localhost:8123?query=INSERT%20INTO%20default.github%20FORMAT%20JSONEachRowСледующие дополнительные параметры относятся к использению коннектора HTTP Sink с ClickHouse. Полный список параметров можно найти здесь:
request.method- Установите значение POSTretry.on.status.codes- Установите 400-500, чтобы повторять попытки при любых кодах ошибок. Скорректируйте значение с учетом ожидаемых ошибок в данных.request.body.format- В большинстве случаев это будет JSON.auth.type- Установите значение BASIC, если для ClickHouse используется BASIC-аутентификация. Другие совместимые с ClickHouse механизмы аутентификации в настоящее время не поддерживаются.ssl.enabled- установите значение true, если используете SSL.connection.user- имя пользователя ClickHouse.connection.password- пароль ClickHouse.batch.max.size- Количество строк, отправляемых в одном батче. Убедитесь, что здесь задано достаточно большое значение. Согласно рекомендациям ClickHouse, значение 1000 следует считать минимальным.tasks.max- Коннектор HTTP Sink поддерживает запуск одной или нескольких задач. Это можно использовать для повышения производительности. Наряду с размером батча это основной способ повысить производительность.key.converter- задайте в соответствии с типами ваших ключей.value.converter- задайте в зависимости от типа данных в вашем топике. Для этих данных схема не требуется. Формат здесь должен соответствовать FORMAT, указанному в параметреhttp.api.url. Проще всего использовать JSON и конвертер org.apache.kafka.connect.json.JsonConverter. Также можно обрабатывать значение как строку с помощью конвертера org.apache.kafka.connect.storage.StringConverter, хотя в этом случае пользователю потребуется извлекать значение в операторе вставки с помощью функций. Формат Avro также поддерживается в ClickHouse при использовании конвертера io.confluent.connect.avro.AvroConverter.
Полный список настроек, включая сведения о том, как настроить прокси, повторные попытки и расширенные параметры SSL, можно найти здесь.
Примеры файлов конфигурации для демонстрационных данных GitHub можно найти здесь при условии, что Connect запущен в автономном режиме, а Kafka размещён в Confluent Cloud.
Создайте таблицу в ClickHouse
Убедитесь, что таблица создана. Ниже приведён пример минимального набора данных GitHub с использованием стандартного движка MergeTree.
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,
labels Array(LowCardinality(String)),
state Enum('none' = 0, 'open' = 1, 'closed' = 2),
assignee LowCardinality(String),
assignees Array(LowCardinality(String)),
closed_at DateTime,
merged_at DateTime,
merge_commit_sha String,
requested_reviewers Array(LowCardinality(String)),
merged_by LowCardinality(String),
review_comments UInt32,
member_login LowCardinality(String)
) ENGINE = MergeTree ORDER BY (event_type, repo_name, created_at)Добавьте данные в Kafka
Запишите сообщения в Kafka. Ниже мы используем kcat, чтобы отправить 10 тыс. сообщений.
head -n 10000 github_all_columns.ndjson | kcat -b <host>:<port> -X security.protocol=sasl_ssl -X sasl.mechanisms=PLAIN -X sasl.username=<username> -X sasl.password=<password> -t githubПростой запрос к целевой таблице "Github" должен подтвердить вставку данных.
SELECT count() FROM default.github;
| count\