Se precisar de ajuda, abra uma issue no repositório ou faça uma pergunta no Slack público do ClickHouse.
Licença
O Kafka Connector Sink é distribuído nos termos da Licença Apache 2.0Requisitos do ambiente
O framework Kafka Connect v2.7 ou posterior deve estar instalado no ambiente e ser executado com Java 11 ou posterior.Matriz de compatibilidade de versões
O conector não inicia em servidores ClickHouse anteriores à versão 23.3. Use sempre o lançamento mais recente, a menos que tenha um motivo para fixar uma versão anterior.
Principais recursos
- Vem com semântica exactly-once pronta para uso. É baseado em um novo recurso central do ClickHouse chamado KeeperMap (usado como armazenamento de estado pelo conector) e permite uma arquitetura minimalista.
- Suporte a armazenamentos de estado de terceiros: atualmente, o padrão é em memória, mas pode usar o KeeperMap (Redis será adicionado em breve).
- Integração principal: desenvolvida, mantida e suportada pela ClickHouse.
- Testado continuamente no ClickHouse Cloud.
- Inserções de dados com schema declarado e sem schema.
- Suporte a todos os tipos de dados do ClickHouse.
Instruções de instalação
Obtenha os detalhes da conexão
Para se conectar ao ClickHouse via HTTP(S), você precisa das seguintes informações:
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:

curl de exemplo.

