- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
本文以 Apache SeaTunnel 开源仓库中的 ClickhouseFile Sink 连接器官方文档为主体,结合仓库内 connector-clickhouse 模块的源码实现,系统讲解该连接器的适用场景、工作原理、全部配置参数与完整配置示例。读完本文,你将掌握如何利用 clickhouse-local 程序将数据先生成 ClickHouse 数据文件、再批量迁移进 ClickHouse 集群(bulkload)的完整链路,并能独立完成可落地的作业配置与排障。
一、ClickhouseFile 是什么
ClickhouseFile 是 SeaTunnel 的 Sink 插件(插件标识ClickhouseFile,定义于 ClickhouseFileSinkFactory.java),它的核心思路是:
- 在数据节点(Spark/SeaTunnel 任务所在节点)本地调用clickhouse-local程序,将 SeaTunnel 收到的数据行批量落盘为 ClickHouse 数据文件(part);
- 通过scp/rsync将生成的数据文件传输到 ClickHouse 服务端的数据目录的
detached/目录下; - 最后通过
ALTER TABLE ... ATTACH PART命令将 detached 分区附加回目标表,完成批量数据装载。
整个过程不经过 ClickHouse 的 JDBC/HTTP 写入接口,而是走“文件生成 + 文件迁移 + ATTACH”的离线路径,也就是官方文档所说的bulk load(批量加载),适合大批量数据快速导入的场景。该 Sink 支持**批模式(Batch)和流模式(Streaming)**两种运行方式。
从源码结构看,ClickhouseFile 是一个具备两阶段提交能力的 Sink:它同时实现了SeaTunnelSink(ClickhouseFileSink.java)、SinkWriter(ClickhouseFileSinkWriter.java)与SinkAggregatedCommitter(ClickhouseFileSinkAggCommitter.java),分别承担“本地写文件”、“聚合提交(ATTACH)”,使文件迁移与 ATTACH 具备事务提交语义。
小提示:如果数据量不需要走文件批量装载,也可以使用 JDBC 方式将数据写入 ClickHouse(对应同一连接器模块中的
ClickhouseSink,配置bulk_size、sql等参数,见 ClickhouseConfig.java)。
二、特性与适用前提
特性说明:
- 精准一次(exactly-once):暂不支持。相关特性能力矩阵可参考 connector-v2-features。
硬性前置条件(不满足则无法正常使用):
- 表引擎必须是
Distributed:ClickhouseFile 依赖 Distributed 表将数据路由分发到各个分片节点。在 ClickhouseFileSink.prepare 中,连接器会通过ClickhouseProxy获取表的引擎信息(table.getEngine())并构建ShardMetadata,据此发现集群的所有分片节点。 internal_replication需要设置为true:当 Distributed 表设置了internal_replication = true时,分片内部的数据副本写入由 ClickHouse 自身负责,SeaTunnel 只需将文件投递到分片主节点,避免重复写入副本带来的数据重复问题。- 每个节点都需要可执行 clickhouse-local 程序:因为每个并行任务(subtask)都会调用该程序生成数据文件,所有数据节点上的 clickhouse-local 程序路径必须一致(见参数
clickhouse_local_path)。
三、工作原理:从数据行到 ATTACH 的完整链路
结合 ClickhouseFileSinkWriter.java 源码,一次完整的写入生命周期如下:
3.1 Writer 初始化:发现分片与本地数据目录
Writer 创建时(ClickhouseFileSinkWriter 构造函数):
- 通过
ClickhouseProxy连接默认节点,用ShardRouter根据 Distributed 表的元数据发现所有分片; - 逐个分片查询其本地表(local table)的
data_paths,得到每台 ClickHouse 服务器上数据文件落盘的真实目录(后续文件传输的目标目录); - 若
node_free_password = false,会执行nodePasswordCheck(),为每个分片节点查找node_pass中配置的密码,找不到则抛出PASSWORD_NOT_FOUND_IN_SHARD_NODE异常,防止后续 scp/rsync 因无凭据而失败。
3.2 write:按分片写入本地临时 CSV 文件
每次write(SeaTunnelRow)(源码 L116-L146)执行:
- 用
shardRouter.getShard(element)决定该行数据发往哪个分片(分片规则见sharding_key一节); - 为该分片在
file_temp_path下创建一个以 UUID(10 位)命名的临时目录,并打开local_data.log文件(CLICKHOUSE_LOCAL_FILE_SUFFIX = "/local_data.log")写入; - 行数据按表字段顺序、以
file_fields_delimiter拼接成一行文本(空字段写空串),通过FileChannel.map内存映射缓冲区(每次映射 128KB)高效追加写入(saveDataToFile)。
3.3 prepareCommit:clickhouse-local 生成数据文件并迁移到服务端
prepareCommit()(源码 L172-L201)是核心步骤:
关闭所有文件通道,确保数据全部落盘;
对每个分片的临时文件调用
generateClickhouseLocalFiles(),通过ProcessBuilder执行 clickhouse-local 命令(源码 L254-L379)。生成的命令大致为:clickhouse local --file <file_temp_path>/<uuid>/local_data.log \ --format_csv_delimiter "<file_fields_delimiter>" \ -S "字段1 类型1,字段2 类型2,..." \ -N "temp_table<uuid>" \ -q "<DDL>; INSERT INTO TABLE <local_table> SELECT ... FROM temp_table<uuid>;" \ --path "<file_temp_path>/<uuid>" # 或 --config-file(compatible_mode 时)其中:
-S指定临时表的 schema(字段类型取自 ClickHouse 表结构);-N指定临时表名,避免并发任务冲突;-q中的 DDL 来自目标本地表的建表语句(经adjustClickhouseDDL()去掉库名前缀与反引号,并过滤storage_policy等不适用于 clickhouse-local 的 SETTINGS,见 源码 L413-L435);- 生成的 part 目录会追加
_<subtaskIndex>后缀重命名,避免不同并行子任务生成的数据 part 名称冲突(对应变更日志中的 “数据 part 名称冲突修复”);
校验生成路径
.../data/_local/<local_table>/下的 part 目录存在(不存在抛FILE_NOT_EXISTS),得到 part 列表;调用
moveClickhouseLocalFileToServer()(源码 L381-L393):根据分片随机选择该节点的一个data_path,通过FileTransferFactory创建的ScpFileTransfer或RsyncFileTransfer把 part 传输到远端.../<data_path>/detached/目录;传输完成后立即清理本地临时目录(
clearLocalFileDirectory)。
其中 scp 传输采用 Apache MINA sshd 客户端(ScpFileTransfer.java),传输后还会在远端执行chown命令:ls -l <父目录> | tail -n 1 | awk '{print $3}' | xargs -i chown -R {}:{} <detached 目录>,将文件属主改为 ClickHouse 服务进程用户——因为只有文件属主与 ClickHouse 服务用户一致,后续ATTACH才能成功。
3.4 AggregatedCommitter commit:ALTER TABLE ATTACH PART
所有并行 Writer 的提交信息(CKFileCommitInfo,即“分片 -> part 文件列表”映射)汇总到聚合提交器后,combine()合并、commit()执行真正的入库(ClickhouseFileSinkAggCommitter.java):
ALTER TABLE <local_table_name> ATTACH PART '<part名称>'即对每个分片、每个 part 文件执行 ATTACH,将 detached 目录中的 part 附加到本地表。至此数据正式进入 ClickHouse,后续由 ClickHouse 的分片/副本机制与 Distributed 表对外提供查询。
四、参数详解
以下参数表与官方文档一致,并补充了源码中的校验逻辑、默认值与取值约束(源码依据:ClickhouseConfig.java,必填/可选规则见 ClickhouseFileSinkFactory.optionRule):
| 名称 | 类型 | 是否必须 | 默认值 |
|---|---|---|---|
| host | string | 是 | - |
| database | string | 是 | - |
| table | string | 是 | - |
| username | string | 是 | - |
| password | string | 是 | - |
| clickhouse_local_path | string | 是 | - |
| sharding_key | string | 否 | - |
| copy_method | string | 否 | scp |
| node_free_password | boolean | 否 | false |
| node_pass | list | 否 | - |
| node_pass.node_address | string | 否 | - |
| node_pass.username | string | 否 | "root" |
| node_pass.password | string | 否 | - |
| compatible_mode | boolean | 否 | false |
| file_fields_delimiter | string | 否 | "\t" |
| file_temp_path | string | 否 | "/tmp/seatunnel/clickhouse-local/file" |
| common-options | - | 否 | - |
必填项校验逻辑在 ClickhouseFileSink.prepare 中通过
CheckConfigUtil.checkAllExists强制执行,缺失任一必填项都会抛出CONFIG_VALIDATION_FAILED异常。
host [string]
ClickHouse集群地址,格式为host:port,允许同时指定多个 host,用逗号分隔。例如:"host1:8123,host2:8123"。连接器会用该地址构建 ClickHouse 节点列表,并以第一个节点作为默认节点来拉取表结构与分片元数据。
database [string]
ClickHouse数据库名。
table [string]
表名称。注意:该表必须是Distributed引擎,且internal_replication置为true。
username [string]
连接ClickHouse的用户名。
password [string]
连接ClickHouse的用户密码。
clickhouse_local_path [string]
数据节点上 clickhouse-local 程序的路径。由于每个任务(每个 subtask)都会调用它,所以每个数据节点上的 clickhouse-local 路径必须完全相同。源码中该路径还会按空格切分后拼入命令(generateClickhouseLocalFiles),因此路径中不要包含空格(示例中使用了"/Users/seatunnel/Tool/clickhouse local",实际生产环境建议使用无空格的绝对路径)。
sharding_key [string]
当 ClickhouseFile 需要对数据分片(split data)时,决定“某一行数据发往哪个分片节点”是关键问题。默认情况下采用随机算法选择分片;指定sharding_key后,该字段将作为分片算法的依据(源码中会从表结构中解析该字段的类型,构造ShardMetadata供ShardRouter使用)。配置后,SeaTunnel 会按照与 ClickHouse Distributed 表分片键一致的语义把数据行路由到对应节点,保证发往各节点本地表的文件与集群分片规则对齐。
copy_method [string]
文件传输方式,默认scp,可选值为scp和rsync(枚举定义见 ClickhouseFileCopyMethod.java,工厂分发见 FileTransferFactory.java)。传入其他值会抛出ILLEGAL_ARGUMENT异常。选择 rsync 时,需要数据节点上装有 rsync 命令。
node_free_password [boolean]
由于 SeaTunnel 需要使用 scp 或 rsync 进行文件传输,因此 SeaTunnel 需要 ClickHouse 服务端的访问权限。如果每个数据节点与 ClickHouse 服务端都配置了免密登录(SSH key 互信),可将此项配置为true;否则必须在node_pass参数中配置对应节点的密码(默认false,即要求提供节点密码)。
node_pass [list]
用来保存所有 ClickHouse 服务器地址及其对应的访问密码。Writer 初始化时会对每个分片节点校验node_pass中是否存在其密码,缺失则任务失败(见 3.1 节nodePasswordCheck)。配置格式为对象列表,每个对象包含node_address、username、password。
node_pass.node_address [string]
ClickHouse 服务器节点地址,需与分片元数据中的节点 host 匹配。
node_pass.username [string]
ClickHouse 服务器节点用户名,默认为root。源码中解析该字段时若缺失则回退为"root"(ClickhouseFileSink.prepare 中 nodeUser 的构建)。
node_pass.password [string]
ClickHouse 服务器节点的访问密码,用于 scp/rsync 认证。
compatible_mode [boolean]
在低版本的 ClickHouse 中,clickhouse-local 程序不支持--path参数,此时需要设置该参数为true,采用替代方式实现--path的功能:连接器会在临时目录下生成一份最小化的config.xml(模板见 ClickhouseFileSinkWriter 的 CK_LOCAL_CONFIG_TEMPLATE,包含<path>配置项),并通过--config-file参数传给 clickhouse-local,从而间接指定数据目录(源码 L305-L319)。
file_fields_delimiter [string]
ClickhouseFile 使用 CSV 格式临时保存数据(即写local_data.log时的字段分隔符)。如果数据本身包含 CSV 的分隔符,可能导致程序解析异常。使用此配置可以避免该情况。配置的值必须正好为一个字符的长度,源码中会在prepare阶段强制校验,长度不为 1 时抛出FILE_FIELDS_DELIMITER must be a single character异常(ClickhouseFileSink.java L169-L173)。默认"\t"(制表符)。
file_temp_path [string]
ClickhouseFile 本地存储临时文件的目录(生成local_data.log、clickhouse-local 数据目录、compatible_mode 下的 config.xml 都位于该目录下),默认/tmp/seatunnel/clickhouse-local/file。请确保该目录所在磁盘空间充足且数据节点间各自独立。
common-options
Sink 插件常用参数(如result_table_name、parallelism等),请参考 Sink 常用选项 获取更多细节。
五、完整配置示例
以下示例取自官方文档(字段、缩进均为可直接使用的 HOCON 格式,ClickhouseFile作为 sink 段名):
ClickhouseFile { host = "192.168.0.1:8123" database = "default" table = "fake_all" username = "default" password = "" clickhouse_local_path = "/Users/seatunnel/Tool/clickhouse local" sharding_key = "age" node_free_password = false node_pass = [{ node_address = "192.168.0.1" password = "seatunnel" }] }结合前文参数说明,一个更贴合生产环境、覆盖多分片与分隔符场景的示例:
ClickhouseFile { host = "ck-node1:8123,ck-node2:8123" database = "dwd" table = "dwd_orders_all" username = "default" password = "your_password" clickhouse_local_path = "/opt/clickhouse/clickhouse-local" sharding_key = "order_id" copy_method = "scp" node_free_password = true file_fields_delimiter = "\t" file_temp_path = "/tmp/seatunnel/clickhouse-local/file" compatible_mode = false }使用说明:若数据节点与 ClickHouse 服务端已配置 SSH 免密,可将
node_free_password置为true并省略node_pass;否则必须为每个分片节点在node_pass中提供node_address(与分片节点 host 一致)与password。
六、注意事项与排障要点
- 表引擎约束:目标表必须是
Distributed引擎且internal_replication = true,否则文件分发到分片的行为不符合预期。 - clickhouse-local 版本一致性:低版本 ClickHouse(clickhouse-local 不支持
--path)请开启compatible_mode;同时所有数据节点上的 clickhouse-local 路径必须一致。 - 文件属主问题:ATTACH 成功的前提是 detached 目录中文件属主与 ClickHouse 服务进程用户一致,连接器通过 scp/rsync 传输后的远端
chown自动处理;若人工介入排查 ATTACH 失败,优先检查属主。 - 分隔符冲突:若业务数据包含默认的
\t或 CSV 逗号,务必设置一个数据中不出现且长度为 1 的file_fields_delimiter。 - 磁盘空间:
file_temp_path会同时存放 CSV 临时文件与 clickhouse-local 生成的数据文件(一份完整数据两份落盘),需要预留足够空间。 - 并行 part 冲突:不同 subtask 生成的 part 已通过追加
_<subtaskIndex>后缀规避命名冲突(对应变更日志的 BugFix),升级到含该修复的版本即可。
七、变更日志
2.2.0-beta 2022-09-26
- 支持将数据写入 ClickHouse 文件并迁移到 ClickHouse 数据目录(本连接器的首次发布,即 bulkload 能力)。
随后版本
- [BugFix] 修复生成的数据 part 名称冲突 Bug,并改进文件提交逻辑(PR #3416)。
- [Feature] 支持
compatible_mode,兼容低版本 ClickHouse(PR #3416)。
- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
相关推荐
SeaTunnel ClickhouseFile Sink 完全指南:基于 clickhouse-local 数据文件的 ClickHouse 批量加载
SeaTunnel ClickhouseFile Sink 完全指南:基于 clickhouse local 数据文件的 ClickHouse 批量加载 本文基
数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel ClickhouseFile Sink Connector 实战指南:基于 clickhouse-local 的分布式批量导入
SeaTunnel ClickhouseFile Sink Connector 实战指南:基于 clickhouse local 的分布式批量导入 Clickh
数据工程大数据批处理流处理SeaTunnel ClickhouseFile 接收器:基于 clickhouse-local 的 ClickHouse Bulk Load 数据接入实战
SeaTunnel ClickhouseFile 接收器:基于 clickhouse local 的 ClickHouse Bulk Load 数据接入实战 本
数据集成ETL大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考