Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Conector HTTP Sink da Confluent

O conector HTTP Sink é agnóstico quanto a tipos de dados e, portanto, não requer um esquema do Kafka, além de oferecer suporte a tipos de dados específicos do ClickHouse, como Map e Array. Essa flexibilidade adicional traz um pequeno aumento na complexidade da configuração.

Abaixo, descrevemos uma instalação simples, extraindo mensagens de um único tópico do Kafka e inserindo linhas em uma tabela do ClickHouse.

Etapas de início rápido

Reúna seus detalhes de 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.

Execute o Kafka Connect e o conector HTTP Sink

Você tem duas opções:

  • Autogerenciado: Baixe o pacote da Confluent e instale-o localmente. Siga as instruções de instalação do conector conforme documentado aqui. Se você usar o método de instalação confluent-hub, seus arquivos de configuração locais serão atualizados.

  • Confluent Cloud: Uma versão totalmente gerenciada do HTTP Sink está disponível para quem usa o Confluent Cloud para hospedar o Kafka. Isso exige que seu ambiente ClickHouse esteja acessível a partir do Confluent Cloud.

Crie a tabela de destino no ClickHouse

Antes do teste de conectividade, vamos começar criando uma tabela de teste no ClickHouse Cloud; essa tabela receberá os dados do Kafka:

CREATE TABLE default.my_table
(
    `side` String,
    `quantity` Int32,
    `symbol` String,
    `price` Int32,
    `account` String,
    `userid` String
)
ORDER BY tuple()

Configure o HTTP Sink

Crie um tópico do Kafka e uma instância do conector HTTP Sink:

Interface do Confluent Cloud mostrando como criar um conector HTTP Sink

Configure o conector HTTP Sink:

  • Informe o nome do tópico que você criou
  • Autenticação
    • HTTP Url - URL do ClickHouse Cloud com uma consulta INSERT especificada: <protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow. Observação: a consulta deve ser codificada.
    • Endpoint Authentication type - BASIC
    • Auth username - nome de usuário do ClickHouse
    • Auth password - senha do ClickHouse
Interface do Confluent Cloud mostrando as configurações de autenticação do conector HTTP Sink

  • Configuração
    • Input Kafka record value format - Depende dos seus dados de origem, mas, na maioria dos casos, será JSON ou Avro. Assumimos JSON nas configurações a seguir.
    • Na seção advanced configurations:
      • HTTP Request Method - Defina como POST
      • Request Body Format - json
      • Batch batch size - De acordo com as recomendações do ClickHouse, defina esse valor como no mínimo 1000.
      • Batch json as array - true
      • Retry on HTTP codes - 400-500, mas ajuste conforme necessário; por exemplo, isso pode mudar se você tiver um proxy HTTP na frente do ClickHouse.
      • Maximum Reties - o padrão (10) é adequado, mas fique à vontade para ajustar se quiser tentativas de repetição mais robustas.
Interface do Confluent Cloud mostrando opções avançadas de configuração do conector HTTP Sink

Testando a conectividade

Crie uma mensagem em um tópico configurado pelo seu HTTP Sink

Interface do Confluent Cloud mostrando como criar uma mensagem de teste em um tópico do Kafka

e verifique se a mensagem criada foi gravada na sua instância do ClickHouse.

Solução de problemas

O HTTP Sink não agrupa mensagens em lote

Da documentação do Sink:

O conector HTTP Sink não agrupa em lote solicitações de mensagens que contêm valores de header do Kafka diferentes.

  1. Verifique se os registros do Kafka têm a mesma chave.
  2. Ao adicionar parâmetros à URL da API HTTP, cada registro pode resultar em uma URL exclusiva. Por esse motivo, o agrupamento em lote é desativado ao usar parâmetros de URL adicionais.

400 requisição inválida

CANNOT_PARSE_QUOTED_STRING

Se o HTTP Sink falhar e exibir a seguinte mensagem ao inserir um objeto JSON em uma coluna String:

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)

