news 2026/9/21 15:47:09

DataX ClickhouseReader 插件深度解析:基于 JDBC 的 ClickHouse 数据抽取原理、配置参数与实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DataX ClickhouseReader 插件深度解析:基于 JDBC 的 ClickHouse 数据抽取原理、配置参数与实战指南

DataX ClickhouseReader 插件深度解析:基于 JDBC 的 ClickHouse 数据抽取原理、配置参数与实战指南

【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX

本篇技术指南以阿里云 DataWorks 数据集成的开源版本 DataX 中的ClickhouseReader 插件为核心,系统讲解其如何通过 JDBC 从 ClickHouse 读取数据并同步到任意 DataX 支持的写入端。文中完整覆盖插件的实现原理、两种作业配置样例(常规表模式与自定义 SQL 模式)、全部核心参数说明、数据类型转换机制与并发切分原理,并结合当前仓库源码给出可验证的调用链与底层实现证据,帮助读者在业务中正确配置、调优并规避一致性、编码、增量同步等典型陷阱。


1 插件定位与快速介绍

ClickhouseReader 是 DataX 生态中的数据库读取插件,负责从 ClickHouse 数据库读取数据。与 DataX 其他 RDBMS 系 reader 插件一致,ClickhouseReader 本身不直接持有数据,而是通过 JDBC 连接远程 ClickHouse 数据库,执行用户配置的 SQL 语句,将 SELECT 结果转换为 DataX 统一的数据抽象(Record),再交由下游 Writer 插件(如 streamwriter、mysqlwriter、hdfswriter 等)落盘。

在代码层面,ClickhouseReader 是一个典型的“薄插件”实现:插件入口 ClickhouseReader.java 几乎不做业务逻辑,而是将Job(任务级:初始化、切分、清理)与Task(通道级:真实读取)全部委托给 RDBMS 体系共用的 CommonRdbmsReader.java,仅通过DataBaseType.ClickHouse标记数据库类型:

private static final DataBaseType DATABASE_TYPE = DataBaseType.ClickHouse; public static class Job extends Reader.Job { public void init() { this.jobConfig = super.getPluginJobConf(); this.commonRdbmsReaderMaster = new CommonRdbmsReader.Job(DATABASE_TYPE); this.commonRdbmsReaderMaster.init(this.jobConfig); } ... }

这种设计意味着 ClickhouseReader 天然继承了 RDBMS 系列 reader 的通用能力(多 JDBC 地址探测、表/自定义 SQL 双模式、splitPk 切分、脏数据收集、类型映射等),这也是理解本文后续所有原理的钥匙。

2 实现原理:从 JDBC 到 DataX Record 的完整链路

简而言之,ClickhouseReader 的工作过程可以概括为以下几步:

  1. 建立连接:通过 JDBC 连接器(驱动类为ru.yandex.clickhouse.ClickHouseDriver,见 DataBaseType.java)连接远程 ClickHouse 数据库;
  2. 生成 SQL
    • 若用户配置了tablecolumnwhere,插件将它们拼接为标准的SELECT语句(模板见 Constant.java:select %s from %s where (%s))发送给 ClickHouse;
    • 若用户配置了querySql,插件原样透传该 SQL 到 ClickHouse,不做任何拼接或改写;
  3. 结果集转换:将 JDBCResultSet逐行读取,按列类型映射为 DataX 自定义数据类型(LongColumnDoubleColumnStringColumnDateColumnBoolColumnBytesColumn),拼装为统一的抽象数据集Record
  4. 数据传递:通过RecordSender.sendToWriter(record)Record交给下游 Writer 处理。

上述第 3、4 步的核心实现在 CommonRdbmsReader.java 的Task.startRead中:先由DBUtil.getConnection建连,随后DBUtil.query(conn, querySql, fetchSize)以指定fetchSize执行查询,再在while (rs.next())循环中通过transportOneRecord逐行构造Record并发送。值得注意的是,读取过程中每条记录的构造与类型转换(buildRecord)会捕获异常并调用taskPluginCollector.collectDirtyRecord(record, e)将问题记录标记为脏数据,而不是直接中断整个任务。

