SeaTunnel JDBC OceanBase 源连接器实战指南:兼容模式、分片并行与多表读取

发布时间:2026/9/19 11:35:59
SeaTunnel JDBC OceanBase 源连接器实战指南:兼容模式、分片并行与多表读取 SeaTunnel JDBC OceanBase 源连接器实战指南兼容模式、分片并行与多表读取【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文围绕 Apache SeaTunnel 的 JDBC OceanBase 源连接器Jdbccom.oceanbase.jdbc.Driver系统讲解如何通过 JDBC 读取 OceanBase 数据包括 MySQL / Oracle 双兼容模式的方言选择机制、完整的源选项参数表、MySQL 与 Oracle 两种模式下的数据类型映射、单表/并行/分片边界/表路径/多表/正则匹配/流式增量等多种读取场景的完整配置以及并行分片在底层是如何实现的。读完本文你将能够独立编写从简单查询到复杂多表并行读取的 OceanBase 数据集成任务。支持这些引擎SparkFlinkSeaTunnel Zeta关键特性特性支持情况批模式✅ 支持流模式❌ 不支持见流式增量区间读取说明精确一次✅ 支持列投影✅ 支持并行性✅ 支持支持用户自定义分片✅ 支持支持多表读取✅ 支持描述通过 JDBC 读取 OceanBase 数据。OceanBase 支持MySQL 兼容模式和Oracle 兼容模式两种租户形态因此 OceanBase 任务应将compatible_mode设置为mysql或oracle。该参数直接决定了 SeaTunnel 加载哪种 JDBC 方言Dialect进而影响 SQL 生成、标识符引用方式、类型映射与行数统计等全部行为。支持的数据源信息数据源支持的版本驱动连接串MavenOceanBase所有 OceanBase 服务器版本com.oceanbase.jdbc.Driverjdbc:oceanbase://localhost:2883/test可从 Maven 中央仓库搜索com.oceanbase:oceanbase-client获取数据库依赖请下载对应 Maven 坐标com.oceanbase:oceanbase-client的驱动包并将其复制到$SEATUNNEL_HOME/plugins/jdbc/lib/工作目录例如cp oceanbase-client-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/注意connector-jdbc是连接器模块而oceanbase-client属于数据库驱动依赖二者都需要就位驱动缺失时作业会在建连阶段直接失败。底层原理compatible_mode 如何决定方言OceanBase 连接器在源码层面同时复用了 MySQL 系与 Oracle 系两条方言实现切换入口就是compatible_mode。相关证据位于 OceanBaseDialectFactory.java工厂通过AutoService(JdbcDialectFactory.class)注册acceptsURL只接受jdbc:oceanbase:开头的连接串create(compatibleMode, fieldIde)中当compatibleMode为oracle不区分大小写时返回OracleDialect其余情况一律返回OceanBaseMysqlDialect若调用不带compatibleMode的create()工厂会直接抛出UnsupportedOperationExceptionCant create JdbcDialect without compatible mode for OceanBase这正是文档中compatible_mode被标记为必填的原因。方言的装载过程在 JdbcDialectLoader.java 中完成通过ServiceLoader扫描 classpath 上所有JdbcDialectFactory优先按显式配置的dialect参数名匹配否则按 URL 前缀匹配当compatible_mode缺失时无法完成方言构造作业配置校验即会失败。以OceanBaseMysqlDialect为例见 OceanBaseMysqlDialect.javaMySQL 兼容模式下还做了这些细节处理标识符使用反引号引用defaultParameter()默认注入rewriteBatchedStatementstrue与allowMultiQueriestrue对应文档提示中JDBC URL 通常会带上rewriteBatchedStatementstrue的说明行数统计优先使用SHOW TABLE STATUS LIKE ...比COUNT(*)快但精度略低仅当配置了带 WHERE 的查询或无法定位表路径时才回退到COUNT(*)子查询统计这直接支撑了下方skip_analyze、use_select_count等分片分析参数的语义。数据类型映射MySQL 模式MySQL 数据类型SeaTunnel 数据类型BIT(1)TINYINT(1)BOOLEANTINYINTBYTETINYINTTINYINT UNSIGNEDSMALLINTSMALLINT UNSIGNEDMEDIUMINTMEDIUMINT UNSIGNEDINTINTEGERYEARINTINT UNSIGNEDINTEGER UNSIGNEDBIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20,0)DECIMAL(x,y)(38)DECIMAL(x,y)DECIMAL(x,y)(38)DECIMAL(38,18)DECIMAL UNSIGNEDDECIMALFLOATFLOAT UNSIGNEDFLOATDOUBLEDOUBLE UNSIGNEDDOUBLECHARVARCHARTINYTEXTMEDIUMTEXTTEXTLONGTEXTJSONENUMSTRINGDATEDATETIMETIMEDATETIMETIMESTAMPTIMESTAMPTINYBLOBMEDIUMBLOBBLOBLONGBLOBBINARYVARBINARBIT(n)GEOMETRYBYTES上述映射与 OceanBaseMySqlTypeConverter.java 的实现一一对应几个值得注意的细节BIT(1)与tinyint(1)都被收窄为BOOLEANBIT(n)、n1则按字节折算为BYTESBIGINT UNSIGNED超出 64 位有符号范围因此映射为DECIMAL(20,0)以保精度TINYINT UNSIGNED因取值范围上探到 255 而提升为SMALLINTint_type_narrowing默认true控制 MySQL 模式下TINYINT(1) → BOOLEAN的收窄行为是否生效。Oracle 模式Oracle 数据类型SeaTunnel 数据类型IntegerDECIMAL(38,0)Number(p), p 9INTNumber(p), p 18BIGINTNumber(p), p 18DECIMAL(38,18)Number(p,s)DECIMAL(p,s)FloatDECIMAL(38,18)REALBINARY_FLOATFLOATBINARY_DOUBLEDOUBLECHARNCHARVARCHARVARCHAR2NVARCHAR2NCLOBCLOBLONGXMLROWIDSTRINGDATETIMESTAMPTIMESTAMPTIMESTAMP WITH LOCAL TIME ZONETIMESTAMPBLOBRAWLONG RAWBFILEBYTESUNKNOWN暂不支持Oracle 兼容模式下decimal_type_narrowing默认true控制在不丢失精度的前提下将DECIMAL收窄为INT/BIGINT的优化是否开启。如果下游对类型精度极其敏感可将其显式置为false。源选项参数名类型必须默认值描述urlString是-JDBC 连接的 URL。参考示例jdbc:oceanbase://localhost:2883/testdriverString是-用于连接到远程数据源的 jdbc 类名应为com.oceanbase.jdbc.Driver。usernameString否-连接实例用户名passwordString否-连接实例密码compatible_modeString是-OceanBase 的兼容模式可以是mysql或oracle。queryString否-查询语句。query、table_path、table_list三者至少配置一个。table_pathString否-完整表路径可替代query使用例如test.source。table_listArray否-待读取的表列表用于多表读取。每个表配置中可以包含table_path、query、partition_column等表级参数。where_conditionString否-所有表或查询共用的过滤条件必须以where开头例如where id 100。connection_check_timeout_secInt否30等待用于验证连接的数据库操作完成的时间秒partition_columnString否-用于并行性分割的列名仅支持数值类型列和字符串类型列。partition_lower_boundBigDecimal否-partition_column 的最小值用于扫描如果未设置SeaTunnel 将查询数据库获取最小值。partition_upper_boundBigDecimal否-partition_column 的最大值用于扫描如果未设置SeaTunnel 将查询数据库获取最大值。partition_numInt否job parallelism分片数量仅支持正整数。使用table_path读取时推荐通过split.size控制单个分片大小。fetch_sizeInt否0对于返回大量对象的查询您可以配置查询中使用的行提取大小以通过减少满足选择条件所需的数据库命中次数来提高性能。零表示使用 jdbc 默认值。split.sizeInt否8096使用table_path读取时每个分片包含的行数。split.even-distribution.factor.lower-boundDouble否0.05判断分片键数据是否均匀分布的下限。split.even-distribution.factor.upper-boundDouble否100判断分片键数据是否均匀分布的上限。split.sample-sharding.thresholdInt否1000数据分布不均时触发采样分片的预估分片数阈值。split.inverse-sampling.rateInt否1000采样分片使用的采样率分母。split.allow-samplingBoolean否true是否允许使用采样分片策略。split.string_split_modeString否sample字符串分片算法可选sample、charset_based。split.string-strategyString否-字符串分片策略可选none、hash、range、auto。split.string_split_mode_collateString否-split.string_split_mode为charset_based时使用的排序规则。use_select_countBoolean否false动态分片阶段使用select count统计行数主要用于 Oracle 兼容读取场景。skip_analyzeBoolean否false动态分片阶段跳过表行数分析主要用于 Oracle 兼容读取场景。use_regexBoolean否false是否将table_path当作正则表达式匹配表。decimal_type_narrowingBoolean否trueOracle 兼容模式下在不丢失精度时将 Decimal 收窄为 INT 或 BIGINT。int_type_narrowingBoolean否trueMySQL 兼容模式下在不丢失精度时将TINYINT(1)收窄为 BOOLEAN。dialectString否-指定 JDBC 方言。OceanBase 通常会根据 URL 自动识别只有特殊兼容场景才需要显式配置。propertiesMap否-其他连接配置参数当 properties 和 URL 具有相同参数时优先级由驱动程序的具体实现确定。例如在 MySQL 中properties 优先于 URL。common-options否-源插件通用参数详见 源通用选项仓库路径source-common-options.md。以上参数在 JdbcSourceFactory.java 的optionRule()中全部注册为可选项除url、driver必填外其中还内嵌了两个提交期校验器值得留意TableListExclusiveValidatortable_list与table_path/query互斥不能同时配置WhereConditionPrefixValidatorwhere_condition必须以where关键字开头不区分大小写否则提交时直接报错避免在运行期拼出畸形 SQL。参数定义细节默认值与描述可在 JdbcSourceOptions.java 中查阅。提示query、table_path、table_list三者至少配置一个。如果未设置partition_column并且 SeaTunnel 无法从表元数据中找到合适的主键或唯一键则源端会以单并发读取。配置了支持的分片列后SeaTunnel 可以并行读取。OceanBase MySQL 模式的 JDBC URL 通常会带上rewriteBatchedStatementstrue等 MySQL 兼容参数OceanBase Oracle 模式需要使用 Oracle 兼容租户并配置compatible_mode oracle。任务示例以下示例均省略了transform块如需了解 transform 插件可参阅 SQL Transform 等文档。示例中的Jdbc即连接器标识符factoryIdentifier。简单最基本的整表查询读取单并发顺序扫描env { parallelism 2 job.mode BATCH } source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue username root password compatible_mode mysql query select * from source } } transform { # 如果您想了解有关如何配置 seatunnel 的更多信息并查看完整的转换插件列表 # 请访问 https://seatunnel.apache.org/docs/transforms/sql } sink { Console {} }并行使用您配置的分片字段和分片数据并行读取查询表。如果您想读取整个表可以这样做env { parallelism 10 job.mode BATCH } source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue username root password compatible_mode mysql query select * from source # 并行分片读取字段 partition_column id # 分片数量 partition_num 10 } } sink { Console {} }配置了partition_column后源端会先探测该列的上下界再将整个区间切成partition_num个子区间并行读取未配置partition_column时退化为单并发读取。并行边界根据您配置的上下边界读取数据源更高效避免全表探测最小值/最大值带来的额外查询source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue username root password compatible_mode mysql query select * from source partition_column id partition_num 10 # 读取开始边界 partition_lower_bound 1 # 读取结束边界 partition_upper_bound 500 } }上界500之外的数据不会被本次作业读取这是将区间扫描与行过滤结合的常用手段。表路径读取希望 SeaTunnel 自动发现表结构并进行分片时可以使用table_path。此时不需要手写 SQL连接器通过 Catalog 读取表元数据并自动完成类型映射与分片source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test username roottest password compatible_mode mysql table_path test.source split.size 8096 } }split.size默认 8096控制每个分片的行数表行数已知后预估分片数 ≈ 行数 /split.size。注意table_path模式下应优先用split.size控制分片粒度而不是partition_num。Oracle 兼容模式连接 Oracle 兼容租户时compatible_mode必须为oracle用户名常带有租户后缀SQL 使用大写标识符更稳妥source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/TESTUSER username TESTUSERtest password compatible_mode oracle query SELECT ID, NAME, CREATE_TIME FROM SOURCE } }此时方言切换为 Oracle 方言标识符引用、Number类型映射见上文 Oracle 模式类型表、行数统计等行为全部按 Oracle 语义执行。多表读取table_list允许一次性注册多张表并可通过where_condition施加公共过滤source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test username roottest password compatible_mode mysql table_list [ { table_path test.source_1 }, { table_path test.source_2 } ] where_condition where id 100 } }where_condition对table_list中所有表统一生效且必须以where开头源码中的WhereConditionPrefixValidator会在提交期强制校验。同时注意table_list与顶层table_path/query互斥。单表自定义 SQL当table_list中的多张表需要不同的 SQL 过滤或投影时可以在每个条目上单独设置query让 SeaTunnel 直接按这条 SQL 读取跳过表元数据查找source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test username roottest password compatible_mode mysql table_list [ { table_path test.orders query select id, amount, status from orders where status PAID }, { table_path test.refunds query select id, order_id, amount from refunds where amount 0 } ] } }每个表条目内的query是独立的可以实现每表各一套投影 过滤的精细化控制同时保留table_list的多表并行能力。表名正则匹配当 OceanBase 租户中存在大量结构相似的表例如按时间分区的orders_2024_q1、orders_2024_q2...设置use_regex true并在table_path中传入正则表达式。SeaTunnel 会枚举匹配的表并按partition_num进行并行读取source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test username roottest password compatible_mode mysql table_path test.orders_2024_q[1-4] use_regex true partition_column id partition_num 8 } }正则按库.表的完整路径匹配test.orders_2024_q[1-4]命中 4 张季度表适合时间分桶、分表等批量接入场景partition_column与partition_num会被应用到每张匹配到的表上。流式增量区间读取OceanBase Source 本质上是一个批连接器。设置job.mode STREAMING只用于开启 checkpoint 以便在失败时恢复作业source 本身仍然是有界的每次作业只会读取一次配置好的partition_lower_bound, partition_upper_bound)区间。如需周期性地拉取新增数据必须在外部重新提交作业例如按计划滑动区间窗口或改用 OceanBase CDC 做持续变更捕获env { parallelism 4 job.mode STREAMING checkpoint.interval 60000 } source { Jdbc { driver com.oceanbase.jdbc.Driver url jdbc:oceanbase://localhost:2883/test username roottest password compatible_mode mysql query select * from orders where id ? and id ? partition_column id partition_lower_bound 1 partition_upper_bound 1000000 partition_num 16 } }该示例利用参数化查询 分片区间实现伪流式增量拉取每次提交作业读取固定区间配合外部调度实现滑动窗口若需要实时捕获 OceanBase 的变更数据应转向 [CDC 相关连接器仓库中存在 connector-cdc-base 等实现。分片并行背后的源码实现上文示例反复出现的partition_column、partition_num、split.*系列参数最终都汇聚到 source 目录 下的分片组件中。从源码结构看分片逻辑大体分为两类固定分片FixedChunkSplitter.java直接按partition_lower_bound/partition_upper_bound/partition_num将数值区间均分对应并行并行边界两个示例动态分片DynamicChunkSplitter.java针对table_path自动发现场景先估算行数近似行数统计在方言层完成例如SHOW TABLE STATUS再按split.size计算分片数随后通过分布因子(MAX - MIN 1) / rowCount判断数据是否均匀上下界即split.even-distribution.factor.*数据倾斜时触发采样分片split.sample-sharding.threshold、split.inverse-sampling.rate、split.allow-sampling共同控制字符串类型分片则由AsciiStringRangeSplitter、CollationBasedSplitter与StringSplitMode/StringSplitStrategy支撑。当未配置任何分片列且表元数据中找不到合适的主键/唯一键时连接器退化为单分片顺序读取——这正是提示章节中单并发行为的实现根因。变更日志本连接器的历史变更记录随 connector-jdbc 变更日志 维护涉及参数增删、行为调整与缺陷修复的版本演进均可在此查阅。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

关于本文作者

来自尧图内容编辑团队

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

尧图内容编辑团队

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

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

延伸阅读

相关资讯与近期热门内容

深度阅读推荐

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

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

网站改版的5个关键决策

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

获取专属建站方案

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

立即免费咨询