Pular para o conteúdo principal
Este é o Conector Sink oficial do Apache Flink, com suporte da ClickHouse. Ele foi desenvolvido com o AsyncSinkBase do Flink e o Java client oficial do ClickHouse. O conector oferece suporte à API DataStream do Apache Flink. O suporte à Table API está planejado para um lançamento futuro.

Requisitos

  • Java 11+ (para o Flink 1.17+) ou 17+ (para o Flink 2.0+)
  • Apache Flink 1.17+
O conector foi dividido em dois artefatos para oferecer suporte ao Flink 1.17+ e ao Flink 2.0+. Escolha o artefato correspondente à versão do Flink que você deseja usar:
O conector não foi testado com versões do Flink anteriores à versão 1.17.2

Instalação e configuração

Importar como dependência

Baixe o binário

O padrão de nomenclatura do arquivo JAR binário é:
onde: Você pode encontrar todos os arquivos JAR disponíveis já lançados no Maven Central Repository.

Como usar a API do DataStream

Trecho

Digamos que você queira inserir dados CSV brutos no ClickHouse:
Mais exemplos e trechos de código podem ser encontrados em nossos testes:

Exemplo de início rápido

Criamos um exemplo baseado em Maven para facilitar os primeiros passos com o ClickHouse Sink: Para instruções mais detalhadas, consulte o Guia de exemplos

Opções de conexão com a API DataStream

Opções do cliente ClickHouse

options e serverSettings devem ser passados ao cliente como Map<String, String>. Um mapa vazio em qualquer um deles usará os padrões do cliente ou do servidor, respectivamente.
Todas as opções disponíveis do Java client estão listadas em ClientConfigProperties.java e nesta página da documentação.Todas as configurações de sessão disponíveis do servidor estão listadas nesta página da documentação.
Por exemplo:

Opções do sink

As opções a seguir vêm diretamente do AsyncSinkBase do Flink:

Tipos de dados compatíveis

A tabela abaixo traz uma referência rápida para a conversão de tipos de dados ao inserir dados do Flink no ClickHouse. Observações:
  • Um ZoneId deve ser fornecido ao realizar operações com data.
  • Precisão e escala devem ser fornecidas ao realizar operações decimais.
  • Para que o ClickHouse consiga interpretar uma String Java como JSON, é necessário habilitar enableJsonSupportAsString em ClickHouseClientConfig.
  • O conector requer um ElementConvertor para mapear elementos no DataStream de entrada para payloads do ClickHouse. Para isso, o conector fornece ClickHouseConvertor e POJOConvertor, que podem ser usados para implementar esse mapeamento com os métodos de serialização de DataWriter acima.

Formatos de entrada suportados

Você pode encontrar a lista de formatos de entrada disponíveis do ClickHouse nesta página da documentação e em ClickHouseFormat.java. Para especificar o formato que o conector deve usar para serializar seu DataStream como payloads para o ClickHouse, use a função setClickHouseFormat. Por exemplo:
Por padrão, o conector usará RowBinaryWithDefaults ou RowBinary caso setSupportDefault em ClickHouseClientConfig seja explicitamente definido como true ou false, respectivamente.

Métricas

O conector expõe as seguintes métricas adicionais, além das métricas já existentes do Flink:

Limitações

  • No momento, o sink oferece uma garantia de entrega at-least-once. O suporte à semântica exactly-once está sendo acompanhado aqui.
  • O sink ainda não oferece suporte a uma fila de dead-letter (DLQ) para armazenar temporariamente registros que não podem ser processados. Enquanto isso, o conector tentará reinserir os registros com falha e os descartará em caso de insucesso. Esse recurso está sendo acompanhado aqui.
  • O sink ainda não oferece suporte à criação por meio da Table API do Flink ou do Flink SQL. Esse recurso está sendo acompanhado aqui.

Compatibilidade de versões do ClickHouse e segurança

  • O conector é testado diariamente, por meio de um workflow de CI, com uma variedade de versões recentes do ClickHouse, incluindo latest e head. As versões testadas são atualizadas periodicamente à medida que novos lançamentos do ClickHouse entram em atividade. Veja aqui as versões com as quais o conector é testado diariamente.
  • Consulte a política de segurança do ClickHouse para ver vulnerabilidades de segurança conhecidas e como relatar uma vulnerabilidade.
  • Recomendamos atualizar o conector continuamente para não perder correções de segurança e outras melhorias.
  • Se você tiver algum problema com a migração, crie uma issue no GitHub e responderemos!
  • Para obter o melhor desempenho, garanta que o tipo de elemento do seu DataStream não seja um tipo genérico — veja aqui a distinção de tipos do Flink. Elementos não genéricos evitam a sobrecarga de serialização do Kryo e melhoram a vazão para o ClickHouse.
  • Recomendamos definir maxBatchSize para pelo menos 1000 e, idealmente, entre 10.000 e 100.000. Veja este guia sobre inserções em massa para mais informações.
  • Para fazer desduplicação no estilo OLTP ou upsert no ClickHouse, consulte esta página da documentação. Observação: isso não deve ser confundido com a desduplicação em lote que ocorre em novas tentativas.

Solução de problemas

CANNOT_READ_ALL_DATA

O erro a seguir pode ocorrer:
Causa: Na maioria dos casos, o erro CANNOT_READ_ALL_DATA significa que o schema da sua table do ClickHouse divergiu do schema do registro no Flink. Isso pode acontecer quando um deles é alterado de uma forma incompatível com versões anteriores. Solução: Atualize o schema da sua table do ClickHouse ou o tipo de dado de entrada do conector (ou ambos) para que sejam compatíveis. Se necessário, consulte o mapeamento de tipos para ver como mapear tipos Java para tipos do ClickHouse. Observação: se ainda houver registros em trânsito, você precisará redefinir o state do Flink ao reiniciar o conector.

Baixa vazão

Você pode notar que a vazão do conector não escala com o paralelismo do job (número de tasks do Flink) ao gravar no ClickHouse. Causa: o processo de merge de parts em segundo plano do ClickHouse pode estar reduzindo a velocidade das inserções. Isso pode acontecer quando o tamanho de lote configurado é muito pequeno, o conector está fazendo flush com muita frequência, ou por uma combinação dos dois fatores. Solução: monitore as métricas numRequestSubmitted e actualRecordsPerBatch para ajudar a determinar como ajustar o tamanho do lote (maxBatchSize) e a frequência de flush. Além disso, consulte Uso avançado e recomendado para recomendações de dimensionamento de lote.

Faltam linhas na minha tabela do ClickHouse

Causa: O(s) lote(s) foi(foram) descartado(s) devido a uma falha não recuperável ou porque não pôde(ram) ser inserido(s) dentro do número configurado de tentativas (configurável via ClickHouseClientConfig.setNumberOfRetries()). Observação: por padrão, o conector tentará reinserir um lote em até 3 tentativas antes de descartá-lo. Solução: Inspecione os logs do TaskManager e/ou os stack traces para identificar a causa raiz.

Contribuição e suporte

Se você quiser contribuir com o projeto ou relatar algum problema, sua colaboração será muito bem-vinda! Visite nosso repositório no GitHub para abrir uma issue, sugerir melhorias ou enviar um pull request. Contribuições são bem-vindas! Consulte o guia de contribuição no repositório antes de começar. Obrigado por ajudar a melhorar o conector do ClickHouse para Flink!
Última modificação em 25 de junho de 2026