Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

Integración entre Apache Beam y ClickHouse

Compatible con ClickHouse

Apache Beam es un modelo de programación unificado y de código abierto que permite a los desarrolladores definir y ejecutar pipelines de procesamiento de datos tanto por lotes como en flujo (continuo). La flexibilidad de Apache Beam radica en su capacidad para admitir una amplia variedad de escenarios de procesamiento de datos, desde operaciones ETL (Extract, Transform, Load) hasta el procesamiento complejo de eventos y la analítica en tiempo real. Esta integración utiliza el conector JDBC oficial de ClickHouse como capa subyacente de inserción.

Paquete de integración

El paquete de integración necesario para integrar Apache Beam y ClickHouse se mantiene y desarrolla en Apache Beam I/O Connectors, un conjunto de integraciones para muchos sistemas populares de almacenamiento de datos y bases de datos. La implementación de org.apache.beam.sdk.io.clickhouse.ClickHouseIO se encuentra en el repositorio de Apache Beam.

Configuración del paquete Apache Beam para ClickHouse

Instalación del paquete

Agregue la siguiente dependencia a su sistema de gestión de paquetes:

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

Los artefactos están disponibles en el repositorio oficial de Maven.

Ejemplo de código

El siguiente ejemplo lee un archivo CSV llamado input.csv como una PCollection, lo convierte en un objeto de tipo Row (usando el esquema definido) y lo inserta en una instancia local de ClickHouse mediante 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) {
        // Crear un 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();

        // Aplicar transformaciones al 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"));

        // Ejecutar el pipeline.
        p.run().waitUntilFinish();
    }
}

Tipos de datos compatibles

ClickHouse Apache Beam Compatible Notas
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 es un LogicalType que representa un array de bytes
de longitud fija ubicado en
org.apache.beam.sdk.schemas.logicaltypes
Schema.TypeName#DECIMAL
Schema.TypeName#MAP

Parámetros de ClickHouseIO.Write

Puede ajustar la configuración de ClickHouseIO.Write con las siguientes funciones setter:

Función setter de parámetro Tipo de argumento Valor predeterminado Descripción
withMaxInsertBlockSize (long maxInsertBlockSize) 1000000 Tamaño máximo de un bloque de filas que se puede insertar.
withMaxRetries (int maxRetries) 5 Número máximo de reintentos para las inserciones fallidas.
withMaxCumulativeBackoff (Duration maxBackoff) Duration.standardDays(1000) Duración acumulada máxima del backoff para los reintentos.
withInitialBackoff (Duration initialBackoff) Duration.standardSeconds(5) Duración del backoff inicial antes del primer reintento.
withInsertDistributedSync (Boolean sync) true Si es true, sincroniza las operaciones de inserción para las tablas distribuidas.
withInsertQuorum (Long quorum) null Número de réplicas necesarias para confirmar una operación de inserción.
withInsertDeduplicate (Boolean deduplicate) true Si es true, la deduplicación está habilitada para las operaciones de inserción.
withTableSchema (TableSchema schema) null Esquema de la tabla de ClickHouse de destino.

Limitaciones

Tenga en cuenta las siguientes limitaciones al usar el conector:

  • A día de hoy, solo se admite la operación Sink. El conector no admite la operación Source.
  • ClickHouse realiza la deduplicación al insertar en una tabla ReplicatedMergeTree o en una tabla Distributed construida sobre ReplicatedMergeTree. Sin replicación, insertar en una MergeTree normal puede dar lugar a duplicados si una inserción falla y luego se reintenta correctamente. Sin embargo, cada bloque se inserta de forma atómica, y el tamaño del bloque se puede configurar mediante ClickHouseIO.Write.withMaxInsertBlockSize(long). La deduplicación se logra mediante el uso de sumas de verificación de los bloques insertados. Para obtener más información sobre la deduplicación, visite Deduplication y Configuración de inserción con deduplicación.
  • El conector no ejecuta ninguna sentencia DDL; por lo tanto, la tabla de destino debe existir antes de la inserción.
Navigation