Se precisar de ajuda, abra uma issue no repositório ou faça uma pergunta no Slack público do ClickHouse.
Licença
Requisitos do ambiente
Matriz de compatibilidade de versões
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
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:
Escolha HTTPS. Os detalhes de conexão são exibidos em um comando
curl de exemplo.
Se você estiver usando ClickHouse autogerenciado, os detalhes de conexão são definidos pelo administrador do seu ClickHouse.
Instruções gerais de instalação
- 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
- 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
Pré-processamento
Tipos de dados suportados
-
(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
Configuração básica
localhost: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
Configuração básica com DLQ
Uso com diferentes formatos de dados
Suporte a schema do Avro
Mapeamento de tipos do Avro
io.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
- tipo lógico
decimalemfixed
- uniões Nullable
- unions em registros
Suporte a schema do Protobuf
Mapeamento de tipos do Protobuf
io.confluent.connect.protobuf.ProtobufConverter, a implementação oficial de serialização/desserialização de Protobuf no Kafka Connect. Consulte a documentação do Kafka Connect para obter 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
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:
Esquemas Protobuf sem suporte
- uniões com várias mensagens (antes da versão 26.1 do CH)
allow_experimental_nullable_tuple_type=1 (consulte esta página da documentação).
Suporte a schema JSON
Suporte a String
Buffering interno
poll() 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
Monitoramento
Métricas específicas do ClickHouse
Métricas de Produtor/Consumidor do Kafka
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
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
Quando o ajuste de desempenho é necessário?
- 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
- 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
Configurações de fetch
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.
Configurações de polling
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.
Configurações recomendadas para alta vazão
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
Quando usar inserções assíncronas
- 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
Como as inserções assíncronas funcionam
- 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
Habilitando inserções assíncronas
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.
Ajuste do comportamento de inserção assíncrona
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
Inserções assíncronas com semântica de exactly-once
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
Tarefas por conector
- 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.
Ignorando partições no agrupamento em lotes
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
topic2TableMap 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
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
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
Monitoramento e solução de problemas de desempenho
- 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
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]”
Esse ajuste pode ter implicações para a semântica exactly-once.
”Quais erros o conector tentará novamente?”
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”
flatten à 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”
value 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):