news 2026/9/19 22:28:15

SeaTunnel JDBC Oracle 源连接器实战指南:配置、数据类型映射与并行分片原理

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
SeaTunnel JDBC Oracle 源连接器实战指南:配置、数据类型映射与并行分片原理

SeaTunnel JDBC Oracle 源连接器实战指南:配置、数据类型映射与并行分片原理

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

导读

本文是 Apache SeaTunnel(当前仓库为 GitHub 精选的se/seatunnel项目)中JDBC Oracle 源连接器(Source Connector)的完整使用与原理指南。通过 JDBC 协议,你可以把 Oracle 数据库中的表数据或自定义 SQL 查询结果高效、并行地读入 SeaTunnel 数据管道,再交给任意 Sink 插件(如控制台、Kafka、JDBC 目标库等)消费。读完本文,你将掌握 Oracle 源连接器的连接配置、类型映射规则、依赖部署、单表/多表读取、四种并行分片模式以及split.*分片策略调优,并能结合实际源码理解分片背后真正发生了什么。

连接器概述

JDBC Oracle 源连接器通过标准 JDBC 接口读取 Oracle 数据源,连接器标识为Jdbc(配置块名称即Jdbc),其工厂实现位于 JdbcSourceFactory.java,插件制品名为connector-jdbc,已默认登记在 config/plugin_config 中。

支持的计算引擎

连接器原生支持以下三种运行引擎:

  • Spark(批)
  • Flink(批)
  • SeaTunnel Zeta(批)

关键特性矩阵

