Skip to main content
Apache Beam é um modelo de programação unificado e de código aberto que permite aos desenvolvedores definir e executar pipelines de processamento de dados, tanto em lote quanto em fluxo (contínuo). A flexibilidade do Apache Beam está na sua capacidade de oferecer suporte a uma ampla variedade de cenários de processamento de dados, desde operações de ETL (Extract, Transform, Load) até o processamento complexo de eventos e analytics em tempo real. Esta integração utiliza o conector JDBC oficial do ClickHouse como camada subjacente de inserção.

Pacote de integração

O pacote de integração necessário para integrar o Apache Beam ao ClickHouse é mantido e desenvolvido em Apache Beam I/O Connectors — um pacote de integrations de vários sistemas populares de armazenamento de dados e bancos de dados. A implementação de org.apache.beam.sdk.io.clickhouse.ClickHouseIO está localizada no repositório do Apache Beam.

Configuração do pacote ClickHouse do Apache Beam

Instalação do pacote

Adicione a seguinte dependência ao seu gerenciador de pacotes:
Versão recomendada do BeamO conector ClickHouseIO é recomendado a partir da versão 2.59.0 do Apache Beam. As versões anteriores podem não oferecer suporte completo a todos os recursos do conector.
Os artefatos podem ser encontrados no repositório oficial do Maven.

Exemplo de código

O exemplo a seguir lê um arquivo CSV chamado input.csv como uma PCollection, converte-o em um objeto Row (usando o schema definido) e o insere em uma instância local do ClickHouse com ClickHouseIO:

Tipos de dados suportados

Parâmetros de ClickHouseIO.Write

Você pode ajustar a configuração de ClickHouseIO.Write com as seguintes funções setter:

Limitações

Considere as seguintes limitações ao usar o conector:
  • Até o momento, apenas a operação Sink é compatível. O conector não oferece suporte à operação Source.
  • O ClickHouse realiza desduplicação ao inserir em uma tabela ReplicatedMergeTree ou em uma tabela Distributed construída sobre uma ReplicatedMergeTree. Sem replicação, inserir em uma tabela MergeTree comum pode resultar em duplicatas se uma inserção falhar e depois for repetida com sucesso. No entanto, cada bloco é inserido atomicamente, e o tamanho do bloco pode ser configurado usando ClickHouseIO.Write.withMaxInsertBlockSize(long). A desduplicação é feita usando checksums dos blocos inseridos. Para mais informações sobre desduplicação, consulte Desduplicação e Configuração de desduplicação na inserção.
  • O conector não executa nenhuma instrução DDL; portanto, a tabela de destino deve existir antes da inserção.
Última modificação em 24 de julho de 2026