Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Como usar o Vector com Kafka e ClickHouse

Usando o Vector com Kafka e ClickHouse

O Vector é um pipeline de dados agnóstico a fornecedores, capaz de ler do Kafka e enviar eventos ao ClickHouse.

Um guia de primeiros passos para usar o Vector com o ClickHouse se concentra no caso de uso de logs e na leitura de eventos de um arquivo. Utilizamos o dataset de exemplo do GitHub, com eventos armazenados em um tópico do Kafka.

O Vector utiliza fontes para recuperar dados por meio de um modelo de envio ou extração. Já os sinks fornecem um destino para os eventos. Portanto, utilizamos a fonte Kafka e o sink do ClickHouse. Observe que, embora o Kafka tenha suporte como sink, não há uma fonte do ClickHouse disponível. Como resultado, o Vector não é apropriado se você quiser transferir dados do ClickHouse para o Kafka.

O Vector também oferece suporte à transformação de dados. Isso está fora do escopo deste guia. Caso precise desse recurso para seu dataset, consulte a documentação do Vector.

Observe que a implementação atual do sink do ClickHouse utiliza a interface HTTP. No momento, o sink do ClickHouse não oferece suporte ao uso de um esquema JSON. Os dados devem ser publicados no Kafka em JSON simples ou como Strings.

Licença

O Vector é distribuído sob a licença MPL-2.0

Reúna os detalhes da conexão

Para se conectar ao ClickHouse via HTTP(S), você precisa das seguintes informações:

Parâmetro(s) Descrição
HOST and PORT Normalmente, a porta é 8443 ao usar TLS ou 8123 quando não se usa TLS.
DATABASE NAME Por padrão, há um banco de dados chamado default; use o nome do banco de dados ao qual você deseja se conectar.
USERNAME and PASSWORD Por padrão, o nome de usuário é default. Use o nome de usuário apropriado para o seu caso de uso.

Os detalhes do seu serviço do ClickHouse Cloud estão disponíveis no console do ClickHouse Cloud. Selecione um serviço e clique em Connect:

botão Connect do serviço do ClickHouse Cloud

Escolha HTTPS. Os detalhes de conexão são exibidos em um comando curl de exemplo.

detalhes de conexão HTTPS do ClickHouse Cloud

Se você estiver usando ClickHouse autogerenciado, os detalhes de conexão são definidos pelo administrador do seu ClickHouse.

Etapas

  1. Crie o tópico github no Kafka e insira o dataset do 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

Este conjunto de dados contém 200.000 linhas focadas no repositório ClickHouse/ClickHouse.

  1. Certifique-se de que a tabela de destino foi criada. Abaixo, usamos o banco de dados padrão.

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);
  1. Baixe e instale o Vector. Crie um arquivo de configuração kafka.toml e ajuste os valores das suas instâncias do Kafka e do 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

Algumas observações importantes sobre esta configuração e o comportamento do Vector:

  • Este exemplo foi testado no Confluent Cloud. Portanto, as opções de segurança sasl.* e ssl.enabled podem não ser adequadas em ambientes autogerenciados.
  • Não é necessário informar um prefixo de protocolo para o parâmetro de configuração bootstrap_servers, por exemplo pkc-2396y.us-east-1.aws.confluent.cloud:9092
  • O parâmetro da fonte decoding.codec = "json" garante que a mensagem seja passada ao sink do ClickHouse como um único objeto JSON. Ao tratar mensagens como strings e usar o valor padrão bytes, o conteúdo da mensagem será anexado ao campo message. Na maioria dos casos, isso exigirá processamento no ClickHouse, conforme descrito no guia Primeiros passos com o Vector.
  • O Vector adiciona vários campos às mensagens. No nosso exemplo, ignoramos esses campos no sink do ClickHouse por meio do parâmetro de configuração skip_unknown_fields = true. Isso ignora campos que não fazem parte do esquema da tabela de destino. Se quiser, ajuste seu esquema para garantir que esses metacampos, como offset, sejam adicionados.
  • Observe como o sink faz referência à fonte de eventos por meio do parâmetro inputs.
  • Observe o comportamento do sink do ClickHouse, conforme descrito aqui. Para obter vazão ideal, talvez seja interessante ajustar os parâmetros buffer.max_events, batch.timeout_secs e batch.max_bytes. Conforme as recomendações do ClickHouse, um valor de 1000 deve ser considerado o mínimo para o número de eventos em um único batch. Para casos de uso com alta vazão uniforme, você pode aumentar o parâmetro buffer.max_events. Vazões mais variáveis podem exigir alterações no parâmetro batch.timeout_secs
  • O parâmetro auto_offset_reset = "smallest" força a fonte do Kafka a começar no início do tópico, garantindo assim que consumamos as mensagens publicadas na etapa (1). Talvez você precise de um comportamento diferente. Veja aqui para mais detalhes.
  1. Inicie o Vector
vector --config ./kafka.toml

Por padrão, uma verificação de integridade é necessária antes de iniciar as inserções no ClickHouse. Isso garante que a conectividade possa ser estabelecida e que o esquema possa ser lido. Adicione VECTOR_LOG=debug no início para obter logs adicionais, o que pode ser útil caso você encontre problemas.

  1. Confirme a inserção dos dados.
SELECT count() AS count FROM github;
count
200000
Navigation