支持 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 语句;因此,目标表必须在插入前已存在。
ClickHouseIO类的文档。- 示例
Github仓库:clickhouse-beam-connector。