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
Prepare
Se você tiver dados carregados em um tópico de destino, poderá adaptar o conteúdo a seguir para uso no seu conjunto de dados. Como alternativa, um conjunto de dados do GitHub de exemplo é fornecido aqui. Esse conjunto de dados é usado nos exemplos abaixo e utiliza um esquema reduzido e um subconjunto das linhas (especificamente, limitamos aos eventos do GitHub relacionados ao repositório do ClickHouse), em comparação com o conjunto de dados completo disponível aqui, por questão de brevidade. Ainda assim, isso é suficiente para que a maioria das consultas publicadas com o conjunto de dados funcione.
2
Configure o ClickHouse
Esta etapa é necessária se você estiver se conectando a um Kafka com segurança habilitada. Essas configurações não podem ser passadas por comandos SQL DDL e devem ser configuradas no config.xml do ClickHouse. Presumimos que você esteja se conectando a uma instância protegida por SASL. Este é o método mais simples ao interagir com a Confluent Cloud.Coloque o trecho acima em um novo arquivo no diretório Depois de criar o banco de dados, você precisará mudar para ele:
conf.d/ ou mescle-o aos arquivos de configuração existentes. Para as configurações que 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
Prepare a sua tabela de destino. No exemplo abaixo, usamos o esquema reduzido do GitHub por questão de brevidade. Observe que, embora usemos um motor de tabela MergeTree, este exemplo pode ser facilmente adaptado para qualquer membro da família MergeTree.
4
Crie e popule o tópico
Em seguida, vamos criar um tópico. Há várias ferramentas que podemos usar para fazer isso. Se estivermos executando o Kafka localmente em nossa máquina ou em um contêiner Docker, o RPK funciona bem. Podemos criar um tópico chamado Se estivermos executando o Kafka no Confluent Cloud, talvez prefiramos usar a Confluent CLI:Agora precisamos popular este tópico com alguns dados, o que faremos usando kcat. Podemos executar um comando semelhante ao seguinte se estivermos executando o Kafka localmente com a autenticação desabilitada:Ou o seguinte, se nosso cluster Kafka usar SASL para autenticação:O conjunto de dados contém 200.000 linhas, portanto a ingestão deve levar apenas alguns segundos. Se você quiser trabalhar com um conjunto de dados maior, dê uma olhada na seção sobre grandes conjuntos de dados do repositório do GitHub ClickHouse/kafka-samples.
github com 5 partições executando o seguinte comando:5
Crie o motor de tabela Kafka
O exemplo abaixo cria um motor de tabela com o mesmo esquema da tabela MergeTree. Isso não é estritamente necessário, pois você pode ter colunas alias ou efêmeras na tabela de destino. As configurações são importantes; no entanto, observe o uso de Discutimos abaixo as configurações do motor e o ajuste de desempenho. Neste ponto, um
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 o nome do grupo de consumidores. Na verdade, os tópicos podem ser especificados como uma lista de valores.select simples na tabela 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
Crie a visão materializada
A visão materializada conectará as duas tabelas criadas anteriormente, lendo dados do motor de tabela Kafka e inserindo-os na tabela MergeTree de destino. Podemos fazer várias transformações nos dados. Faremos uma leitura e inserção simples. O uso de * pressupõe que os nomes das colunas sejam idênticos (diferenciam maiúsculas de minúsculas).No momento da criação, a visão materializada se conecta ao motor Kafka e começa a ler, inserindo linhas na tabela de destino. Esse processo continuará indefinidamente, consumindo as mensagens inseridas posteriormente no Kafka. Sinta-se à vontade para executar novamente o script de inserção para inserir mais mensagens no Kafka.
7
Confirme se as linhas foram inseridas
Confirme que existem dados na tabela de destino:Você deve ver 200.000 linhas:
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
Inserção direta de linhas
Primeiro, confirme o número de linhas da tabela de destino.Você deve ter 200.000 linhas:Agora insira linhas da tabela de destino GitHub de volta para o motor de tabela Kafka github_queue. Observe como utilizamos o formato JSONEachRow e limitamos o select a 100.Conte novamente as linhas na tabela GitHub para confirmar que o total aumentou em 100. Como mostrado no diagrama acima, as linhas foram inseridas no Kafka por meio do motor de tabela Kafka antes de serem relidas pelo mesmo motor e inseridas na tabela de destino GitHub pela nossa visão materializada!Você deve ver mais 100 linhas:
2
Usando visões materializadas
Podemos usar visões materializadas para enviar mensagens a um motor Kafka (e a um tópico) quando documentos são inseridos em uma tabela. Quando linhas são inseridas na tabela GitHub, uma visão materializada é acionada, fazendo com que as linhas sejam reinseridas em um motor Kafka e em um novo tópico. Mais uma vez, isso fica mais claro na ilustração:Crie um novo tópico Kafka Agora crie uma nova visão materializada Se você inserir no tópico github original, criado como parte de Kafka to ClickHouse, os documentos aparecerão como mágica no tópico “github_clickhouse”. Confirme isso com ferramentas nativas do Kafka. Por exemplo, abaixo, inserimos 100 linhas no tópico github usando kcat em um tópico hospedado no Confluent Cloud:Uma leitura no tópico Embora seja um exemplo mais elaborado, isso ilustra o poder das visões materializadas quando usadas em conjunto com o motor Kafka.
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.