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 的工作过程可以概括为以下几步:
- 建立连接:通过 JDBC 连接器(驱动类为
ru.yandex.clickhouse.ClickHouseDriver,见 DataBaseType.java)连接远程 ClickHouse 数据库; - 生成 SQL:
- 若用户配置了
table、column、where,插件将它们拼接为标准的SELECT语句(模板见 Constant.java:select %s from %s where (%s))发送给 ClickHouse; - 若用户配置了
querySql,插件原样透传该 SQL 到 ClickHouse,不做任何拼接或改写;
- 若用户配置了
- 结果集转换:将 JDBC
ResultSet逐行读取,按列类型映射为 DataX 自定义数据类型(LongColumn、DoubleColumn、StringColumn、DateColumn、BoolColumn、BytesColumn),拼装为统一的抽象数据集Record; - 数据传递:通过
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-core、datax-common与plugin-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用于控制任务可容忍的错误记录数(record与percentage取其一即生效)。
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 检查:
- 若
table与querySql同时配置或同时缺失,直接抛出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.ClickHouse的appendJDBCSuffixForReader不做额外后缀拼接(ClickHouse 分支为空),因此连接串完全由用户控制。
username
- 描述:数据源的用户名。
- 必选:是;默认值:无。
password
- 描述:数据源指定用户名的密码。
- 必选:是;默认值:无。
username、password为必填项这一点在源码中有硬性校验: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中读取配置时使用的兜底值为1000(getInt(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'" ]- 注意:
"是"的转义字符串,实际书写 JSON 时按标准 JSON 语法即可。 - 必选:否;默认值:无。
从源码看,session 配置由 CommonRdbmsReader.java 的Task.startRead在建立连接后调用DBUtil.dealWithSessionConfig(conn, readerSliceConfig, dataBaseType, basicMsg)执行,属于从 RDBMS 通用读取框架继承的能力。
4 类型转换机制
ClickhouseReader 支持大部分 ClickHouse 类型,但仍存在部分类型未支持的情况,使用时请注意检查字段类型。官方文档给出的 ClickHouse 类型 → DataX 内部类型转换列表如下:
| DataX 内部类型 | Clickhouse 数据类型 |
|---|---|
| Long | UInt8, UInt16, UInt32, UInt64, UInt128, UInt256, Int8, Int16, Int32, Int64, Int128, Int256 |
| Double | Float32, Float64, Decimal |
| String | String, FixedString |
| Date | DATE, Date32, DateTime, DateTime64 |
| Boolean | Boolean |
| Bytes | BLOB, BFILE, RAW, LONG RAW |
请注意:除上述罗列字段类型外,其他类型均不支持。
该映射在源码层面与 JDBC 标准类型一一对应(见 CommonRdbmsReader.java 的buildRecord,以及同逻辑的 ResultSetReadProxy.java):
CHAR / NCHAR / VARCHAR / LONGVARCHAR / NVARCHAR / LONGNVARCHAR / CLOB / NCLOB→StringColumn;SMALLINT / TINYINT / INTEGER / BIGINT→LongColumn(通过rs.getString取值后构造);NUMERIC / DECIMAL / FLOAT / REAL / DOUBLE→DoubleColumn;TIME / DATE / TIMESTAMP→DateColumn(其中 DATE 类型若列类型名为year则映射为LongColumn);BINARY / VARBINARY / BLOB / LONGVARBINARY→BytesColumn;BOOLEAN / BIT→BoolColumn;NULL→StringColumn;- 其他类型:抛出
UNSUPPORTED_TYPE错误,提示“请尝试使用数据库函数将其转换 DataX 支持的类型,或者不同步该字段”。
因此,对于文档类型转换表中未覆盖的 ClickHouse 类型,一个可行的规避方式是:在column或querySql中使用 ClickHouse 函数将目标列显式转换为受支持的类型后再同步。
5 数据切分(splitPk)与并发机制
ClickhouseReader 的并发能力来源于 RDBMS 框架的通用切分逻辑,调用链为:
ClickhouseReader.Job.split(mandatoryNumber)→CommonRdbmsReader.Job.split→ReaderSplitUtil.doSplit(originalConfig, adviceNumber)
在 ReaderSplitUtil.java 中:
- 表模式且配置了 splitPk:当每个表应切分的份数
> 1时,进入切分分支,调用SingleTableSplitUtil.splitSingleTable对每个表做范围切分; - 表模式未配置 splitPk:直接为每张表生成一条完整查询 SQL(
select %s from %s where (%s)模板); - querySql 模式:将配置的每条 querySql 直接作为一个切分片,不做任何改写。
SingleTableSplitUtil.splitSingleTable(见 SingleTableSplitUtil.java)是切分的核心实现,其原理可以概括为:
- 先执行
SELECT MIN(splitPk), MAX(splitPk) FROM table [WHERE ... AND splitPk IS NOT NULL]取得切分主键的取值范围(getPkRange); - 校验切分主键类型:仅接受整型(
BIGINT/INTEGER/SMALLINT/TINYINT)或字符串(CHAR/VARCHAR/NVARCHAR等)类型,否则抛ILLEGAL_SPLIT_PK错误——这与文档中“splitPk 仅支持整形数据切分,不支持浮点、日期等其他类型”的约束一致; - 依据范围将数据等分为若干区间,为每个区间生成带
where条件的查询 SQL; - 额外生成一个
splitPk IS NULL的兜底分片,保证主键为空的数据也不会丢失; - 最终将每个分片作为独立的
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)、观察网卡流量与两端负载,重点验证fetchSize与speed.channel对吞吐的影响,同时监控 DataX 进程内存以避免fetchSize过大导致的 OOM。
7 约束限制
7.1 主备同步数据恢复问题
Clickhouse 使用主从灾备时,备库会从主库不间断地通过 binlog 恢复数据。由于主备数据同步存在时间差(尤其在某些特定情况如网络延迟下),备库同步恢复的数据可能与主库有较大差别,导致从备库同步的数据不是一份当前时间的完整镜像。
针对这一问题,文档提及计划提供preSql功能以在读取前执行前置 SQL 保障数据可用性,但该功能当前待补充实现(文档原文标注“该功能待补充”)。读者在从备库抽取数据时应自行评估数据新鲜度与一致性要求。
7.2 一致性约束
Clickhouse 在数据存储划分中属于 RDBMS 系统,对外可提供强一致性数据查询接口。例如:当一次同步任务启动运行过程中,若该库存在其他数据写入方写入数据,ClickhouseReader完全不会获取到写入更新数据——这是由数据库本身的快照特性(MVCC,多版本并发控制)决定的。
但上述特性仅在ClickhouseReader 单线程模型下成立。当用户配置splitPk启用并发抽取后,ClickhouseReader 会先后启动多个并发任务:由于多个并发任务相互之间不属于同一个读事务,且并发任务之间存在时间间隔,因此这份数据并不是完整的、一致的数据快照。
针对多线程下的一致性快照需求,技术上目前无法直接实现,只能从工程角度解决,文档给出两种可选思路(各有取舍,用户自行权衡):
- 使用单线程同步,即不再进行数据切分。缺点是速度较慢,但能够很好地保证一致性;
- 关闭其他数据写入方,保证当前数据为静态数据,例如锁表、关闭备库同步等。缺点是可能影响在线业务。
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 卫绾的经验总结):
- SQL 执行计划异常导致抽取时间长:抽取时尽可能使用全表扫描代替索引扫描;
- 合理设置 SQL 的并发度:通过
speed.channel与splitPk合理增大并发,减少抽取时间; - 抽取 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),仅供参考