特性支持情况
批处理✅ 支持
流处理❌ 不支持(Oracle 源本质是批连接器,详见"流式增量 ID 区间读取"一节)
精确一次(Exactly-Once)✅ 支持
列投影✅ 支持(通过自定义query即可实现投影)
并行性✅ 支持
用户自定义 Split✅ 支持
多表读取✅ 支持(通过table_list

其中"列投影"无需额外配置:你可以在query中只SELECT需要的字段,SeaTunnel 会按查询结果集推断 schema,天然实现投影效果。

支持的数据源与数据库依赖

支持的数据源信息

数据源支持的版本驱动连接串Maven 坐标
Oracle不同依赖版本对应不同驱动类oracle.jdbc.OracleDriverjdbc:oracle:thin:@datasource01:1523:xecom.oracle.database.jdbc:ojdbc8

驱动 jar 的部署位置(按引擎区分)

⚠️ 不同引擎下 JDBC 驱动 jar 的放置目录不同,务必按引擎选择。

对于 Spark / Flink 引擎:

  1. 将 ojdbc8 驱动 jar 放入${SEATUNNEL_HOME}/plugins/目录;
  2. 如需支持 i18n 字符集(如中文等多字节字符),额外将orai18n.jar复制到${SEATUNNEL_HOME}/plugins/

对于 SeaTunnel Zeta 引擎:

  1. 将 ojdbc8 驱动 jar 放入${SEATUNNEL_HOME}/lib/目录;
  2. 如需 i18n 字符集,额外将orai18n.jar复制到${SEATUNNEL_HOME}/lib/

说明:orai18n.jar是 Oracle JDBC 驱动国际化支持包,缺失时可能导致 NLS 字符集相关字段读取异常。两个 jar 可从 Maven 中央仓库获取。

数据类型映射

连接器的 Oracle 类型映射由 OracleTypeConverter.java 实现,Oracle 侧类型识别由 OracleTypeMapper.java 完成(读取ResultSetMetaDatagetColumnTypeNamegetPrecisiongetScale等信息构造类型定义)。整体映射规则如下:

Oracle 数据类型SeaTunnel 数据类型
INTEGERDECIMAL(38,0)
FLOATDECIMAL(38, 18)
NUMBER(precision <= 9, scale == 0)INT
NUMBER(9 < precision <= 18, scale == 0)BIGINT
NUMBER(18 < precision, scale == 0)DECIMAL(38, 0)
NUMBER(scale != 0)DECIMAL(38, 18)
BINARY_DOUBLEDOUBLE
BINARY_FLOAT、REALFLOAT
CHAR、NCHAR、VARCHAR、NVARCHAR2、VARCHAR2、LONG、ROWID、NCLOB、CLOB、XML、INTERVALSTRING
DATETIMESTAMP
TIMESTAMP、TIMESTAMP WITH LOCAL TIME ZONETIMESTAMP
BLOB、RAW、LONG RAW、BFILEBYTES

映射细节与边界(来自源码实现)

阅读 OracleTypeConverter.java 的convert方法可以确认几个容易踩坑的细节:

  • NUMBER 的默认精度/尺度:当 JDBC 元数据未给出 precision 时按 38 处理,scale 未给出时按 127 处理(127 是 Oracle 驱动用于表示"未指定尺度"的哨兵值);scale <= 0时按precision - scale重新计算有效精度。
  • INTERVAL 类型名清理:驱动可能返回INTERVAL DAY(2) TO SECOND(6)这类带精度修饰符的类型名,转换器会先剥离括号内的精度限定再精确匹配,最终映射为 STRING。
  • NVARCHAR2/NCHAR 长度换算:OracleTypeMapper 中会把 NVARCHAR2/NCHAR 的双字节长度换算为 SeaTunnel 侧的 4 字节长度(charToDoubleByteLength/doubleByteTo4ByteLength),避免下游 schema 长度失真。
  • BLOB 的字符串化开关handle_blob_as_stringtrue时,BLOB 会被映射为 STRING(而不是 BYTES),该开关当前仅对 Oracle 生效。
  • 字符串列长度上限:CHAR 最大 2000 字节、VARCHAR2 最大 4000 字节、CLOB/NCLOB 最大 4GB-1、LONG 最大 2GB-1,这些常量在转换器中均有定义。

反转换(reconvert)说明

转换器同时实现了reconvert(SeaTunnel 类型 → Oracle DDL 类型),用于把 SeaTunnel 的 DECIMAL 映射回NUMBER(p,s)、BOOLEAN 映射回NUMBER(1)、STRING 按长度选择VARCHAR2(n)CLOB、BYTES 按长度选择RAW(n)BLOB等。这说明同一套 Oracle 类型体系同时服务于 Source(读)与 Sink/Catalog(写)场景。

源选项(Source Options)完整参考

所有参数在 JdbcSourceOptions.java 与 JdbcCommonOptions.java 中定义,工厂通过 JdbcSourceFactory.optionRule() 完成提交期校验(如table_listtable_path/query互斥、where_condition必须以where开头)。

参数名类型必须默认值描述
urlString-JDBC 连接 URL,如jdbc:oracle:thin:@datasource01:1523:xe
driverString-JDBC 驱动类名,Oracle 为oracle.jdbc.OracleDriver
usernameString-连接实例用户名
passwordString-连接实例密码
queryString-查询语句;未配置table_pathtable_list时必填
connection_check_timeout_secInt30等待连接验证数据库操作完成的时间(秒)
partition_columnString-用于并行分片的列名,仅支持数值类型主键,只能配置一列
partition_lower_boundBigDecimal-partition_column的最小值;未设置时 SeaTunnel 自动查询数据库获取
partition_upper_boundBigDecimal-partition_column的最大值;未设置时 SeaTunnel 自动查询数据库获取
partition_numIntjob parallelism分割数量,仅支持正整数,默认等于任务并行度
fetch_sizeInt0查询行提取大小(fetch size);0 表示使用 JDBC 默认值。注意 Oracle 方言在 fetch_size <= 0 时实际使用默认 128(见下文源码说明)
propertiesMap-其他连接配置参数;当 properties 与 URL 中参数相同时,优先级由驱动实现决定,Oracle 下 properties 优先于 URL
use_regexBooleanfalse控制table_path是否按正则表达式匹配;true时 table_path 视为正则,否则视为精确路径
table_pathString-表的完整路径,可替代query,如"test_schema.table1"
table_listArray-要读取的表列表,可替代table_path;元素支持table_path与可选query字段
where_conditionString-所有表/查询的公共行过滤条件,必须以where开头
split.sizeInt8096一个分片包含多少行
split.even-distribution.factor.lower-boundDouble0.05分片键分布因子下限,用于判断数据是否均匀分布
split.even-distribution.factor.upper-boundDouble100分片键分布因子上限,用于判断数据是否均匀分布
split.sample-sharding.thresholdInt1000触发采样分片策略的估算分片数阈值
split.inverse-sampling.rateInt1000采样分片策略的采样率倒数,1000表示 1/1000 采样
split.allow-samplingBooleantrue是否允许对分布不均匀的分片键使用采样分片策略;false时回退到迭代式不均匀分片
use_select_countBooleanfalse分片前是否使用select count(*)估算表行数
skip_analyzeBooleanfalse是否跳过分片前的表行数分析
decimal_type_narrowingBooleantrueDecimal 类型收窄;为true时在无精度损失的前提下把 Oracle Decimal 收窄为 Int/Long,当前仅对 Oracle 生效
common-options--源插件通用参数,参考 源通用选项

fetch_size 的 Oracle 特化行为

一般 JDBC 方言在fetch_size = 0时不做任何设置,但 Oracle 方言在 OracleDialect.creatPreparedStatement 中做了特化:fetchSize > 0时使用配置值,否则强制setFetchSize(128)(常量DEFAULT_ORACLE_FETCH_SIZE = 128)。也就是说 Oracle 下即使不配置fetch_size,游标也会按 128 行一批拉取,这能显著减少网络往返次数。

decimal_type_narrowing 详解

该参数定义于 JdbcCommonOptions.java,描述明确标注"当前仅在 Oracle 上生效"。它在 OracleTypeConverter.convert 的 NUMBER 分支中体现:当scale <= 0且有效精度newPrecision = precision - scale满足条件时触发收窄。

decimal_type_narrowing = true时:

OracleSeaTunnel
NUMBER(1, 0)Boolean
NUMBER(6, 0)INT
NUMBER(10, 0)BIGINT

decimal_type_narrowing = false时:

OracleSeaTunnel
NUMBER(1, 0)Decimal(1, 0)
NUMBER(6, 0)Decimal(6, 0)
NUMBER(10, 0)Decimal(10, 0)

从源码看,收窄逻辑为:有效精度等于 1 映射为 BOOLEAN、小于等于 9 映射为 INT、小于等于 18 映射为 BIGINT;超出 18 且小于 38 则保留Decimal(newPrecision, 0),否则取Decimal(38, 0)。开启收窄可减少下游小数运算开销,但如果下游逻辑强依赖 DECIMAL 语义,建议关闭。

并行读取与分片原理

JDBC Source 连接器支持并行读取表数据。SeaTunnel 会按一定规则把表数据拆分为多个分片(Split),交给 Reader 并行读取,Reader 数量由parallelism选项决定。分片的核心实现位于 ChunkSplitter.java 及其子类(DynamicChunkSplitterFixedChunkSplitter),分片枚举与下发由 JdbcSourceSplitEnumerator.java 负责,实际读取由 JdbcSourceReader.java 完成。

分片键(Split Key)选择规则

  1. 如果设置了partition_column,则用该列计算分片。该列必须在支持的分片数据类型中。
  2. 如果partition_column未设置,SeaTunnel 会读取表结构并取主键和唯一索引。如果主键或唯一索引由多个列组成,则取第一个落在支持的分片数据类型中的列。例如表的主键是(guid, name varchar),由于guid不在支持的分片类型中,会使用name列进行分片。

支持的分片数据类型:

  • String
  • Number(int、bigint、decimal 等)
  • Date

分片相关选项详解

split.size

一个分片包含多少行;读取表时,捕获的表会按split.size拆分为多个分片。默认8096

split.even-distribution.factor.lower-bound

⚠️ 不推荐修改

分片键分布因子的下限。分布因子用于判断表数据是否均匀分布,计算公式为(MAX(id) - MIN(id) + 1) / 行数。如果计算得到的分布因子大于等于该下限,则按均匀分布拆分;否则视为分布不均匀,并在估算分片数超过split.sample-sharding.threshold时使用采样分片策略。默认0.05

split.even-distribution.factor.upper-bound

⚠️ 不推荐修改

分片键分布因子的上限。如果分布因子小于等于该上限,按均匀分布拆分;否则视为不均匀分布,使用采样分片策略。默认100.0

split.sample-sharding.threshold

触发采样分片策略的估算分片数阈值。当分布因子超出上下限区间,且估算分片数(行数 /split.size)超过该阈值时,会使用采样分片策略。默认1000

split.inverse-sampling.rate

采样分片策略的采样率倒数。例如1000表示 1/1000 采样率。默认1000

split.allow-sampling

是否允许对分布不均匀的分片键使用采样分片策略。设置为false时,SeaTunnel 会回退到迭代式不均匀分片(iterative query 方式逐个拉取分片边界)。默认true

partition_column

用于分片的列名。仅支持数值类型,只能配置一列。

partition_upper_bound

partition_column的扫描最大值。未设置时 SeaTunnel 会查询数据库获取。

partition_lower_bound

partition_column的扫描最小值。未设置时 SeaTunnel 会查询数据库获取。

partition_num

⚠️ 不推荐修改,推荐使用split.size控制分片大小。

将数据拆分为多少个分片,仅支持正整数。默认等于任务并行度。它是分片数量上限:实际分片数会在满足partition_num的前提下,由边界值按范围均匀切分得出。

均匀分布判断与采样分片的整体流程

结合 JdbcSourceOptions.java 中的选项描述与源码可以还原如下决策链:

  1. 计算分片键分布因子(MAX(id) - MIN(id) + 1) / 行数
  2. 若分布因子落在[lower-bound, upper-bound]区间内 → 判定数据均匀,使用均匀分片优化(直接按范围等分);
  3. 若分布因子超出区间且估算分片数超过split.sample-sharding.threshold→ 触发采样分片策略:按1 / split.inverse-sampling.rate采样率抽样分片键值,基于采样分布切分;
  4. 若不允许采样(split.allow-sampling = false)→ 回退到迭代式不均匀分片,通过循环执行queryNextChunkMax(Oracle 实现见 OracleDialect.queryNextChunkMax,内部使用ORDER BY ... ASC+ROWNUM <= chunkSize取块内最大值)逐个确定分片边界。

Oracle 行数估算的特殊实现

OracleDialect.approximateRowCntStatement 体现了 Oracle 分片的行数估算策略:

  • 未配置query,或query不含 WHERE 且配置了table_path时,优先走Oracle 表统计信息:执行analyze table <表> compute statistics for table(除非skip_analyze = true),然后查询all_tablesNUM_ROWS
  • 其余情况(query 含 WHERE、或只有 query 无 table_path)使用select count(*)
  • 显式设置use_select_count = true时强制走count(*)路径。

这也解释了skip_analyzeuse_select_count的用途:前者跳过 ANALYZE 语句避免对大表造成统计开销,后者强制使用精确 count 而非统计值。

提示与限制

  • 如果表无法拆分(例如表没有主键或唯一索引,且未设置partition_column),将以单并发运行。
  • 单表读取可使用table_path代替query。需要读取多张表时,请使用table_list
  • table_listtable_path/query是互斥的(由JdbcSourceFactoryTableListExclusiveValidator在提交期强制校验),同一作业只能选择一种表选择模式。

任务示例

以下示例均来自 docs/zh/connectors/source/Oracle.md,配置格式为 HOCON,可直接放入 SeaTunnel 配置文件(如config/v2.batch.config.template的样式)执行。

简单示例(单并行全量查询)

该示例从 Oracle 的 test 数据库查询TEST_TABLE的 16 条数据,以单并行方式运行,并查询其所有字段;你也可以指定要查询的字段(投影),最终输出到控制台。

env { parallelism = 4 job.mode = "BATCH" } source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" username = "root" password = "123456" query = "SELECT * FROM TEST_TABLE" } } transform { # 更多 transform 插件配置请参考 transforms 相关文档 } sink { Console {} }

说明:示例中env.parallelism = 4定义了任务的整体并行度;query未配置分片键时以单并发读取,parallelism在这里决定的是下游算子并行度。

按 partition_column 并行

通过配置分片字段和分片数据,可以并行读取查询表中的数据。需要读取整张表时可以使用这种方式。

env { parallelism = 4 job.mode = "BATCH" } source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" connection_check_timeout_sec = 100 username = "root" password = "123456" # 按需定义查询逻辑 query = "SELECT * FROM TEST_TABLE" # 用于并行分片读取的字段 partition_column = "ID" # 分片数量 partition_num = 10 properties { database.oracle.jdbc.timezoneAsRegion = "false" } } } sink { Console {} }

说明:properties.database.oracle.jdbc.timezoneAsRegion = "false"是 Oracle JDBC 驱动的连接属性,用于控制驱动对 TIMESTAMP WITH TIME ZONE 的时区处理方式。当properties与 URL 中出现同名参数时,Oracle 驱动下properties优先。

按主键或唯一索引并行

配置table_path会开启自动分片(自动选取主键/唯一索引中支持分片数据类型的列),可通过split.*选项调整分片策略。

env { parallelism = 4 job.mode = "BATCH" } source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" connection_check_timeout_sec = 100 username = "root" password = "123456" table_path = "DA.SCHEMA1.TABLE1" query = "select * from SCHEMA1.TABLE1" split.size = 10000 } } sink { Console {} }

说明:table_path格式为库.模式.表(如DA.SCHEMA1.TABLE1),此处同时给出query用于约束读取内容。配置table_path后分片键自动解析,建议用split.size(而非partition_num)控制分片粒度。

并行边界(显式上下界)

显式指定查询的上下界可以更高效地读取数据源(避免额外的MIN/MAX查询开销)。

source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" connection_check_timeout_sec = 100 username = "root" password = "123456" # 按需定义查询逻辑 query = "SELECT * FROM TEST_TABLE" partition_column = "ID" # 读取起点 partition_lower_bound = 1 # 读取终点 partition_upper_bound = 500 partition_num = 10 } }

说明:这里把[1, 500]区间交给partition_num = 10个分片,每个分片负责约 50 个 ID 的范围。partition_lower_bound/partition_upper_bound类型为 BigDecimal。

多表读取(table_list)

配置table_list会开启自动分片,可通过split.*选项调整分片策略。table_list的每个元素支持独立配置table_path与可选query

env { job.mode = "BATCH" parallelism = 4 } source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" connection_check_timeout_sec = 100 username = "root" password = "123456" "table_list" = [ { "table_path" = "XE.TEST.USER_INFO" }, { "table_path" = "XE.TEST.YOURTABLENAME" } ] #where_condition = "where id > 100" split.size = 10000 #split.even-distribution.factor.upper-bound = 100 #split.even-distribution.factor.lower-bound = 0.05 #split.sample-sharding.threshold = 1000 #split.inverse-sampling.rate = 1000 } } sink { Console {} }

说明:table_listtable_path/query互斥,不可同时配置。被注释掉的split.*选项展示了对多表分片策略进行调优的入口。

流式增量 ID 区间读取

Oracle Source 本质上是一个批连接器。设置job.mode = "STREAMING"只用于开启 checkpoint 以便在失败时恢复作业;source 本身仍然是有界的,每次作业只会读取一次配置好的[partition_lower_bound, partition_upper_bound)区间。如需周期性地拉取新增数据,必须在外部重新提交作业(例如按计划滑动区间窗口),或改用 Oracle-CDC 做持续变更捕获。

env { parallelism = 4 job.mode = "STREAMING" checkpoint.interval = 60000 } source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" username = "root" password = "123456" query = "SELECT * FROM ORDERS WHERE ORDER_ID >= ? AND ORDER_ID < ?" partition_column = "ORDER_ID" partition_lower_bound = 1 partition_upper_bound = 1000000 partition_num = 16 } }

说明:checkpoint.interval = 60000让作业每分钟做一次 checkpoint;一旦失败恢复,作业会从上次 checkpoint 状态重新读取(分片边界与已读进度会随 checkpoint 持久化),这正是"流式模式下 checkpoint 恢复"的用法。查询语句中的?占位符由 SeaTunnel 用分片边界绑定。

使用 TNS 连接串

如果 Oracle 部署只暴露 TNS 别名,可以把url指向 TNS 别名。TNS 名称由 classpath 上的oracle.net.tns_admin解析。

source { Jdbc { url = "jdbc:oracle:thin:@tns_alias" driver = "oracle.jdbc.OracleDriver" username = "root" password = "123456" properties { oracle.net.tns_admin = "/etc/oracle" } table_path = "SCHEMA.ORDERS" split.size = 10000 } }

说明:oracle.net.tns_admin指定 tnsnames.ora 所在目录,@tns_alias中的tns_alias即 tnsnames.ora 中定义的连接别名。该场景下 Oracle 侧通常不暴露主机端口,因此 url 不再使用host:port:service形式。

使用 where_condition 过滤行

通过where_condition可以为table_listquery中的所有条目应用一个公共过滤条件。字符串必须以where开头,以便拼接到自定义查询或table_path自动生成的查询后。

source { Jdbc { url = "jdbc:oracle:thin:@datasource01:1523:xe" driver = "oracle.jdbc.OracleDriver" username = "root" password = "123456" table_path = "SCHEMA.ORDERS" where_condition = "where status = 'ACTIVE' and created_at >= DATE '2026-01-01'" split.size = 10000 } }

说明:提交时JdbcSourceFactoryWhereConditionPrefixValidator会校验该字符串必须以where开头(大小写不敏感),防止拼出非法 SQL。该条件会作用于所有表/查询条目,适合多表场景下的统一增量或过滤需求。

变更日志

连接器的历史变更记录维护在连接器变更日志中(docs/zh/connectors/changelog/connector-jdbc.md),可在仓库内查看connector-jdbc的版本演进与行为调整。

进一步阅读

  • JDBC Oracle 源连接器英文文档
  • 连接器通用特性说明
  • 源连接器通用选项
  • SeaTunnel 配置示例
  • 源码入口:JdbcSourceFactory.java、OracleDialect.java、OracleTypeConverter.java

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

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

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

ISO/TS 30431人力资本报告XML数据字典解析与Python提取实践

简介&#xff1a;ISO/TS 30431:2021 是国际标准化组织发布的人力资源管理领域技术规范&#xff0c;围绕领导力指标簇的标准化定义展开&#xff0c;面向人力资源管理者、HR数据分析师、企业培训与组织发展负责人&#xff0c;帮助组织解决领导力评估指标口径不一、难以横向比较的…

作者头像 李华