十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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 的工作过程可以概括为以下几步建立连接通过 JDBC 连接器驱动类为ru.yandex.clickhouse.ClickHouseDriver见 DataBaseType.java连接远程 ClickHouse 数据库生成 SQL若用户配置了table、column、where插件将它们拼接为标准的SELECT语句模板见 Constant.javaselect %s from %s where (%s)发送给 ClickHouse若用户配置了querySql插件原样透传该 SQL 到 ClickHouse不做任何拼接或改写结果集转换将 JDBCResultSet逐行读取按列类型映射为 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/serrorLimit用于控制任务可容忍的错误记录数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中读取配置时使用的兜底值为1000getInt(Constant.FETCH_SIZE, 1000)实际生效值以作业配置为准建议在文档建议的 1024 上下根据内存与网络情况调整。session描述控制数据读取时的时间格式、时区等会话配置。如果表中有时间字段可配置该值以明确告知 Clickhouse 读取的时间格式如 NLS_DATE_FORMAT、NLS_TIME_FORMAT 等。其配置值为 JSON 数组格式例如session: [ alter session set NLS_DATE_FORMATyyyy-mm-dd hh24:mi:ss, alter session set NLS_TIMESTAMP_FORMATyyyy-mm-dd hh24:mi:ss, alter session set NLS_TIMESTAMP_TZ_FORMATyyyy-mm-dd hh24:mi:ss, alter session set TIME_ZONEUS/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, Int256DoubleFloat32, Float64, DecimalStringString, FixedStringDateDATE, Date32, DateTime, DateTime64BooleanBooleanBytesBLOB, BFILE, RAW, LONG RAW请注意除上述罗列字段类型外其他类型均不支持。该映射在源码层面与 JDBC 标准类型一一对应见 CommonRdbmsReader.java 的buildRecord以及同逻辑的 ResultSetReadProxy.javaCHAR / NCHAR / VARCHAR / LONGVARCHAR / NVARCHAR / LONGNVARCHAR / CLOB / NCLOB→StringColumnSMALLINT / TINYINT / INTEGER / BIGINT→LongColumn通过rs.getString取值后构造NUMERIC / DECIMAL / FLOAT / REAL / DOUBLE→DoubleColumnTIME / DATE / TIMESTAMP→DateColumn其中 DATE 类型若列类型名为year则映射为LongColumnBINARY / VARBINARY / BLOB / LONGVARBINARY→BytesColumnBOOLEAN / BIT→BoolColumnNULL→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直接为每张表生成一条完整查询 SQLselect %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的兜底分片保证主键为空的数据也不会丢失最终将每个分片作为独立的ConfigurationTask 配置返回由 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常见问题Q1ClickhouseReader 同步报错报错信息为 XXX如何处理A多为网络或权限问题请先使用 Clickhouse 命令行clickhouse-client测试同样的连接与查询。如果命令行同样报错即可证实是环境问题请联系你的 DBA。Q2ClickhouseReader 抽取速度很慢怎么办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),仅供参考
返回列表