DataX GaussDbReader 插件详解:基于 JDBC 的 GaussDB 数据抽取实战指南

发布时间:2026/9/21 23:15:21
DataX GaussDbReader 插件详解:基于 JDBC 的 GaussDB 数据抽取实战指南 DataX GaussDbReader 插件详解基于 JDBC 的 GaussDB 数据抽取实战指南【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX导读GaussDbReader 是阿里云 DataWorks 数据集成开源版本 DataX 中的关系型数据库读取插件用于将 GaussDBopenGauss数据库中的数据通过 JDBC 连接批量抽取出来并以 DataX 统一的抽象数据集Record形式传递给下游 Writer从而打通 GaussDB 与其他数据源之间的数据同步链路。读完本文你将掌握 GaussDbReader 的完整配置方法、每个参数的语义与取值范围、表模式与自定义 SQL 两种抽取方式的取舍、splitPk 并行切分的底层原理以及 GaussDB 类型与 DataX 内部类型的映射规则可直接据此编写可运行的 DataX 同步作业。1 快速介绍GaussDbReader 插件实现了从 GaussDB 数据库读取数据的能力。在底层实现上GaussDbReader 通过 JDBC 连接远程 GaussDB 数据库并执行相应的 SQL 语句将数据从 GaussDB 库中 SELECT 出来随后交给下游的 Writer 插件继续处理。从插件代码结构看GaussDbReader 本身非常精简它没有重新实现一套读取逻辑而是基于plugin-rdbms-util模块提供的通用 RDBMS 读取框架完成工作。核心实现位于 GaussDbReader.java其类定义如下public class GaussDbReader extends Reader { private static final DataBaseType DATABASE_TYPE DataBaseType.GaussDB; ... }其中DATABASE_TYPE DataBaseType.GaussDB是关键它在 DataBaseType.java 中定义为GaussDB(gaussdb, org.opengauss.Driver)即 GaussDbReader 使用org.opengauss.Driver作为 JDBC 驱动类。对应地插件在 pom.xml 中声明了依赖org.opengauss:opengauss-jdbc:3.0.0。2 实现原理简而言之GaussDbReader 通过 JDBC 连接器连接到远程的 GaussDB 数据库并根据用户配置的信息生成查询 SELECT SQL 语句并发送到远程 GaussDB 数据库并将该 SQL 执行返回结果使用 DataX 自定义的数据类型拼装为抽象的数据集并传递给下游 Writer 处理。具体处理逻辑分两种模式对于用户配置 Table、Column、Where 的信息GaussDbReader 将其拼接为 SQL 语句发送到 GaussDB 数据库对于用户配置 querySql 信息GaussDbReader 直接将其发送到 GaussDB 数据库。从源码角度看GaussDbReader 的 Job 与 Task 阶段全部委托给通用框架CommonRdbmsReaderJob 阶段init()首先校验并回写fetchSize配置若小于 1 直接抛出REQUIRED_VALUE异常随后实例化CommonRdbmsReader.Job(DATABASE_TYPE)并调用其init()由 OriginalConfPretreatmentUtil.java 完成配置预处理校验username/password必填、去除where末尾分号、识别 table 模式还是 querySql 模式、在多个jdbcUrl中探测可用连接并回写最终 URLsplit(adviceNumber)阶段调用ReaderSplitUtil.doSplit完成切片。Task 阶段startRead()调用CommonRdbmsReader.Task.startRead()内部通过DBUtil.getConnection(...)建立连接、DBUtil.query(conn, querySql, fetchSize)执行查询然后逐行读取 ResultSet由buildRecord()按 JDBC 类型把列值映射为 DataX 的LongColumn、DoubleColumn、StringColumn、DateColumn、BoolColumn、BytesColumn最终通过recordSender.sendToWriter(record)传递给下游。对应的通用查询 SQL 模板定义在 Constant.javapublic static String QUERY_SQL_TEMPLATE_WITHOUT_WHERE select %s from %s ; public static String QUERY_SQL_TEMPLATE select %s from %s where (%s);即 table 模式下生成的 SQL 形如select column1,column2 from table where (where条件)。3 功能说明3.1 配置样例3.1.1 表模式从 GaussDB 同步抽取数据到本地以下配置从 GaussDB 数据库按表同步数据到本地Writer 使用streamwriter将结果打印/输出{ job: { setting: { speed: { byte: 1048576 }, errorLimit: { record: 0, percentage: 0.02 } }, content: [ { reader: { name: gaussdbreader, parameter: { username: xx, password: xx, column: [ id, name ], splitPk: id, connection: [ { table: [ table ], jdbcUrl: [ jdbc:opengauss://host:port/database ] } ] } }, writer: { name: streamwriter, parameter: { print: true } } } ] } }配置要点说明setting.speed.byte设置传输速度单位为 byte/sDataX 运行会尽可能达到该速度但不超过它setting.errorLimit出错限制record为出错的 record 条数上限大于该值即报错percentage为出错 record 百分比上限1.0 表示 100%0.02 表示 2%reader.name固定为gaussdbreader。3.1.2 自定义 SQL 模式按 querySql 抽取在 join、复杂过滤等场景下可以直接配置自定义查询 SQL{ job: { setting: { speed: 1048576 }, content: [ { reader: { name: gaussdbreader, parameter: { username: xx, password: xx, where: , connection: [ { querySql: [ select db_id,on_line_flag from db_info where db_id 10; ], jdbcUrl: [ jdbc:opengauss://host:port/database, jdbc:opengauss://host:port/database ] } ] } }, writer: { name: streamwriter, parameter: { print: false, encoding: UTF-8 } } } ] } }注意querySql模式与表模式互斥配置了querySql后table、column、where都会被忽略详见 3.2 中 querySql 一节。3.2 参数说明jdbcUrl描述到对端数据库的 JDBC 连接信息使用 JSON 数组描述并支持一个库填写多个连接地址。之所以使用 JSON 数组描述连接信息是因为支持多个 IP 探测如果配置了多个地址GaussDbReader 可以依次探测 IP 的可连接性直到选择一个合法的 IP如果全部连接失败GaussDbReader 报错。注意jdbcUrl 必须包含在connection配置单元中。对于外部使用情况JSON 数组填写一个 JDBC 连接即可。jdbcUrl 按照 GaussDB 官方规范并可以填写连接附加控制信息如字符编码、超时等参数具体格式可参考 GaussDB 官方 JDBC 连接说明文档。必选是默认值无从源码看多地址探测逻辑由DBUtil.chooseJdbcUrl(...)完成OriginalConfPretreatmentUtil.java选定后还会调用DATABASE_TYPE.appendJDBCSuffixForReader(jdbcUrl)追加连接后缀。对 GaussDB 而言该方法不追加任何参数见 DataBaseType.java 中 GaussDB 分支为空实现因此连接串完全由用户控制。username描述数据源的用户名。必选是默认值无username在配置预处理阶段会作为必填项校验缺失时抛出REQUIRED_VALUE错误对应originalConfig.getNecessaryValue(Key.USERNAME, ...)。password描述数据源指定用户名的密码。必选是默认值无与username一样password缺失时会在 Job 初始化阶段直接报错。table描述所选取的需要同步的表。使用 JSON 的数组描述因此支持多张表同时抽取。当配置为多张表时用户自己需保证多张表是同一 schema 结构GaussDbReader 不予检查表是否为同一逻辑表。注意table必须包含在connection配置单元中。必选是table 模式下默认值无column描述所配置的表中需要同步的列名集合使用 JSON 的数组描述字段信息。用户使用*代表默认使用所有列配置例如[*]。支持列裁剪即列可以挑选部分列进行导出。支持列换序即列可以不按照表 schema 信息进行导出。支持常量配置用户需要按照 GaussDB 语法格式例如[id, hello::varchar, true, 2.5::real, power(2,3)]其中id为普通列名hello::varchar为字符串常量true为布尔值2.5为浮点数power(2,3)为函数。column 必须用户显式指定同步的列集合不允许为空必选是默认值无splitPk描述GaussDbReader 进行数据抽取时如果指定splitPk表示用户希望使用 splitPk 代表的字段进行数据分片DataX 因此会启动并发任务进行数据同步这样可以大大提高数据同步的效能。推荐 splitPk 使用表主键因为表主键通常情况下比较均匀因此切分出来的分片也不容易出现数据热点。文档明确目前 splitPk 仅支持整型数据切分不支持浮点、字符串型、日期等其他类型。如果用户指定其他非支持类型GaussDbReader 将报错。splitPk 设置为空底层将视作用户不允许对单表进行切分因此使用单通道进行抽取。必选否默认值空从源码结构看切分实现位于 SingleTableSplitUtil.java它先通过getPkRange(configuration)查询切分主键的 min/max 范围再调用RdbmsRangeSplitWrap.splitAndWrap(...)将区间按并发通道数等分为每个切片生成带where splitPk start and splitPk end的查询 SQL当主键区间两侧出现 NULL 时会退化为单切片执行。另外值得注意从该通用框架的代码看字符串类型主键PK_TYPE_STRING同样具备切分路径实际支持情况以当前仓库代码为准。where描述筛选条件GaussDbReader 根据指定的 column、table、where 条件拼接 SQL并根据这个 SQL 进行数据抽取。在实际业务场景中往往会选择当天的数据进行同步可以将 where 条件指定为gmt_create $bizdate。注意不可以将 where 条件指定为limit 10limit 不是 SQL 合法的 where 子句。where 条件可以有效地进行业务增量同步。where 条件不配置或者为空视作全表同步数据。必选否默认值无补充一个源码细节在配置预处理阶段OriginalConfPretreatmentUtil.dealWherewhere 条件末尾的中英文分号;或会被自动去除避免拼接 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 时GaussDbReader 直接忽略 table、column、where 条件的配置。必选否默认值无fetchSize描述该配置项定义了插件和数据库服务器端每次批量数据获取条数该值决定了 DataX 和服务器端的网络交互次数能够较大地提升数据抽取性能。注意该值过大2048可能造成 DataX 进程 OOM。必选否默认值文档描述为 1024当前仓库源码中 Constant.java 定义的DEFAULT_FETCH_SIZE为 1000实际生效默认值以所使用 DataX 版本源码为准。补充源码校验逻辑GaussDbReader.javaJob 初始化时若用户配置的fetchSize小于 1会直接抛出异常提示根据 DataX 的设计fetchSize 设置值不能小于 1。读取阶段该值作为DBUtil.query(conn, querySql, fetchSize)的批量拉取参数传入控制每次与 GaussDB 服务器交互时取回的行数。4 类型转换目前 GaussDbReader 支持大部分 GaussDB 类型但也存在部分个别类型没有支持的情况请注意检查你的类型。下面列出 GaussDbReader 针对 GaussDB 类型转换列表DataX 内部类型GaussDB 数据类型Longbigint, bigserial, integer, smallint, serialDoubledouble precision, money, numeric, realStringvarchar, char, text, bit, inetDatedate, time, timestampBooleanboolBytesbytea该映射表与通用读取框架CommonRdbmsReader.buildRecord()CommonRdbmsReader.java中的 JDBC 类型分支一一对应CHAR/NCHAR/VARCHAR/LONGVARCHAR/NVARCHAR/LONGNVARCHAR→StringColumnSMALLINT/TINYINT/INTEGER/BIGINT→LongColumnNUMERIC/DECIMAL/FLOAT/REAL/DOUBLE→DoubleColumnTIME/DATE/TIMESTAMP→DateColumnBOOLEAN/BIT→BoolColumnBINARY/VARBINARY/BLOB/LONGVARBINARY→BytesColumn。请注意除上述罗列字段类型外其他类型均不支持money、inet、bit需用户使用a_inet::varchar类似的语法转换后再抽取。对于无法直接映射的类型读取框架会抛出UNSUPPORTED_TYPE异常并提示请尝试使用数据库函数将其转换成 DataX 支持的类型 或者不同步该字段此时可在column中借助 GaussDB 的::type强转或内置函数完成转换。若某列读取出错该行会被识别为脏数据记录taskPluginCollector.collectDirtyRecord并由作业级errorLimit控制任务是否中止。5 性能报告以下为 GaussDbReader 在特定软硬件环境下的实测性能数据供读者作为并发通道数与分片策略配置的参考基准数据来自仓库文档具体数值随环境而异。5.1 环境准备5.1.1 数据特征建表语句create table pref_test( id serial, a_bigint bigint, a_bit bit(10), a_boolean boolean, a_char character(5), a_date date, a_double double precision, a_integer integer, a_money money, a_num numeric(10,2), a_real real, a_smallint smallint, a_text text, a_time time, a_timestamp timestamp )该测试表覆盖了 GaussDB 中 Long、Double、String、Date、Boolean 等主要类型可作为类型转换适配性验证的参考样例。5.1.2 机器参数执行 DataX 的机器参数cpu16 核 Intel(R) Xeon(R) CPU E5620 2.40GHzmemMemTotal: 24676836kBMemFree: 6365080kBnet百兆双网卡GaussDB 数据库机器参数D12 24 逻辑核、192G 内存、12 * 480G SSD 阵列5.2 测试报告5.2.1 单表测试报告通道数是否按照主键切分DataX速度(Rec/s)DataX流量(MB/s)DataX机器运行负载1否102110.630.21是102110.630.24否102110.630.24是400002.480.58否102110.630.28是780484.840.8说明这里的单表主键类型为 serial数据分布均匀。对单表如果没有按照主键切分那么配置通道个数不会提升速度效果与 1 个通道一样。这组数据直观印证了splitPk的价值只有配置了切分主键增加通道数才能真正提升吞吐4 通道约 4 倍、8 通道约 7.6 倍不切分时无论配置多少通道都退化为单通道串行抽取。这与 3.2 节中splitPk 为空视作单通道抽取的实现是一致的。6 常见问题与使用建议table 模式与 querySql 模式如何选择单表、按列裁剪/换序、增量 where 过滤用 table 模式多表 join、复杂子查询用 querySql 模式。二者不可混用。多张表抽取的前提table数组配置多张表时必须保证表结构schema一致插件不做一致性校验。常量列与函数列column中支持 GaussDB 语法常量hello::varchar、true、2.5::real与函数power(2,3)可实现抽取时的列级加工。类型不支持时怎么办优先在column中用::varchar等强转语法转换无法转换的字段建议不同步。避免 OOMfetchSize不建议超过 2048推荐保持默认值约 1000。增量同步利用where配合业务时间字段如gmt_create $bizdate实现按天/按批次的增量抽取不要使用limit作为 where 条件。连接可靠性可在jdbcUrl数组配置多个 GaussDB 地址插件会依次探测可用连接全部失败才报错。通过以上配置与原理说明你可以基于 GaussDbReader 快速搭建 GaussDB 到任意 DataX 支持目标端文本、HDFS、MySQL、PostgreSQL 等的数据同步作业并借助 splitPk 并行切片获得可扩展的抽取吞吐。【免费下载链接】DataXDataX是阿里云DataWorks数据集成的开源版本。项目地址: https://gitcode.com/gh_mirrors/da/DataX创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

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

尧图内容编辑团队

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

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

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

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

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

网站改版的5个关键决策

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

获取专属建站方案

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

立即免费咨询