SeaTunnel ClickhouseFile Sink 完全指南:基于 clickhouse-local 数据文件的 ClickHouse 批量加载

发布时间:2026/9/16 14:58:05
SeaTunnel ClickhouseFile Sink 完全指南:基于 clickhouse-local 数据文件的 ClickHouse 批量加载 SeaTunnel ClickhouseFile Sink 完全指南基于 clickhouse-local 数据文件的 ClickHouse 批量加载【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 官方文档 ClickhouseFile 与连接器源码展开系统讲解 ClickhouseFile Sink 连接器的工作原理、全部配置项含义、clickhouse-local 数据文件的生成与传输机制以及基于ATTACH PART的提交逻辑。读完本文你能够正确配置该连接器完成对 ClickHouse 分布式表的大规模批量写入bulk load理解从数据落盘、远端文件传输到数据 part 挂载的完整链路并掌握compatible_mode、node_pass、file_fields_delimiter等关键选项背后的实现细节与适用约束。概述与适用前提ClickhouseFile 是 SeaTunnel 提供的一种面向 ClickHouse 的文件型 Sink 连接器。它的核心思路与常规 JDBC 逐行写入不同先在 SeaTunnel 计算节点本地调用clickhouse-local程序将数据写成 ClickHouse 本地数据文件local data part再把这些文件通过 scp/rsync 传输到 ClickHouse 服务器节点的数据目录detached/目录下最后执行ATTACH PART语句把 part 挂载进表即官方文档所称的 bulk load 批量加载模式。根据官方文档的描述该连接器有以下明确的适用前提与能力边界仅支持引擎为Distributed的 ClickHouse 表且建表时internal_replication选项应为true。这是因为连接器最终操作的是Distributed表背后的本地表local table文件需要挂载到本地表的数据目录中同时支持 Batch 与 Streaming 模式不支持 exactly-once 语义官方文档的 Key features 中该项未勾选。文件是先转移、后挂载的两阶段过程失败重试可能产生重复 part如果只需要常规的 JDBC 写入方式SeaTunnel 同样提供了名为Clickhouse的 JDBC Sink 连接器对应源码中的 ClickhouseSink可作为替代方案。从源码结构看该连接器在 SeaTunnel 插件体系中的注册名为ClickhouseFileClickhouseFileSink 的getPluginName()返回ClickhouseFile而 ClickhouseFileSinkFactory 通过factoryIdentifier()返回同一标识两者一致才能保证作业配置中的ClickhouseFile { ... }被正确路由到该工厂。端到端工作原理整个数据流可以分为四个阶段每个阶段都能在源码中找到对应实现阶段一行数据落盘为 CSV 临时文件ClickhouseFileSinkWriter 的write()方法L122-L151处理每条SeaTunnelRow先经ShardRouter根据sharding_key未配置则随机将行路由到目标 shard每个 shard 维护一个临时文件file_temp_path/uuid/local_data.loguuid 为 10 位随机串行数据按表字段顺序拼接、以file_fields_delimiter分隔并追加换行写入使用MappedByteBuffer内存映射缓冲缓冲区大小固定为1024 * 128字节见 L79 的bufferSize缓冲写满后重新映射文件尾部继续写入。从源码结构看这里把行数据组织为 CSV 文本而非 ClickHouse 原生二进制格式是为了让后续clickhouse-local命令以--format_csv_delimiter方式统一解析字段类型则通过命令的-S参数显式声明这一点在生成命令时会用到完整表 schema见下节。阶段二调用 clickhouse-local 生成数据文件prepareCommit()L178-L206在 checkpoint 触发时执行关闭各 shard 的文件通道后调用generateClickhouseLocalFiles()拼出如下形式的命令并通过bash -c执行L259-L337clickhouse_local_path local \ --file file_temp_path/uuid/local_data.log \ --format_csv_delimiter file_fields_delimiter \ -S field1 type1,field2 type2,... \ -N temp_tableuuid \ -d _local \ -n \ -q 建表DDL; INSERT INTO TABLE local_table SELECT 字段列表 FROM temp_tableuuid; \ --path file_temp_path/uuid # 或 --config-filecompatible_mode 下命令中几个值得注意的细节均来自 ClickhouseFileSinkWriter 源码-S参数按目标表的 schema 声明所有列及其 ClickHouse 类型-q中INSERT ... SELECT的字段列表里凡是表中有而输入行中没有的列会以NULL占位即支持向含物化列的表写入只写部分字段建表 DDL 会经过adjustClickhouseDDL()L421-L443处理去掉库名前缀与反引号并从SETTINGS子句中过滤掉storage_policy这类与本地临时实例不兼容的设置避免 clickhouse-local 建表失败若clickhouse_local_path配置了空格分隔的多个可执行文件路径命令前缀会带上local子命令以适配不同部署形态的 clickhouse-local命令执行完成后连接器会检查file_temp_path/uuid/data/_local/local_table/目录是否存在并列出其中的 part 子目录排除detached然后把每个 part 目录重命名为part名_subtask序号以降低并发任务间的文件名冲突重命名失败仅告警继续compatible_mode的实现细节老版本 clickhouse-local 不支持--path参数开启该模式后连接器不再传--path而是向临时目录生成一份config.xml模板见 L65-L67 的CK_LOCAL_CONFIG_TEMPLATE其中path指向当前 uuid 临时目录再通过--config-file传入以此间接实现指定数据目录的效果。阶段三传输文件到 ClickHouse 服务器moveClickhouseLocalFileToServer()L388-L401把每个 part 目录上传到目标 shard 节点传输方式由copy_method决定FileTransferFactory 仅创建ScpFileTransfer或RsyncFileTransfer两种实现目标路径不是随意指定的ClickhouseFileSinkWriter构造时会逐 shard 连接对应节点查询该 shard 本地表的data_paths然后随机挑一个数据目录把文件传进其下的detached/目录ScpFileTransfer 基于 Apache SSHD 客户端实现固定使用22 端口L46支持密码认证node_pass中的 password与公钥认证key_path通过FileKeyPairProvider加载 RSA 私钥上传选项为Recursive TargetIsDirectory PreserveAttributes上传后还会在远端执行一段 shell 命令对detached/目录chown为 ClickHouse 服务端运行用户——源码注释说明得清楚只有文件属主与 ClickHouse 服务用户一致时ATTACH命令才能生效L107-L123。阶段四ATTACH 挂载数据 part写入器提交后返回CKFileCommitInfo记录每个 shard 上的 detached 文件清单。真正的“提交”发生在聚合提交器 ClickhouseFileSinkAggCommitter 中combine()先按 shard 合并多个 writer 的文件列表commit()再对每个文件执行L125-L140ALTER TABLE local_table ATTACH PART part文件名;执行成功即数据正式进入表中。该连接器实现了SinkAggregatedCommitter接口见 ClickhouseFileSink 的createAggregatedCommitter()因此挂载动作是全局一次性的而不是每个 writer 各自执行。此外Writer 初始化时会做节点密码预检nodePasswordCheck()L153-L175在node_free_password false时要求node_pass必须覆盖每一个 shard 节点按 hostname 或 host 匹配否则直接抛出PASSWORD_NOT_FOUND_IN_SHARD_NODE异常——这解释了为什么多节点集群必须配全所有节点的密码以及免密登录时为什么要设置node_free_password true。配置项详解官方文档给出的完整选项表如下默认值已与 ClickhouseFileSinkOptions、ClickhouseBaseOptions 源码核对一致NameTypeRequiredDefaulthoststringyes-databasestringyes-tablestringyes-usernamestringyes-passwordstringyes-clickhouse_local_pathstringyes-sharding_keystringno-copy_methodstringnoscpnode_free_passwordbooleannofalsenode_passlistno-node_pass.node_addressstringno-node_pass.usernamestringnorootnode_pass.passwordstringno-compatible_modebooleannofalsefile_fields_delimiterstringno\tfile_temp_pathstringno/tmp/seatunnel/clickhouse-local/filekey_pathstringno/tmp/id_rsacommon-options-no-各选项的含义与实现要点hostClickHouse 集群地址格式为host:port支持逗号分隔多节点如host1:8123,host2:8123。必填项在 ClickhouseFileSinkFactory.optionRule() 中声明为HOST, TABLE, DATABASE, USERNAME, PASSWORD, CLICKHOUSE_LOCAL_PATH六项与文档的 Required 列完全对应。database / table / username / password目标库表名与 ClickHouse 账号。注意 username/password 是 ClickHouse 服务账号与后面传输文件用的 Linux 账号是两回事。clickhouse_local_pathclickhouse-local程序在 SeaTunnel 计算节点上的可执行路径。由于每个并行任务都会直接调用该程序文档特别强调它必须位于每个计算节点的相同路径。源码中该值可包含空格分隔的多个路径generateClickhouseLocalFiles按空格拆分以适配不同的可执行文件命名。sharding_key数据拆分写入时用于选择目标节点的字段。不配置时随机路由配置后按该字段值做分片路由对应ShardRouter。该字段会在工厂中被解析出类型并放入ShardMetadata见 ClickhouseFileSinkFactory。copy_method文件传输方式可选scp默认与rsync。ClickhouseFileCopyMethod 为枚举类型from()方法忽略大小写匹配传入其他值会抛出Unknown ClickhouseFileCopyMethod异常。node_free_passwordSeaTunnel 需要访问 ClickHouse 服务器端文件系统来放置数据文件。若计算节点到各 ClickHouse 节点已配置免密登录SSH key可置为true并省略node_pass否则必须在node_pass中配置每个节点的登录密码。node_pass列表保存所有 ClickHouse 服务器节点的地址与登录凭据。子项node_address为节点地址username为该节点的 Linux 用户名默认root——源码中工厂会把未显式给出username的条目统一补为rootClickhouseFileSinkFactory非 root 场景必须显式指定changelog 中“support non-root username for fileTransfer”即此需求password为对应密码。配置对象对应 NodePassConfig其中node_address通过JsonProperty显式映射 snake_case 配置键。compatible_mode低版本 ClickHouse 的 clickhouse-local 不支持--path参数时置为true改用生成临时config.xml--config-file的方式实现等价效果实现见前文阶段二。file_fields_delimiter临时 CSV 文件的字段分隔符默认制表符\t。若数据值本身包含分隔符字符会导致解析异常此时应更换为数据中不出现的字符。源码中强制校验该值必须恰好为 1 个字符否则抛CONFIG_VALIDATION_FAILEDClickhouseFileSinkFactory。file_temp_path临时文件根目录默认/tmp/seatunnel/clickhouse-local/file。每次提交在其下创建uuid/子目录存放local_data.log与生成的 part 文件传输完成后整个目录会被清理clearLocalFileDirectory注意磁盘空间需能容纳单个并行任务的批量数据。key_pathscp/rsync 使用的 SSH 私钥文件路径默认/tmp/id_rsa。ScpFileTransfer 在init()中将其加载为 RSA 公钥身份未配置时仅走密码认证。common-optionsSink 插件通用参数参见官方文档中的 Sink Common Options 说明位于 docs/en/connectors/common-options/ 目录。另外从 ClickhouseBaseOptions 与工厂的 optional 列表可以确认该连接器还支持可选的server_time_zoneClickHouse 会话时区不设置时使用系统默认时区虽然文档选项表未单列实际配置中可以使用。完整配置示例官方文档给出的 HOCON 示例已保留原样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 }] }结合前文的机制说明落地部署前建议按以下清单自查目标表是Distributed引擎且internal_replication true每个计算节点SeaTunnel 实际运行任务的节点都存在clickhouse_local_path指定的可执行文件且该节点可执行 clickhouse-local计算节点可通过 SSH22 端口登录各 ClickHouse 节点要么免密配node_free_password true要么node_pass覆盖了所有节点nodePasswordCheck预检node_pass.username指定的 Linux 用户对 ClickHouse 数据目录有写权限连接器随后会 chown 给 ClickHouse 服务用户临时目录file_temp_path所在磁盘有足够空间且分隔符file_fields_delimiter不出现在数据中低版本 clickhouse-local 记得开启compatible_mode。与 ClickhouseJDBCSink 的关系同一连接器模块中还实现了基于 JDBC 批量的ClickhouseSinkClickhouseSink二者选择可以简单归纳为常规在线写入、要求语义更强时用 JDBC Sink超大吞吐的批量装载、且目标是Distributed分布式表时用 ClickhouseFile。官方文档也在 Key features 中提示“Write data to Clickhouse can also be done using JDBC”两者可视为互补而非互斥。相关源码与文档索引内容路径官方文档docs/en/connectors/sink/ClickhouseFile.md连接器 Changelogdocs/en/connectors/changelog/connector-clickhouse.md工厂与选项校验ClickhouseFileSinkFactory选项定义ClickhouseFileSinkOptions / ClickhouseBaseOptions写入与文件生成ClickhouseFileSinkWriterATTACH 提交ClickhouseFileSinkAggCommitter文件传输实现ScpFileTransfer / RsyncFileTransfer / FileTransferFactory从 连接器 Changelog 可以看到ClickhouseFile 在 2.3.9 左右经历了一轮集中改进统一文件生成路径、rsync 指定用户、attach SQL 日志、直接连接各 shard 节点获取数据路径、支持公钥身份认证等这些能力都已体现在当前源码中使用前建议对照 changelog 确认自己的版本包含所需特性。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考