关于 JDBC 依赖,clickhousereader/pom.xml 中声明了ru.yandex.clickhouse:clickhouse-jdbc:0.2.4,同时依赖datax-coredatax-commonplugin-rdbms-util,这印证了插件对通用 RDBMS 读取框架的复用关系。

3 功能说明与配置实战

3.1 配置样例一:常规表模式(column + table + where 自动拼 SQL)

以下作业从 Clickhouse 数据库同步抽取数据到本地(writer 使用 streamwriter 打印内容),是 ClickhouseReader 最典型的用法:

{ "job": { "setting": { "speed": { //设置传输速度 byte/s 尽量逼近这个速度但是不高于它. // channel 表示通道数量,byte表示通道速度,如果单通道速度1MB,配置byte为1048576表示一个channel "byte": 1048576 }, //出错限制 "errorLimit": { //先选择record "record": 0, //百分比 1表示100% "percentage": 0.02 } }, "content": [ { "reader": { "name": "clickhousereader", "parameter": { // 数据库连接用户名 "username": "root", // 数据库连接密码 "password": "root", "column": [ "id","name" ], "connection": [ { "table": [ "table" ], "jdbcUrl": [ "jdbc:clickhouse://[HOST_NAME]:PORT/[DATABASE_NAME]" ] } ] } }, "writer": { //writer类型 "name": "streamwriter", // 是否打印内容 "parameter": { "print": true } } } ] } }

说明:speed.byte为字节速率上限(如 1048576 表示单通道 1MB/s),errorLimit用于控制任务可容忍的错误记录数(recordpercentage取其一即生效)。

3.2 配置样例二:自定义 SQL 模式(querySql 直接透传)

在某些复杂场景(如多表 JOIN、子查询、聚合)下,where不足以描述筛选条件,可改用querySql完全自定义抽取语句:

{ "job": { "setting": { "speed": { "channel": 5 } }, "content": [ { "reader": { "name": "clickhousereader", "parameter": { "username": "root", "password": "root", "where": "", "connection": [ { "querySql": [ "select db_id,on_line_flag from db_info where db_id < 10" ], "jdbcUrl": [ "jdbc:clickhouse://1.1.1.1:8123/default" ] } ] } }, "writer": { "name": "streamwriter", "parameter": { "visible": false, "encoding": "UTF-8" } } } ] } }

两种模式在源码中被严格区分:OriginalConfPretreatmentUtil.recognizeTableOrQuerySqlMode(见 OriginalConfPretreatmentUtil.java)会逐个 connection 检查:

  • tablequerySql同时配置或同时缺失,直接抛出TABLE_QUERYSQL_MIXED/TABLE_QUERYSQL_MISSING错误;
  • 若多个 connection 混用两种模式,同样报错;
  • 当用户配置 querySql 时,ClickhouseReader 直接忽略 table、column、where、splitPk 的配置(源码中会打印 warn 日志并移除这些项)。

3.3 参数详解

jdbcUrl
  • 描述:对端 ClickHouse 数据库的 JDBC 连接信息,使用 JSON 数组描述,支持一个库填写多个连接地址。之所以使用数组,是因为阿里集团内部支持多 IP 探测:配置多个地址时,ClickhouseReader 会依次探测 IP 的可连接性,直到选择一个合法的 IP;若全部连接失败则报错。对外部使用者,填写一个 JDBC 连接即可。该参数必须包含在 connection 配置单元内
  • jdbcUrl 需遵循 ClickHouse 官方 JDBC 规范,可附带连接控制参数(如jdbc:clickhouse://host:8123/default)。
  • 必选:是;默认值:无。

源码侧,多地址探测在DBUtil.chooseJdbcUrl中完成(由 OriginalConfPretreatmentUtil.java 的dealJdbcAndTable调用),选定后会回写回connection[i].jdbcUrl。另外可注意到DataBaseType.ClickHouseappendJDBCSuffixForReader不做额外后缀拼接(ClickHouse 分支为空),因此连接串完全由用户控制。

username
  • 描述:数据源的用户名。
  • 必选:是;默认值:无。
password
  • 描述:数据源指定用户名的密码。
  • 必选:是;默认值:无。

usernamepassword为必填项这一点在源码中有硬性校验:OriginalConfPretreatmentUtil.doPretreatment中调用originalConfig.getNecessaryValue(Key.USERNAME, ...)getNecessaryValue(Key.PASSWORD, ...),缺失即抛REQUIRED_VALUE错误。

table
  • 描述:所选取的需要同步的表,使用 JSON 数组描述,支持多张表同时抽取。当配置多张表时,用户需自己保证多张表是同一 schema 结构,ClickhouseReader不检查各表是否为同一逻辑表。该参数必须包含在 connection 配置单元内
  • 必选:是;默认值:无。
column
  • 描述:所配置的表中需要同步的列名集合,使用 JSON 数组描述字段信息。支持:
    • 全部列:使用*代表默认使用所有列,如["*"]
    • 列裁剪:可以挑选部分列进行导出;
    • 列换序:可以不按照表 schema 信息进行导出;
    • 常量与表达式:按 JSON 格式填写,例如["id", "\table`", "1", "'bazhen.csy'", "null", "to_char(a + 1)", "2.3" , "true"]其中id为普通列名,``table`` 为包含保留字的列名,1为整型数字常量,'bazhen.csy'为字符串常量,null为空指针,to_char(a + 1)为表达式,2.3为浮点数,true` 为布尔值。
  • Column 必须显式填写,不允许为空
  • 必选:是;默认值:无。

源码侧,表模式下column缺失或为空会直接报REQUIRED_VALUE;若配置为单个"*"则回填为*(全列);若混用*与其他列名则报ILLEGAL_VALUE,这些逻辑均位于 OriginalConfPretreatmentUtil.java 的dealColumnConf中。

splitPk
  • 描述:数据抽取时的切分主键。指定 splitPk 后,DataX 会启动并发任务进行数据同步,显著提升同步效能。推荐使用表主键作为 splitPk,因为主键分布通常较均匀,切分出的分片不容易出现数据热点。
  • 限制:目前 splitPk 仅支持整型数据切分不支持浮点、日期等其他类型。若指定了非支持类型,ClickhouseReader 将报错。
  • splitPk 不填写时,视为不对单表切分,ClickhouseReader 使用单通道同步全量数据
  • 必选:否;默认值:无。
where
  • 描述:筛选条件。ClickhouseReader 根据指定的 column、table、where 条件拼接 SQL 再进行抽取。实际业务中常用于同步当天数据,如where: "gmt_create > $bizdate"
  • 注意:不可将 where 条件指定为limit 10,因为 limit 不是 SQL 合法的 where 子句。
  • where 条件可以有效地进行业务增量同步。
  • 必选:否;默认值:无。

源码中dealWhere会对 where 做预处理:去掉首尾空白,并剔除结尾的分号;或全角),避免拼接 SQL 时产生语法问题。

querySql
  • 描述:在某些业务场景下,where 不足以描述筛选条件,可通过该配置项自定义筛选 SQL。配置此项后,DataX 系统忽略 table、column 这些配置项,直接使用该配置项的内容进行筛选。例如多表 join 后同步数据:select a,b from table_a join table_b on table_a.id = table_b.id
  • 当用户配置 querySql 时,ClickhouseReader 直接忽略 table、column、where 条件的配置
  • 必选:否;默认值:无。
fetchSize
  • 描述:定义插件与数据库服务端每次批量获取数据的条数,该值决定了 DataX 与服务端的网络交互次数,能够较大地提升数据抽取性能。
  • 注意:该值过大(>2048)可能造成 DataX 进程 OOM
  • 必选:否。
  • 默认值:文档标注为1024;从当前仓库源码看,ClickhouseReader.java 的Task.startRead中读取配置时使用的兜底值为1000getInt(Constant.FETCH_SIZE, 1000)),实际生效值以作业配置为准,建议在文档建议的 1024 上下根据内存与网络情况调整。
session
  • 描述:控制数据读取时的时间格式、时区等会话配置。如果表中有时间字段,可配置该值以明确告知 Clickhouse 读取的时间格式(如 NLS_DATE_FORMAT、NLS_TIME_FORMAT 等)。其配置值为 JSON 数组格式,例如:
"session": [ "alter session set NLS_DATE_FORMAT='yyyy-mm-dd hh24:mi:ss'", "alter session set NLS_TIMESTAMP_FORMAT='yyyy-mm-dd hh24:mi:ss'", "alter session set NLS_TIMESTAMP_TZ_FORMAT='yyyy-mm-dd hh24:mi:ss'", "alter session set TIME_ZONE='US/Pacific'" ]
  • 注意:&quot;"的转义字符串,实际书写 JSON 时按标准 JSON 语法即可。
  • 必选:否;默认值:无。

从源码看,session 配置由 CommonRdbmsReader.java 的Task.startRead在建立连接后调用DBUtil.dealWithSessionConfig(conn, readerSliceConfig, dataBaseType, basicMsg)执行,属于从 RDBMS 通用读取框架继承的能力。

4 类型转换机制

ClickhouseReader 支持大部分 ClickHouse 类型,但仍存在部分类型未支持的情况,使用时请注意检查字段类型。官方文档给出的 ClickHouse 类型 → DataX 内部类型转换列表如下:

DataX 内部类型Clickhouse 数据类型
LongUInt8, UInt16, UInt32, UInt64, UInt128, UInt256, Int8, Int16, Int32, Int64, Int128, Int256
DoubleFloat32, Float64, Decimal
StringString, FixedString
DateDATE, Date32, DateTime, DateTime64
BooleanBoolean
BytesBLOB, BFILE, RAW, LONG RAW

请注意:除上述罗列字段类型外,其他类型均不支持。

该映射在源码层面与 JDBC 标准类型一一对应(见 CommonRdbmsReader.java 的buildRecord,以及同逻辑的 ResultSetReadProxy.java):

  • CHAR / NCHAR / VARCHAR / LONGVARCHAR / NVARCHAR / LONGNVARCHAR / CLOB / NCLOBStringColumn
  • SMALLINT / TINYINT / INTEGER / BIGINTLongColumn(通过rs.getString取值后构造);
  • NUMERIC / DECIMAL / FLOAT / REAL / DOUBLEDoubleColumn
  • TIME / DATE / TIMESTAMPDateColumn(其中 DATE 类型若列类型名为year则映射为LongColumn);
  • BINARY / VARBINARY / BLOB / LONGVARBINARYBytesColumn
  • BOOLEAN / BITBoolColumn
  • NULLStringColumn
  • 其他类型:抛出UNSUPPORTED_TYPE错误,提示“请尝试使用数据库函数将其转换 DataX 支持的类型,或者不同步该字段”。

因此,对于文档类型转换表中未覆盖的 ClickHouse 类型,一个可行的规避方式是:在columnquerySql中使用 ClickHouse 函数将目标列显式转换为受支持的类型后再同步。

5 数据切分(splitPk)与并发机制

ClickhouseReader 的并发能力来源于 RDBMS 框架的通用切分逻辑,调用链为:

ClickhouseReader.Job.split(mandatoryNumber)CommonRdbmsReader.Job.splitReaderSplitUtil.doSplit(originalConfig, adviceNumber)

在 ReaderSplitUtil.java 中:

  • 表模式且配置了 splitPk:当每个表应切分的份数> 1时,进入切分分支,调用SingleTableSplitUtil.splitSingleTable对每个表做范围切分;
  • 表模式未配置 splitPk:直接为每张表生成一条完整查询 SQL(select %s from %s where (%s)模板);
  • querySql 模式:将配置的每条 querySql 直接作为一个切分片,不做任何改写。

SingleTableSplitUtil.splitSingleTable(见 SingleTableSplitUtil.java)是切分的核心实现,其原理可以概括为:

  1. 先执行SELECT MIN(splitPk), MAX(splitPk) FROM table [WHERE ... AND splitPk IS NOT NULL]取得切分主键的取值范围(getPkRange);
  2. 校验切分主键类型:仅接受整型BIGINT/INTEGER/SMALLINT/TINYINT)或字符串CHAR/VARCHAR/NVARCHAR等)类型,否则抛ILLEGAL_SPLIT_PK错误——这与文档中“splitPk 仅支持整形数据切分,不支持浮点、日期等其他类型”的约束一致;
  3. 依据范围将数据等分为若干区间,为每个区间生成带where条件的查询 SQL;
  4. 额外生成一个splitPk IS NULL的兜底分片,保证主键为空的数据也不会丢失;
  5. 最终将每个分片作为独立的Configuration(Task 配置)返回,由 DataX 调度为并发通道执行。

另外,从ReaderSplitUtil中可以看到,单表切分份数会受splitFactor(默认 5)放大(eachTableShouldSplittedNumber * splitFactor),这是为了避免导入下游(如 Hive)时产生过多小文件而做的工程权衡。理解这一点有助于在实际生产中预估并发任务数。

6 性能测试方法与报告框架

原文档中保留了性能测试的章节框架(clickhousereader/doc/clickhousereader.md 第 4 节),包含:

  • 环境准备:为模拟线上真实数据设计两个 ClickHouse 数据表;分别给出执行 DataX 的机器参数与 ClickHouse 数据库机器参数;
  • 测试报告:以表 1 为例,按并发任务数记录DataX速度(Rec/s)DataX流量网卡流量DataX运行负载DB运行负载等指标。

需要说明的是:当前仓库的该文档中,上述各子章节的具体数值(数据表 DDL、机器规格、实测吞吐量)为留空占位状态,并未提供可引用的实测数据。因此本文不臆造任何性能数字。若读者需要自行评估性能,建议参照上述框架:固定数据表与机器环境、按并发任务数梯度(如 1、2、4、8)记录 DataX 吞吐(Rec/s)与流量(Byte/s)、观察网卡流量与两端负载,重点验证fetchSizespeed.channel对吞吐的影响,同时监控 DataX 进程内存以避免fetchSize过大导致的 OOM。

7 约束限制

7.1 主备同步数据恢复问题

Clickhouse 使用主从灾备时,备库会从主库不间断地通过 binlog 恢复数据。由于主备数据同步存在时间差(尤其在某些特定情况如网络延迟下),备库同步恢复的数据可能与主库有较大差别,导致从备库同步的数据不是一份当前时间的完整镜像

针对这一问题,文档提及计划提供preSql功能以在读取前执行前置 SQL 保障数据可用性,但该功能当前待补充实现(文档原文标注“该功能待补充”)。读者在从备库抽取数据时应自行评估数据新鲜度与一致性要求。

7.2 一致性约束

Clickhouse 在数据存储划分中属于 RDBMS 系统,对外可提供强一致性数据查询接口。例如:当一次同步任务启动运行过程中,若该库存在其他数据写入方写入数据,ClickhouseReader完全不会获取到写入更新数据——这是由数据库本身的快照特性(MVCC,多版本并发控制)决定的。

但上述特性仅在ClickhouseReader 单线程模型下成立。当用户配置splitPk启用并发抽取后,ClickhouseReader 会先后启动多个并发任务:由于多个并发任务相互之间不属于同一个读事务,且并发任务之间存在时间间隔,因此这份数据并不是完整的、一致的数据快照

针对多线程下的一致性快照需求,技术上目前无法直接实现,只能从工程角度解决,文档给出两种可选思路(各有取舍,用户自行权衡):

  1. 使用单线程同步,即不再进行数据切分。缺点是速度较慢,但能够很好地保证一致性;
  2. 关闭其他数据写入方,保证当前数据为静态数据,例如锁表、关闭备库同步等。缺点是可能影响在线业务。

7.3 数据库编码问题

ClickhouseReader 底层使用 JDBC 进行数据抽取,JDBC 天然适配各类编码并在底层完成编码转换,因此ClickhouseReader 不需要用户指定编码,可以自动获取编码并转码。

但对于 ClickHouse 底层写入编码与其设定编码不一致的混乱情况,ClickhouseReader 无法识别,也无法提供解决方案——这类情况下导出有可能为乱码,需要用户从数据写入源头保证编码一致。

7.4 增量数据同步

ClickhouseReader 使用 JDBC SELECT 语句完成数据抽取,因此可以借助SELECT ... WHERE ...实现增量抽取,常见方式有两种:

  • 基于更新时间戳:数据库在线应用写入数据时填充 modify 字段为更改时间戳(覆盖新增、更新、删除逻辑删)。这类应用只需在where条件中带上一次同步阶段的时间戳即可;
  • 基于自增 ID:对于新增流水型数据,可在where条件中带上一次同步阶段的最大自增 ID作为下界。

对于业务上没有字段区分新增/修改数据的情况,ClickhouseReader 无法进行增量数据同步,只能同步全量数据。

7.5 SQL 安全性

ClickhouseReader 提供querySql让用户自行实现 SELECT 抽取语句,但本身对 querySql 不做任何安全性校验(如防注入、防危险语句等),这块需要由 DataX 使用方自行保证。生产环境中应通过白名单、SQL 审计等手段管控作业配置的来源。

8 FAQ(常见问题)

Q1:ClickhouseReader 同步报错,报错信息为 XXX,如何处理?

A:多为网络或权限问题,请先使用 Clickhouse 命令行(clickhouse-client)测试同样的连接与查询。如果命令行同样报错,即可证实是环境问题,请联系你的 DBA。

Q2:ClickhouseReader 抽取速度很慢怎么办?

A:影响抽取时间的原因大致有以下几个(来自专业 DBA 卫绾的经验总结):

  1. SQL 执行计划异常导致抽取时间长:抽取时尽可能使用全表扫描代替索引扫描;
  2. 合理设置 SQL 的并发度:通过speed.channelsplitPk合理增大并发,减少抽取时间;
  3. 抽取 SQL 要简单:尽量不用replace等函数,这类函数非常消耗 CPU,会严重影响抽取速度。

9 总结

ClickhouseReader 是 DataX 中一个“实现轻量但能力完备”的 RDBMS 系读取插件:它以 JDBC 为唯一通道,通过复用plugin-rdbms-util的通用框架(配置校验、多地址探测、SQL 拼接、splitPk 切分、类型映射、脏数据处理)实现了对 ClickHouse 的高效读取,并支持常规表模式与自定义 SQL 模式两种用法。

在实际使用中,建议读者重点关注以下四点:其一,column必须显式配置,且只支持文档类型转换表中列出的类型,其余类型需用 ClickHouse 函数转换;其二,splitPk仅支持整型切分,启用后注意多并发任务之间不具备一致快照;其三,fetchSize是吞吐与内存的平衡点,切忌过大(>2048 有 OOM 风险);其四,增量同步依赖业务侧提供时间戳或自增 ID 字段,无此类字段时只能全量同步。结合本文给出的 ClickhouseReader.java、CommonRdbmsReader.java、OriginalConfPretreatmentUtil.java、ReaderSplitUtil.java 与 SingleTableSplitUtil.java 等源码路径,读者可以进一步深入验证插件的每一处行为细节。

【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/21 15:25:21

Python面向对象编程核心技术与工程实践

1. 为什么需要面向对象编程&#xff1f;十五年前我刚接触Python时&#xff0c;所有代码都是线性脚本。直到接手一个电商库存管理系统&#xff0c;3000行代码挤在同一个文件里&#xff0c;修改价格计算逻辑需要排查几十个函数——那天起我真正理解了OOP的价值。面向对象编程&…

作者头像 李华