Kafka para ClickHouse
Se você usa ClickHouse Cloud, recomendamos usar o ClickPipes. O ClickPipes oferece suporte nativo a conexões de rede privadas, ao dimensionamento independente da ingestão e dos recursos do cluster, além de monitoramento abrangente para streaming de dados do Kafka para o ClickHouse.
Visão geral
Etapas
1. Preparar
2. Configure o ClickHouse
config.xml do ClickHouse. Presumimos que você esteja se conectando a uma instância protegida por SASL. Esse é o método mais simples ao interagir com o Confluent Cloud.
conf.d/ ou mescle-o aos arquivos de configuração existentes. Para ver quais configurações podem ser definidas, consulte aqui.
Também vamos criar um banco de dados chamado KafkaEngine para usar neste tutorial:
3. Crie a tabela de destino
4. Criar e preencher o tópico
github com 5 partições executando o seguinte comando:
5. Crie o motor de tabela Kafka
JSONEachRow como tipo de dado para consumir JSON de um tópico Kafka. Os valores github e clickhouse representam, respectivamente, o nome do tópico e os nomes dos grupos de consumidores. Na prática, os tópicos podem ser uma lista de valores.
github_queue deve ler algumas linhas. Observe que isso avançará os offsets do consumidor, impedindo que essas linhas sejam lidas novamente sem um reset. Observe o limite e o parâmetro obrigatório stream_like_engine_allow_direct_select.
6. Criar a visão materializada
7. Confirme que as linhas foram inseridas
Operações comuns
Interrompendo & reiniciando o consumo de mensagens
Adicionando metadados do Kafka
SELECT da visão materializada.
Primeiro, executamos a operação de parada descrita acima antes de adicionar colunas à nossa tabela de destino.
_.
Uma lista completa das colunas virtuais pode ser encontrada aqui.
Para atualizar nossa tabela com as colunas virtuais, precisaremos remover a visão materializada, reanexar a tabela do motor Kafka e recriar a visão materializada.
Modificar as configurações do motor Kafka
Depuração de problemas
Tratando mensagens malformadas
- Trate os campos da mensagem como strings. É possível usar funções na instrução da visão materializada para fazer limpeza e conversão de tipo, se necessário. Isso não deve ser considerado uma solução de produção, mas pode ajudar em uma ingestão pontual.
- Se você estiver consumindo JSON de um tópico usando o formato JSONEachRow, use a configuração
input_format_skip_unknown_fields. Ao gravar dados, por padrão, o ClickHouse lança uma exceção se os dados de entrada contiverem colunas que não existem na tabela de destino. No entanto, se essa opção estiver habilitada, essas colunas excedentes serão ignoradas. Novamente, isso não é uma solução adequada para produção e pode confundir outras pessoas. - Considere a configuração
kafka_skip_broken_messages. Isso exige que o usuário especifique o nível de tolerância por bloco para mensagens malformadas, considerado no contexto dekafka_max_block_size. Se essa tolerância for excedida (medida em número absoluto de mensagens), o comportamento normal de exceção será retomado, e as demais mensagens serão ignoradas.
Semântica de entrega e desafios com duplicatas
Inserções com quórum
ClickHouse para Kafka
Etapas
1. Inserindo linhas diretamente
2. Usando visões materializadas
github_out ou equivalente. Certifique-se de que um motor de tabela Kafka github_out_queue aponte para esse tópico.
github_out_mv que aponte para a tabela GitHub, inserindo linhas no mecanismo acima quando ela for acionada. Como resultado, as adições à tabela GitHub serão enviadas ao nosso novo tópico Kafka.
github_out deve confirmar a entrega das mensagens.
Clusters e desempenho
Trabalhando com clusters do ClickHouse
Ajuste de desempenho
- O desempenho varia conforme o tamanho da mensagem, o formato e os tipos da tabela de destino. É razoável esperar 100 mil linhas/s em um único motor de tabela. Por padrão, as mensagens são lidas em blocos, controlados pelo parâmetro
kafka_max_block_size. Por padrão, ele é definido como max_insert_block_size, cujo valor padrão é 1.048.576. A menos que as mensagens sejam extremamente grandes, quase sempre vale a pena aumentar esse valor. Valores entre 500 mil e 1 milhão não são incomuns. Teste e avalie o efeito sobre a taxa de transferência. - O número de consumers de um motor de tabela pode ser aumentado com
kafka_num_consumers. No entanto, por padrão, os inserts serão linearizados em uma única thread, a menos quekafka_thread_per_consumerseja alterado do valor padrão 1. Defina-o como 1 para garantir que os flushes sejam realizados em paralelo. Observe que criar uma tabela com motor Kafka com N consumers (ekafka_thread_per_consumer=1) é logicamente equivalente a criar N motores Kafka, cada um com uma visão materializada ekafka_thread_per_consumer=0. - Aumentar o número de consumers não é gratuito. Cada consumer mantém seus próprios buffers e threads, aumentando a sobrecarga no servidor. Fique atento a essa sobrecarga e, se possível, primeiro faça o scale linearmente no cluster.
- Se a taxa de transferência das mensagens do Kafka variar e atrasos forem aceitáveis, considere aumentar
stream_flush_interval_mspara garantir que blocos maiores sejam gravados. - background_message_broker_schedule_pool_size define o número de threads que executam tasks em segundo plano. Essas threads são usadas para streaming do Kafka. Essa configuração é aplicada na inicialização do servidor ClickHouse e não pode ser alterada em uma sessão de usuário; o valor padrão é 16. Se você observar timeouts nos logs, pode ser apropriado aumentá-la.
- Para a comunicação com o Kafka, é usada a biblioteca librdkafka, que por sua vez cria threads. Assim, um grande número de tabelas Kafka ou de consumers pode resultar em muitas trocas de contexto. Distribua essa carga pelo cluster, replicando apenas as tabelas de destino, se possível, ou considere usar um motor de tabela para ler de vários topics — uma lista de valores é compatível. Várias visões materializadas podem ler de uma única tabela, cada uma filtrando os dados de um topic específico.
Configurações adicionais
- Kafka_max_wait_ms - O tempo de espera, em milissegundos, para ler mensagens do Kafka antes de tentar novamente. É definido no nível do perfil do usuário e o padrão é 5000.