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
ReplicatedMergeTreeou em uma tabelaDistributedconstruída sobre umaReplicatedMergeTree. 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 usandoClickHouseIO.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.
- documentação da classe
ClickHouseIOdocumentation. - repositório do GitHub com exemplos clickhouse-beam-connector.