Instruções gerais de instalação
O conector é distribuído como um único arquivo JAR contendo todos os arquivos de classe necessários para executar o plugin. Para instalar o plugin, siga estas etapas:- Baixe um arquivo ZIP contendo o arquivo JAR do conector na página de Releases do repositório ClickHouse Kafka Connect Sink.
- Extraia o conteúdo do arquivo ZIP e copie-o para o local desejado.
- Adicione à configuração plugin.path, no arquivo de propriedades do Connect, o caminho para o diretório do plugin, para que o Confluent Platform possa encontrá-lo.
- Forneça um nome de tópico, o hostname da instância do ClickHouse e a senha na configuração.
- Reinicie a Confluent Platform.
- Se você usa a Confluent Platform, faça login na UI do Confluent Control Center para verificar se o ClickHouse Sink está disponível na lista de conectores.
Opções de configuração
Para conectar o ClickHouse Sink ao servidor ClickHouse, você precisa fornecer:- detalhes da conexão: hostname (obrigatório) e porta (opcional)
- credenciais do usuário: senha (obrigatória) e nome de usuário (opcional)
- classe do conector:
com.clickhouse.kafka.connect.ClickHouseSinkConnector(obrigatória) - topics ou topics.regex: os tópicos do Kafka a serem consumidos — os nomes dos tópicos devem corresponder aos nomes das tabelas (obrigatório)
- conversores de chave e valor: defina-os com base no tipo de dados do seu tópico. Obrigatórios se ainda não estiverem definidos na configuração do worker.
Tabelas de destino
O ClickHouse Connect Sink lê mensagens de tópicos do Kafka e as grava nas tabelas apropriadas. O ClickHouse Connect Sink grava dados em tabelas já existentes. Certifique-se de que uma tabela de destino com o schema adequado tenha sido criada no ClickHouse antes de começar a inserir dados nela. Cada tópico requer uma tabela de destino dedicada no ClickHouse. O nome da tabela de destino deve corresponder ao nome do tópico de origem.Pré-processamento
Se você precisar transformar mensagens de saída antes de enviá-las para o ClickHouse Kafka Connect Sink, use Transformações do Kafka Connect.Tipos de dados suportados
Com um schema declarado:-
(1) - JSON é suportado apenas quando as configurações do ClickHouse incluem
input_format_binary_read_json_as_string=1. Isso funciona apenas para a família de formatos RowBinary, e a configuração afeta todas as colunas na requisição de insert, portanto todas elas devem ser do tipo string. Nesse caso, o conector converterá STRUCT em uma string JSON. -
(2) - Quando struct tem unions como
oneof, o converter deve ser configurado para NÃO adicionar prefixo/sufixo aos nomes de campo. Há a configuraçãogenerate.index.for.unions=falsepara oProtobufConverter.
Receitas de configuração
Estas são algumas receitas de configuração comuns para você começar rapidamente.Configuração básica
A configuração mais simples para começar pressupõe que você esteja executando o Kafka Connect no modo distribuído e tenha um servidor ClickHouse em execução emlocalhost:8443 com SSL habilitado; os dados estão em JSON sem schema.
A configuração do conector acima exige que você habilite as substituições do cliente na configuração do worker por meio de
connector.client.config.override.policy=All. Consulte a documentação do Kafka Connect para mais informações.Configuração básica com vários tópicos
O conector pode consumir dados de vários tópicosConfiguração básica com DLQ
Suporte a schema do Avro
Mapeamento de tipos do Avro
O mapeamento de tipos abaixo é definido porio.confluent.connect.avro.AvroConverter, a implementação oficial do serializador/desserializador Avro no Kafka Connect. Consulte a documentação do Kafka Connect para informações avançadas sobre a lógica de conversão.
✅: Compatível
❌: Não compatível
️⚠️: Parcialmente compatível
Consulte Tipos de dados suportados para ver o mapeamento entre os tipos do Kafka Connect e os tipos do ClickHouse.
Schemas Avro sem suporte
Os seguintes schemas Avro não têm suporte no conector:- tipo lógico
decimalemfixed
- unions Nullable
- unions em registros
Suporte a schema do Protobuf
Mapeamento de tipos do Protobuf
O mapeamento de tipos abaixo é definido porio.confluent.connect.protobuf.ProtobufConverter, a implementação oficial de serialização/desserialização do Protobuf no Kafka Connect. Consulte a documentação do Kafka Connect para informações avançadas sobre a lógica de conversão.
✅: Suportado
❌: Não suportado
️⚠️: Parcialmente suportado
Consulte Tipos de dados suportados para ver o mapeamento entre os tipos do Kafka Connect e os tipos do ClickHouse.
Observação sobre a tradução de campos oneof para colunas do ClickHouse
O conector não oferece suporte à tradução de unions (oneof) do Protobuf para o tipo Variant do ClickHouse. Em vez disso, liste os campos oneof como campos Nullable individuais no schema da sua tabela do ClickHouse.
Por exemplo:
Schemas Protobuf sem suporte
Os seguintes schemas Protobuf não são suportados pelo conector em versões mais antigas do ClickHouse:- uniões com várias mensagens (antes da versão 26.1 do CH e, antes da versão 26.9 do CH, a menos que um SETTING esteja habilitado, veja abaixo)
allow_experimental_nullable_tuple_type = 1 (consulte esta página da documentação).
Suporte a schema JSON
Suporte a String
O conector oferece suporte ao String Converter em diferentes formatos do ClickHouse: JSON, CSV e TSV.Buffering interno
O buffering interno permite que a tarefa do sink acumule registros de várias chamadaspoll() e os envie ao ClickHouse em batches maiores. Isso pode melhorar a vazão em workloads em que cada poll() produz muitos batches pequenos por partição.
Comportamento principal:
bufferCountcontrola quantos registros são mantidos em buffer antes do flush.bufferFlushTimedefine um tempo máximo de espera (em milissegundos) antes de fazer o flush dos registros em buffer.bufferFlushTimesó tem efeito quandobufferCount > 0.bufferCount=0ebufferFlushTime=0mantêm o buffering desabilitado (comportamento padrão).- O buffering não é compatível quando
exactlyOnce=true.
exactlyOnce=false na configuração do conector ou desative o buffering com bufferCount=0.
Exemplo:
Logging
O logging é fornecido automaticamente pela plataforma Kafka Connect. O destino e o formato dos logs podem ser configurados por meio do arquivo de configuração do Kafka Connect. Se estiver usando o Confluent Platform, os logs poderão ser visualizados executando um comando de CLI:Monitoramento
O ClickHouse Kafka Connect reporta métricas de runtime por meio do Java Management Extensions (JMX). O JMX vem habilitado no Kafka Connector por padrão.Métricas específicas do ClickHouse
O conector expõe métricas personalizadas por meio do seguinte nome de MBean:Métricas de Produtor/Consumidor do Kafka
O conector expõe métricas padrão de produtor e consumidor do Kafka que fornecem informações sobre o fluxo de dados, a taxa de transferência e o desempenho. Métricas no Nível do Tópico:records-sent-total: Número total de registros enviados para o tópicobytes-sent-total: Total de bytes enviados para o tópicorecord-send-rate: Taxa média de registros enviados por segundobyte-rate: Taxa média de bytes enviados por segundocompression-rate: Taxa de compressão obtida
records-sent-total: Total de registros enviados para a partiçãobytes-sent-total: Total de bytes enviados para a partiçãorecords-lag: Lag atual na partiçãorecords-lead: Lead atual na partiçãoreplica-fetch-lag: Informações de lag das réplicas
connection-creation-total: Total de conexões criadas com o nó do Kafkaconnection-close-total: Total de conexões encerradasrequest-total: Total de solicitações enviadas ao nóresponse-total: Total de respostas recebidas do nórequest-rate: Taxa média de solicitações por segundoresponse-rate: Taxa média de respostas por segundo
- Taxa de transferência: Acompanhar as taxas de ingestão de dados
- Lag: Identificar gargalos e atrasos no processamento
- Compressão: Medir a eficiência da compressão de dados
- Saúde da conexão: Monitorar a conectividade e a estabilidade da rede
Métricas do Kafka Connect Framework
O conector se integra ao framework do Kafka Connect e expõe métricas sobre o ciclo de vida das tarefas e o rastreamento de erros. Métricas de status das tarefas:task-count: Número total de tarefas no conectorrunning-task-count: Número de tarefas em execução no momentopaused-task-count: Número de tarefas pausadas no momentofailed-task-count: Número de tarefas que falharamdestroyed-task-count: Número de tarefas destruídasunassigned-task-count: Número de tarefas não atribuídas
running, paused, failed, destroyed, unassigned
Métricas de erro:
deadletterqueue-produce-failures: Número de gravações na DLQ que falharamdeadletterqueue-produce-requests: Total de tentativas de gravação na DLQlast-error-timestamp: Timestamp do último errorecords-skip-total: Número total de registros ignorados devido a errosrecords-retry-total: Número total de registros que passaram por nova tentativaerrors-total: Número total de erros encontrados
offset-commit-failures: Número de commits de offset que falharamoffset-commit-avg-time-ms: Tempo médio dos commits de offsetoffset-commit-max-time-ms: Tempo máximo dos commits de offsetput-batch-avg-time-ms: Tempo médio para processar um loteput-batch-max-time-ms: Tempo máximo para processar um lotesource-record-poll-total: Total de registros coletados
Boas práticas de monitoramento
- Monitore o lag do consumidor: Acompanhe
records-lagpor partição para identificar gargalos de processamento - Acompanhe as taxas de erro: Observe
errors-totalerecords-skip-totalpara detectar problemas de qualidade dos dados - Observe a integridade das tarefas: Monitore as métricas de status das tarefas para garantir que estejam em execução corretamente
- Meça a vazão: Use
records-send-rateebyte-ratepara acompanhar o desempenho da ingestão - Monitore a integridade da conexão: Verifique as métricas de conexão no nível do nó para identificar problemas de rede
- Acompanhe a eficiência da compressão: Use
compression-ratepara otimizar a transferência de dados
Limitações
- Não há suporte a exclusões.
- O tamanho do batch é herdado das propriedades do consumer do Kafka.
- Ao usar o KeeperMap para exactly-once, se o offset for alterado ou recuado, será necessário excluir o conteúdo do KeeperMap para esse tópico específico. (Consulte o guia de solução de problemas abaixo para mais detalhes)
Ajuste de desempenho e otimização da vazão
Esta seção aborda estratégias de ajuste de desempenho para o ClickHouse Kafka Connect Sink. O ajuste de desempenho é essencial ao lidar com casos de uso de alta vazão ou quando é necessário otimizar o uso de recursos e minimizar o lag.Quando o ajuste de desempenho é necessário?
O ajuste de desempenho normalmente é necessário nos seguintes cenários:- Cargas de trabalho de alta vazão: ao processar milhões de eventos por segundo de tópicos do Kafka
- Consumer lag: quando seu conector não consegue acompanhar a taxa de produção de dados, causando um atraso cada vez maior
- Restrições de recursos: quando você precisa otimizar o uso de CPU, memória ou rede
- Múltiplos tópicos: ao consumir simultaneamente vários tópicos de alto volume
- Mensagens pequenas: ao lidar com muitas mensagens pequenas que se beneficiariam do agrupamento em lotes no lado do servidor
- Você está processando volumes baixos a moderados (< 10.000 mensagens/segundo)
- O consumer lag é estável e aceitável para o seu caso de uso
- As configurações padrão do conector já atendem aos seus requisitos de vazão
- Seu cluster ClickHouse consegue lidar facilmente com a carga de entrada
Entendendo o fluxo de dados
Antes de ajustar, é importante entender como os dados fluem pelo conector:- Kafka Connect Framework busca mensagens dos tópicos do Kafka em segundo plano
- O conector faz polling de mensagens no buffer interno do framework
- O conector agrupa as mensagens em lotes com base no tamanho do polling
- O ClickHouse recebe a inserção em lote via HTTP/S
- O ClickHouse processa a inserção (de forma síncrona ou assíncrona)
Ajuste do tamanho do lote no Kafka Connect
O primeiro nível de otimização é controlar a quantidade de dados que o conector recebe por lote do Kafka. O Kafka Connect (o framework) busca mensagens de tópicos do Kafka em segundo plano, independentemente do conector:fetch.min.bytes: Quantidade mínima de dados antes de o framework repassar os dados ao conector (padrão: 1 byte)fetch.max.bytes: Quantidade máxima de dados a buscar em uma única solicitação (padrão: 52428800 / 50 MB)fetch.max.wait.ms: Tempo máximo de espera antes de retornar os dados sefetch.min.bytesnão for atingido (padrão: 500 ms)
No Confluent Cloud, para ajustar essas configurações, é necessário abrir um chamado de suporte pelo Confluent Cloud.
max.poll.records: Número máximo de registros retornados em uma única consulta de polling (padrão: 500)max.partition.fetch.bytes: Quantidade máxima de dados por partição (padrão: 1048576 / 1 MB)
No Confluent Cloud, para ajustar essas configurações, é necessário abrir um chamado de suporte pelo Confluent Cloud.
As propriedades acima exigem que você habilite overrides de cliente na configuração do worker por meio de
connector.client.config.override.policy=All. Consulte a documentação do Kafka Connect para mais informações.- Lotes maiores = Melhor desempenho de ingestão no ClickHouse, menos partes, menor sobrecarga
- Lotes maiores = Maior uso de memória, com possível aumento da latência de ponta a ponta
- Lotes grandes demais = Risco de timeouts, erros de OutOfMemory ou de exceder
max.poll.interval.ms
Inserções assíncronas
As inserções assíncronas são um recurso poderoso quando o conector envia lotes relativamente pequenos ou quando você quer otimizar ainda mais a ingestão, transferindo para o ClickHouse a responsabilidade pelo agrupamento em lotes. Considere ativar inserções assíncronas quando:- Muitos lotes pequenos: Seu conector envia pequenos lotes com frequência (< 1000 linhas por lote)
- Alta concorrência: Várias tarefas do conector estão gravando na mesma tabela
- Implantação distribuída: Você executa muitas instâncias do conector em hosts diferentes
- Sobrecarga na criação de partes: Você está enfrentando erros de “too many partes”
- Carga de trabalho mista: Combinação de ingestão em tempo real com cargas de trabalho de consulta
- Você já estiver enviando lotes grandes (> 10.000 linhas por lote) com frequência controlada
- Você precisar de visibilidade imediata dos dados (as consultas precisam ver os dados instantaneamente)
- A semântica exactly-once com
wait_for_async_insert=0entrar em conflito com seus requisitos - Seu caso de uso puder se beneficiar, em vez disso, de melhorias no batching no lado do cliente
- Recebe a consulta INSERT do conector
- Grava os dados em um buffer na memória (em vez de gravá-los imediatamente no disco)
- Retorna sucesso ao conector (se
wait_for_async_insert=0) - Grava o buffer no disco quando uma destas condições é atendida:
- O buffer atinge
async_insert_max_data_size(padrão: 100 MB) async_insert_busy_timeout_msmilissegundos se passaram desde a primeira inserção (padrão: 1000 ms)- Número máximo de consultas acumuladas (
async_insert_max_query_number, padrão: 100)
- O buffer atinge
clickhouseSettings:
async_insert=1: Habilita inserções assíncronaswait_for_async_insert=1(recomendado): O conector espera os dados serem gravados no armazenamento do ClickHouse antes de confirmar o recebimento. Isso oferece garantias de entrega.wait_for_async_insert=0: O conector confirma o recebimento imediatamente após armazenar os dados em buffer. Melhor desempenho, mas os dados podem ser perdidos se o servidor falhar antes da gravação.
async_insert_max_data_size(padrão: 104857600 / 100 MB): Tamanho máximo do buffer antes da gravaçãoasync_insert_busy_timeout_ms(padrão: 1000): Tempo máximo (ms) antes da gravaçãoasync_insert_stale_timeout_ms(padrão: 0): Tempo (ms) desde o último insert antes da gravaaçãoasync_insert_max_query_number(padrão: 100): Número máximo de consultas antes da gravação
- Benefícios: Menos partes, melhor desempenho de merge, menor sobrecarga de CPU, maior throughput em cenários de alta concorrência
- Considerações: Os dados não ficam imediatamente disponíveis para consulta, latência de ponta a ponta um pouco maior
- Riscos: Perda de dados em caso de falha do servidor se
wait_for_async_insert=0, possível pressão de memória com buffers grandes
exactlyOnce=true com inserções assíncronas:
wait_for_async_insert=1 com exactly-once para garantir que os commits de offset ocorram somente depois que os dados forem persistidos.
Para mais informações sobre async inserts, consulte a documentação de async inserts do ClickHouse.
Paralelismo do conector
Aumente o paralelismo para melhorar a vazão:- O número máximo efetivo de tarefas = número de partições do tópico
- Cada tarefa mantém sua própria conexão com o ClickHouse
- Mais tarefas = maior sobrecarga e possível contenção de recursos
tasks.max igual ao número de partições do tópico e depois ajuste com base nas métricas de CPU e vazão.
Por padrão, o conector agrupa as mensagens em lotes por partição. Para obter maior vazão, você pode agrupar em lotes entre partições:
exactlyOnce=false. Essa configuração pode melhorar a vazão ao criar batches maiores, mas elimina as garantias de ordenação em cada partição.
Múltiplos tópicos de alta vazão
Se o conector estiver configurado para assinar vários tópicos, você estiver usandotopic2TableMap para mapear tópicos para tabelas e estiver enfrentando um gargalo na inserção, o que resulta em consumer lag, considere criar um conector por tópico.
O principal motivo para isso acontecer é que, atualmente, os lotes são inseridos em cada tabela em série.
Recomendação: para vários tópicos de alto volume, implante uma instância de conector por tópico para maximizar a vazão de inserção em paralelo.
Considerações sobre o engine de tabela do ClickHouse
Escolha o engine de tabela do ClickHouse adequado para seu caso de uso:MergeTree: Melhor para a maioria dos casos de uso; equilibra o desempenho de consulta e inserçãoReplicatedMergeTree: Necessário para alta disponibilidade; adiciona sobrecarga de replicação*MergeTreecomORDER BYadequado: Otimize para seus padrões de consulta
insert em nível de conector:
Pool de conexões e timeouts
O conector mantém conexões HTTP com o ClickHouse. Ajuste os timeouts para redes com alta latência:socket_timeout(padrão: 30000 ms): Tempo máximo para operações de leituraconnection_timeout(padrão: 10000 ms): Tempo máximo para estabelecer a conexão
Buffer de rede do cliente
Ao usar o cliente V2 ("client_version": "V2"), os dados são copiados entre o socket e a memória da aplicação por meio de um buffer controlado por client_network_buffer_size (em bytes). O tamanho padrão é de 300000 bytes. Um buffer maior pode melhorar a vazão em conexões com alta latência ou alta largura de banda, mas aumenta o uso de memória por tarefa de conexão.
Todas as opções do cliente são configuradas por meio de jdbcConnectionProperties, e não de clickhouseSettings. Por exemplo:
&:
Monitoramento e solução de problemas de desempenho
Monitore estas métricas principais:- Consumer lag: Use ferramentas de monitoramento do Kafka para acompanhar o lag por partição
- Métricas do conector: Monitore
receivedRecords,recordProcessingTime,taskProcessingTimevia JMX (consulte Monitoring) - Métricas do ClickHouse:
system.asynchronous_inserts: Monitore o uso do buffer de async insertsystem.parts: Monitore o número de partes para detectar problemas de mergesystem.merges: Monitore merges ativossystem.events: AcompanheInsertedRows,InsertedBytes,FailedInsertQuery
Resumo das boas práticas
- Comece com os padrões, depois meça e ajuste com base no desempenho real
- Prefira lotes maiores: Busque 10.000-100.000 linhas por insert, quando possível
- Use inserções assíncronas ao enviar muitos lotes pequenos ou em cenários de alta concorrência
- Sempre use
wait_for_async_insert=1com semântica de exactly-once - Escale horizontalmente: Aumente
tasks.maxaté o número de partições - Um conector por tópico de alto volume para obter vazão máxima
- Monitore continuamente: Acompanhe o consumer lag, a contagem de partes e a atividade de merge
- Teste minuciosamente: Sempre teste alterações de configuração sob carga realista antes da implantação em produção
Exemplo: Configuração de alta vazão
Aqui está um exemplo completo otimizado para alta vazão:A configuração do conector acima exige que você habilite sobrescritas de cliente na configuração do seu worker por meio de
connector.client.config.override.policy=All. Consulte a documentação do Kafka Connect para mais informações.- Processa até 10.000 registros por poll
- Agrupa partições em lotes para inserts maiores
- Usa inserções assíncronas com buffer de 16 MB
- Executa 8 tarefas em paralelo (corresponda à sua quantidade de partições)
- Otimizada para vazão em vez de ordenação estrita
Solução de problemas
“Inconsistência de estado para o tópico [someTopic] partição [0]”
Isso acontece quando o offset armazenado no KeeperMap é diferente do offset armazenado no Kafka, geralmente quando um tópico foi excluído
ou quando o offset foi ajustado manualmente.
Para corrigir isso, você precisará excluir os valores antigos armazenados para esse tópico + partição:
Esse ajuste pode ter implicações para a semântica exactly-once.
“Quais erros o conector tentará novamente?”
No momento, o foco está em identificar erros transitórios para os quais é possível fazer nova tentativa, incluindo:ClickHouseException- Esta é uma exceção genérica que pode ser lançada pelo ClickHouse. Em geral, ela é lançada quando o servidor está sobrecarregado, e os seguintes códigos de erro são considerados especialmente transitórios:- 3 - UNEXPECTED_END_OF_FILE
- 107 - FILE_DOESNT_EXIST
- 159 - TIMEOUT_EXCEEDED
- 164 - READONLY
- 202 - TOO_MANY_SIMULTANEOUS_QUERIES
- 203 - NO_FREE_CONNECTION
- 209 - SOCKET_TIMEOUT
- 210 - NETWORK_ERROR
- 241 - MEMORY_LIMIT_EXCEEDED
- 242 - TABLE_IS_READ_ONLY
- 252 - TOO_MANY_PARTS
- 285 - TOO_FEW_LIVE_REPLICAS
- 319 - UNKNOWN_STATUS_OF_INSERT
- 425 - SYSTEM_ERROR
- 999 - KEEPER_EXCEPTION
SocketTimeoutException- Esta é lançada quando ocorre timeout no socket.UnknownHostException- Esta é lançada quando não é possível resolver o host.IOException- Esta é lançada quando há um problema de rede.
“Todos os meus dados estão em branco/zerados”
Provavelmente, os campos nos seus dados não correspondem aos campos da tabela — isso é especialmente comum com CDC (e o formato Debezium). Uma solução comum é adicionar a transformaçãoflatten à configuração do seu conector:
_ como delimitador). Os campos na tabela passarão então a seguir o formato “campo1_campo2_campo3” (ou seja, “before_id”, “after_id”, etc.).
“Quero usar minhas chaves do Kafka no ClickHouse”
As chaves do Kafka não são armazenadas no campovalue por padrão, mas você pode usar a transformação KeyToValue para mover a chave para o campo value (com o novo nome de campo _key):