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 软件包设置
将以下依赖项添加到包管理框架中:
推荐的 Beam 版本建议从 Apache Beam 2.59.0 版本起使用 ClickHouseIO 连接器。
更早的版本可能无法完全支持该连接器的功能。
这些制品可在官方 Maven 仓库中获取。
以下示例将名为 input.csv 的 CSV 文件读取为 PCollection,再将其转换为 Row 对象 (使用已定义的 schema) ,并通过 ClickHouseIO 将其插入本地 ClickHouse 实例中:
你可以使用以下 setter 函数调整 ClickHouseIO.Write 配置:
使用该连接器时,请注意以下限制:
- 截至目前,仅支持 Sink 操作,不支持 Source 操作。
- 向
ReplicatedMergeTree 或基于 ReplicatedMergeTree 构建的 Distributed 表插入数据时,ClickHouse 会执行去重。未启用复制时,如果插入失败后重试成功,向普通 MergeTree 表插入可能会产生重复数据。不过,每个块的插入都是原子的,并且可以使用 ClickHouseIO.Write.withMaxInsertBlockSize(long) 配置块大小。去重是通过对已插入块的校验和进行比对来实现的。有关去重的更多信息,请参阅 去重 和 插入去重配置。
- 该连接器不会执行任何 DDL 语句;因此,目标表必须在插入前已存在。