Использование Vector с Kafka и ClickHouse
Vector — это не зависящий от поставщика конвейер данных, который может читать данные из Kafka и отправлять события в ClickHouse.
Руководство по началу работы для Vector с ClickHouse посвящено сценарию работы с логами и чтению событий из файла. Мы используем пример датасета GitHub, в котором события хранятся в топике Kafka.
Vector использует источники для получения данных по модели push или pull. Приёмники, в свою очередь, выступают пунктом назначения для событий. Поэтому мы используем источник Kafka и приёмник ClickHouse. Обратите внимание: хотя Kafka поддерживается в качестве приёмника, источник для ClickHouse недоступен. Поэтому Vector не подходит, если вам нужно передавать данные из ClickHouse в Kafka.
Vector также поддерживает преобразование данных. Это выходит за рамки данного руководства. Если это требуется для вашего датасета, обратитесь к документации Vector.
Обратите внимание, что текущая реализация приёмника ClickHouse использует HTTP-интерфейс. Приёмник ClickHouse в настоящее время не поддерживает использование схемы JSON. Данные должны публиковаться в Kafka либо в обычном формате JSON, либо в виде String.
Лицензия
Vector распространяется на условиях лицензии MPL-2.0
Подготовьте сведения о подключении
Чтобы подключиться к 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
githubи вставьте датасет GitHub.
cat /opt/data/github/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Этот набор данных включает 200 000 строк, относящихся к репозиторию ClickHouse/ClickHouse.
- Убедитесь, что целевая таблица создана. Ниже используется база данных по умолчанию.
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);- Скачайте и установите Vector. Создайте файл конфигурации
kafka.tomlи измените значения для ваших экземпляров Kafka и ClickHouse.
[sources.github]
type = "kafka"
auto_offset_reset = "smallest"
bootstrap_servers = "<kafka_host>:<kafka_port>"
group_id = "vector"
topics = [ "github" ]
tls.enabled = true
sasl.enabled = true
sasl.mechanism = "PLAIN"
sasl.username = "<username>"
sasl.password = "<password>"
decoding.codec = "json"
[sinks.clickhouse]
type = "clickhouse"
inputs = ["github"]
endpoint = "http://localhost:8123"
database = "default"
table = "github"
skip_unknown_fields = true
auth.strategy = "basic"
auth.user = "username"
auth.password = "password"
buffer.max_events = 10000
batch.timeout_secs = 1Несколько важных замечаний об этой конфигурации и о поведении Vector:
- Этот пример был протестирован с Confluent Cloud. Поэтому параметры безопасности
sasl.*иssl.enabledмогут не подходить для самоуправляемых развертываний. - Для параметра конфигурации
bootstrap_serversпрефикс протокола не требуется, например:pkc-2396y.us-east-1.aws.confluent.cloud:9092 - Параметр источника
decoding.codec = "json"гарантирует, что сообщение будет передано в приёмник ClickHouse как единый объект JSON. Если обрабатывать сообщения как строки и использовать значениеbytesпо умолчанию, содержимое сообщения будет добавлено в полеmessage. В большинстве случаев это потребует дополнительной обработки в ClickHouse, как описано в руководстве Начало работы с Vector. - Vector добавляет ряд полей в сообщения. В нашем примере мы игнорируем эти поля в приёмнике ClickHouse с помощью параметра конфигурации
skip_unknown_fields = true. При этом игнорируются поля, которых нет в схеме целевой таблицы. При необходимости скорректируйте схему, чтобы добавить эти метаполя, напримерoffset. - Обратите внимание, как приёмник ссылается на источник событий через параметр
inputs. - Обратите внимание на поведение приёмника ClickHouse, описанное здесь. Для оптимальной пропускной способности можно настроить параметры
buffer.max_events,batch.timeout_secsиbatch.max_bytes. Согласно рекомендациям ClickHouse, значение 1000 следует считать минимальным количеством событий в одном батче. Для сценариев с равномерно высокой пропускной способностью можно увеличить параметрbuffer.max_events. При более переменной пропускной способности может потребоваться изменить параметрbatch.timeout_secs. - Параметр
auto_offset_reset = "smallest"принудительно задает для источника Kafka чтение с начала топика, гарантируя, что будут прочитаны сообщения, опубликованные на шаге (1). В вашем случае может потребоваться другое поведение. Подробнее см. здесь.
- Запустите Vector
vector --config ./kafka.tomlПо умолчанию перед началом вставки в ClickHouse требуется выполнить проверку работоспособности. Это позволяет убедиться, что соединение устанавливается и схему можно прочитать. Добавьте в начало VECTOR_LOG=debug, чтобы включить более подробное логирование — это может быть полезно, если возникнут проблемы.
- Подтвердите вставку данных.
SELECT count() AS count FROM github;| количество |
|---|
| 200000 |