SeaTunnel ClickhouseFile Sink 连接器详解:基于 clickhouse-local 的海量数据批量文件导入

发布时间:2026/10/9 2:35:06
SeaTunnel ClickhouseFile Sink 连接器详解:基于 clickhouse-local 的海量数据批量文件导入 数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文以 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它同时实现了SeaTunnelSinkClickhouseFileSink.java、SinkWriterClickhouseFileSinkWriter.java与SinkAggregatedCommitterClickhouseFileSinkAggCommitter.java分别承担“本地写文件”、“聚合提交ATTACH”使文件迁移与 ATTACH 具备事务提交语义。小提示如果数据量不需要走文件批量装载也可以使用 JDBC 方式将数据写入 ClickHouse对应同一连接器模块中的ClickhouseSink配置bulk_size、sql等参数见 ClickhouseConfig.java。二、特性与适用前提特性说明精准一次exactly-once暂不支持。相关特性能力矩阵可参考 connector-v2-features。硬性前置条件不满足则无法正常使用表引擎必须是DistributedClickhouseFile 依赖 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下创建一个以 UUID10 位命名的临时目录并打开local_data.log文件CLICKHOUSE_LOCAL_FILE_SUFFIX /local_data.log写入行数据按表字段顺序、以file_fields_delimiter拼接成一行文本空字段写空串通过FileChannel.map内存映射缓冲区每次映射 128KB高效追加写入saveDataToFile。3.3 prepareCommitclickhouse-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_tableuuid \ -q DDL; INSERT INTO TABLE local_table SELECT ... FROM temp_tableuuid; \ --path file_temp_path/uuid # 或 --config-filecompatible_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 commitALTER TABLE ATTACH PART所有并行 Writer 的提交信息CKFileCommitInfo即“分片 - part 文件列表”映射汇总到聚合提交器后combine()合并、commit()执行真正的入库ClickhouseFileSinkAggCommitter.javaALTER TABLE local_table_name ATTACH PART part名称即对每个分片、每个 part 文件执行 ATTACH将 detached 目录中的 part 附加到本地表。至此数据正式进入 ClickHouse后续由 ClickHouse 的分片/副本机制与 Distributed 表对外提供查询。四、参数详解以下参数表与官方文档一致并补充了源码中的校验逻辑、默认值与取值约束源码依据ClickhouseConfig.java必填/可选规则见 ClickhouseFileSinkFactory.optionRule名称类型是否必须默认值hoststring是-databasestring是-tablestring是-usernamestring是-passwordstring是-clickhouse_local_pathstring是-sharding_keystring否-copy_methodstring否scpnode_free_passwordboolean否falsenode_passlist否-node_pass.node_addressstring否-node_pass.usernamestring否rootnode_pass.passwordstring否-compatible_modeboolean否falsefile_fields_delimiterstring否\tfile_temp_pathstring否/tmp/seatunnel/clickhouse-local/filecommon-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。源码中解析该字段时若缺失则回退为rootClickhouseFileSink.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-optionsSink 插件常用参数如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 版本一致性低版本 ClickHouseclickhouse-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兼容低版本 ClickHousePR #3416。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐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大数据批处理流处理变更数据捕获上一篇hyperframes 的 frame.md 设计规范从解析优先级到帧级品牌系统的完整指南下一篇Phase {N} — UI Review创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

尧图内容编辑团队 内容团队

尧图内容编辑团队

本文由尧图网络内容编辑团队执笔。团队由资深项目经理、前端工程师与设计师组成,所有内容均来自亲手交付的真实项目,先讲清问题、再给出可落地的解法。尧图深耕北京网站建设十年,服务过京华建材集团、智造科技等各行业客户,把一线经验沉淀为可复用的行业观察。

  • 十年建站经验,覆盖建材、制造、服务、文创等
  • 项目经理把关选题与事实准确性
  • 工程师与设计师联合撰写专业细节
  • 统一编辑规范,保证文风与排版一致
  • 每月复盘转化数据,迭代选题方向

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

建站决策前值得细读的三篇

网站改版的5个关键决策
2024-08-12

网站改版的5个关键决策

什么时候该改版、改到什么程度、如何避免流量掉光,京华建材集团改版复盘给出答案。

获取专属建站方案

看完文章,把您的行业与预算告诉我们,免费获取一份量身定制的官网建设方案与报价。

立即免费咨询