Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Integração do Apache Beam com o ClickHouse

Suportado pelo ClickHouse

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:

<dependency>
    <groupId>org.apache.beam</groupId>
    <artifactId>beam-sdks-java-io-clickhouse</artifactId>
    <version>${beam.version}</version>
</dependency>

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:


package org.example;

import org.apache.beam.sdk.Pipeline;
import org.apache.beam.sdk.io.TextIO;
import org.apache.beam.sdk.io.clickhouse.ClickHouseIO;
import org.apache.beam.sdk.schemas.Schema;
import org.apache.beam.sdk.transforms.DoFn;
import org.apache.beam.sdk.transforms.ParDo;
import org.apache.beam.sdk.values.PCollection;
import org.apache.beam.sdk.values.Row;
import org.joda.time.DateTime;

public class Main {

    public static void main(String[] args) {
        // Cria um objeto Pipeline.
        Pipeline p = Pipeline.create();

        Schema SCHEMA =
                Schema.builder()
                        .addField(Schema.Field.of("name", Schema.FieldType.STRING).withNullable(true))
                        .addField(Schema.Field.of("age", Schema.FieldType.INT16).withNullable(true))
                        .addField(Schema.Field.of("insertion_time", Schema.FieldType.DATETIME).withNullable(false))
                        .build();

        // Aplica transformações ao pipeline.
        PCollection<String> lines = p.apply("ReadLines", TextIO.read().from("src/main/resources/input.csv"));

        PCollection<Row> rows = lines.apply("ConvertToRow", ParDo.of(new DoFn<String, Row>() {
            @ProcessElement
            public void processElement(@Element String line, OutputReceiver<Row> out) {

                String[] values = line.split(",");
                Row row = Row.withSchema(SCHEMA)
                        .addValues(values[0], Short.parseShort(values[1]), DateTime.now())
                        .build();
                out.output(row);
            }
        })).setRowSchema(SCHEMA);

        rows.apply("Write to ClickHouse",
                        ClickHouseIO.write("jdbc:clickhouse://localhost:8123/default?user=default&password=******", "test_table"));

        // Executa o pipeline.
        p.run().waitUntilFinish();
    }
}

Tipos de dados suportados

ClickHouse Apache Beam Compatível Observações
TableSchema.TypeName.FLOAT32 Schema.TypeName#FLOAT
TableSchema.TypeName.FLOAT64 Schema.TypeName#DOUBLE
TableSchema.TypeName.INT8 Schema.TypeName#BYTE
TableSchema.TypeName.INT16 Schema.TypeName#INT16
TableSchema.TypeName.INT32 Schema.TypeName#INT32
TableSchema.TypeName.INT64 Schema.TypeName#INT64
TableSchema.TypeName.STRING Schema.TypeName#STRING
TableSchema.TypeName.UINT8 Schema.TypeName#INT16
TableSchema.TypeName.UINT16 Schema.TypeName#INT32
TableSchema.TypeName.UINT32 Schema.TypeName#INT64
TableSchema.TypeName.UINT64 Schema.TypeName#INT64
TableSchema.TypeName.DATE Schema.TypeName#DATETIME
TableSchema.TypeName.DATETIME Schema.TypeName#DATETIME
TableSchema.TypeName.ARRAY Schema.TypeName#ARRAY
TableSchema.TypeName.ENUM8 Schema.TypeName#STRING
TableSchema.TypeName.ENUM16 Schema.TypeName#STRING
TableSchema.TypeName.BOOL Schema.TypeName#BOOLEAN
TableSchema.TypeName.TUPLE Schema.TypeName#ROW
TableSchema.TypeName.FIXEDSTRING FixedBytes FixedBytes é um LogicalType que representa um array de bytes
de comprimento fixo, localizado em
org.apache.beam.sdk.schemas.logicaltypes
Schema.TypeName#DECIMAL
Schema.TypeName#MAP

Parâmetros de ClickHouseIO.Write

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

Função setter de parâmetro Tipo de argumento Valor padrão Descrição
withMaxInsertBlockSize (long maxInsertBlockSize) 1000000 Tamanho máximo de um bloco de linhas a serem inseridas.
withMaxRetries (int maxRetries) 5 Número máximo de tentativas para inserções com falha.
withMaxCumulativeBackoff (Duration maxBackoff) Duration.standardDays(1000) Duração máxima acumulada de backoff para tentativas.
withInitialBackoff (Duration initialBackoff) Duration.standardSeconds(5) Duração do backoff inicial antes da primeira tentativa.
withInsertDistributedSync (Boolean sync) true Se true, sincroniza as operações de inserção em tabelas distribuídas.
withInsertQuorum (Long quorum) null Número de réplicas necessário para confirmar uma operação de inserção.
withInsertDeduplicate (Boolean deduplicate) true Se true, a desduplicação é ativada para operações de inserção.
withTableSchema (TableSchema schema) null Schema da tabela ClickHouse de destino.

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.
Navigation