Skip to content
ClickHouse Docs
ClickHouse DocsClickHouse Docs

将 Apache Beam 与 ClickHouse 集成

支持 ClickHouse

Apache Beam 是一种开源的统一编程模型,使开发者能够定义并执行批次和 stream (连续) 数据处理管道。Apache Beam 的灵活性在于,它支持广泛的数据处理场景,从 ETL (提取、转换、加载) 操作到复杂事件处理和实时分析。 此集成使用 ClickHouse 官方的 JDBC 连接器 作为底层插入机制。

集成包

用于集成 Apache Beam 和 ClickHouse 的集成包由 Apache Beam I/O Connectors 维护和开发;这是一个包含多种流行数据存储系统和数据库的集成组件集合。 org.apache.beam.sdk.io.clickhouse.ClickHouseIO 的实现位于 Apache Beam repo 中。

Apache Beam ClickHouse 软件包设置

软件包安装

将以下依赖项添加到包管理框架中:

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

这些制品可在官方 Maven 仓库中获取。

代码示例

以下示例将名为 input.csv 的 CSV 文件读取为 PCollection,再将其转换为 Row 对象 (使用已定义的 schema) ,并通过 ClickHouseIO 将其插入本地 ClickHouse 实例中:


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) {
        // 创建一个 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();

        // 对管道应用转换操作。
        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"));

        // 运行管道。
        p.run().waitUntilFinish();
    }
}

支持的数据类型

ClickHouse Apache Beam 是否支持 说明
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 是一种 LogicalType,表示固定长度的
字节数组,位于
org.apache.beam.sdk.schemas.logicaltypes
Schema.TypeName#DECIMAL
Schema.TypeName#MAP

ClickHouseIO.Write 参数

你可以使用以下 setter 函数调整 ClickHouseIO.Write 配置:

参数设置函数 参数类型 默认值 说明
withMaxInsertBlockSize (long maxInsertBlockSize) 1000000 单次插入的行块最大大小。
withMaxRetries (int maxRetries) 5 插入失败时的最大重试次数。
withMaxCumulativeBackoff (Duration maxBackoff) Duration.standardDays(1000) 重试的最大累计退避耗时。
withInitialBackoff (Duration initialBackoff) Duration.standardSeconds(5) 首次重试前的初始退避耗时。
withInsertDistributedSync (Boolean sync) true 如果为 true,则会同步分布式表的插入操作。
withInsertQuorum (Long quorum) null 确认一次插入操作所需的副本数。
withInsertDeduplicate (Boolean deduplicate) true 如果为 true,则会对插入操作启用去重。
withTableSchema (TableSchema schema) null 目标 ClickHouse 表的 schema。

限制

使用该连接器时,请注意以下限制:

  • 截至目前,仅支持 Sink 操作,不支持 Source 操作。
  • ReplicatedMergeTree 或基于 ReplicatedMergeTree 构建的 Distributed 表插入数据时,ClickHouse 会执行去重。未启用复制时,如果插入失败后重试成功,向普通 MergeTree 表插入可能会产生重复数据。不过,每个块的插入都是原子的,并且可以使用 ClickHouseIO.Write.withMaxInsertBlockSize(long) 配置块大小。去重是通过对已插入块的校验和进行比对来实现的。有关去重的更多信息,请参阅 去重插入去重配置
  • 该连接器不会执行任何 DDL 语句;因此,目标表必须在插入前已存在。
Navigation