Defina a configuração input_format_json_read_objects_as_strings=1 na URL como uma string codificada SETTINGS%20input_format_json_read_objects_as_strings%3D1

Carregue o conjunto de dados do GitHub (opcional)

Observe que este exemplo preserva os campos Array do conjunto de dados do GitHub. Pressupomos que, nos exemplos, você tenha um tópico github vazio e use o kcat para inserir mensagens no Kafka.

Preparar a configuração

Siga estas instruções para configurar o Connect de acordo com o seu tipo de instalação, observando as diferenças entre um cluster standalone e um distribuído. Se estiver usando o Confluent Cloud, a configuração distribuída é a aplicável.

O parâmetro mais importante é o http.api.url. A interface HTTP do ClickHouse exige que você codifique a instrução INSERT como um parâmetro na URL. Isso deve incluir o formato (JSONEachRow, neste caso) e o banco de dados de destino. O formato deve ser compatível com os dados do Kafka, que serão convertidos em uma string no payload HTTP. Esses parâmetros devem ser escapados na URL. Um exemplo desse formato para o conjunto de dados do GitHub (supondo que você esteja executando o ClickHouse localmente) é mostrado abaixo:

<protocol>://<clickhouse_host>:<clickhouse_port>?query=INSERT%20INTO%20<database>.<table>%20FORMAT%20JSONEachRow

http://localhost:8123?query=INSERT%20INTO%20default.github%20FORMAT%20JSONEachRow

Os seguintes parâmetros adicionais são relevantes para usar o HTTP Sink com o ClickHouse. Uma lista completa de parâmetros pode ser encontrada aqui:

  • request.method - Defina como POST
  • retry.on.status.codes - Defina como 400-500 para repetir a tentativa em qualquer código de erro. Ajuste conforme os erros esperados nos dados.
  • request.body.format - Na maioria dos casos, será JSON.
  • auth.type - Defina como BASIC se você usar segurança com o ClickHouse. Outros mecanismos de autenticação compatíveis com o ClickHouse não são compatíveis no momento.
  • ssl.enabled - defina como true se estiver usando SSL.
  • connection.user - nome de usuário do ClickHouse.
  • connection.password - senha do ClickHouse.
  • batch.max.size - O número de linhas a serem enviadas em um único lote. Certifique-se de definir um número adequadamente alto. De acordo com as recomendações do ClickHouse, um valor de 1000 deve ser considerado o mínimo.
  • tasks.max - O conector HTTP Sink oferece suporte à execução de uma ou mais tarefas. Isso pode ser usado para aumentar o desempenho. Junto com o tamanho do lote, esse é o principal meio de melhorar o desempenho.
  • key.converter - defina de acordo com os tipos das suas chaves.
  • value.converter - defina com base no tipo de dados no seu tópico. Esses dados não precisam de um esquema. O formato aqui deve ser consistente com o FORMAT especificado no parâmetro http.api.url. A opção mais simples é usar JSON e o conversor org.apache.kafka.connect.json.JsonConverter. Também é possível tratar o valor como uma string, por meio do conversor org.apache.kafka.connect.storage.StringConverter, embora isso exija que o usuário extraia um valor na instrução insert usando funções. O formato Avro também é compatível com o ClickHouse ao usar o conversor io.confluent.connect.avro.AvroConverter.

Uma lista completa de configurações, incluindo como configurar um proxy, tentativas e SSL avançado, pode ser encontrada aqui.

Exemplos de arquivos de configuração para os dados de amostra do GitHub podem ser encontrados aqui, considerando que o Connect esteja em execução no modo standalone e que o Kafka esteja hospedado no Confluent Cloud.

Criar a tabela no ClickHouse

Certifique-se de que a tabela foi criada. Um exemplo de um conjunto de dados mínimo do GitHub usando um MergeTree padrão é mostrado abaixo.

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)

Adicionar dados ao Kafka

Insira mensagens no Kafka. A seguir, usamos kcat para inserir 10 mil mensagens.

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

Uma simples consulta à tabela de destino "Github" deve confirmar a inserção dos dados.

SELECT count() FROM default.github;

| count\
Navigation