Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Confluent коннектор HTTP Sink

Коннектор 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:

Кнопка подключения сервиса ClickHouse Cloud

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

Сведения о подключении к ClickHouse Cloud по HTTPS

Если вы используете самоуправляемый 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:

Интерфейс Confluent Cloud, показывающий, как создать коннектор 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 - BASIC
    • Auth username - имя пользователя ClickHouse
    • Auth password - пароль ClickHouse
Интерфейс Confluent Cloud, показывающий настройки аутентификации для коннектора HTTP Sink

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

Проверка подключения

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

Интерфейс Confluent Cloud, показывающий, как создать тестовое сообщение в топике Kafka

и убедитесь, что созданное сообщение было записано в ваш экземпляр ClickHouse.

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

HTTP Sink не выполняет батчинг сообщений

Из документации по Sink:

Коннектор HTTP Sink не выполняет батчинг запросов для сообщений с разными значениями заголовков Kafka.

  1. Убедитесь, что у ваших записей Kafka одинаковый ключ.
  2. При добавлении параметров в 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 - Установите значение POST
  • retry.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\
Navigation