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.OracleDriver | jdbc:oracle:thin:@datasource01:1523:xe | com.oracle.database.jdbc:ojdbc8 |
驱动 jar 的部署位置(按引擎区分)
⚠️ 不同引擎下 JDBC 驱动 jar 的放置目录不同,务必按引擎选择。
对于 Spark / Flink 引擎:
- 将 ojdbc8 驱动 jar 放入
${SEATUNNEL_HOME}/plugins/目录; - 如需支持 i18n 字符集(如中文等多字节字符),额外将
orai18n.jar复制到${SEATUNNEL_HOME}/plugins/。
对于 SeaTunnel Zeta 引擎:
- 将 ojdbc8 驱动 jar 放入
${SEATUNNEL_HOME}/lib/目录; - 如需 i18n 字符集,额外将
orai18n.jar复制到${SEATUNNEL_HOME}/lib/。
说明:
orai18n.jar是 Oracle JDBC 驱动国际化支持包,缺失时可能导致 NLS 字符集相关字段读取异常。两个 jar 可从 Maven 中央仓库获取。
数据类型映射
连接器的 Oracle 类型映射由 OracleTypeConverter.java 实现,Oracle 侧类型识别由 OracleTypeMapper.java 完成(读取ResultSetMetaData的getColumnTypeName、getPrecision、getScale等信息构造类型定义)。整体映射规则如下:
| Oracle 数据类型 | SeaTunnel 数据类型 |
|---|---|
| INTEGER | DECIMAL(38,0) |
| FLOAT | DECIMAL(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_DOUBLE | DOUBLE |
| BINARY_FLOAT、REAL | FLOAT |
| CHAR、NCHAR、VARCHAR、NVARCHAR2、VARCHAR2、LONG、ROWID、NCLOB、CLOB、XML、INTERVAL | STRING |
| DATE | TIMESTAMP |
| TIMESTAMP、TIMESTAMP WITH LOCAL TIME ZONE | TIMESTAMP |
| BLOB、RAW、LONG RAW、BFILE | BYTES |
映射细节与边界(来自源码实现)
阅读 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_string为true时,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_list与table_path/query互斥、where_condition必须以where开头)。
| 参数名 | 类型 | 必须 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL,如jdbc:oracle:thin:@datasource01:1523:xe |
| driver | String | 是 | - | JDBC 驱动类名,Oracle 为oracle.jdbc.OracleDriver |
| username | String | 否 | - | 连接实例用户名 |
| password | String | 否 | - | 连接实例密码 |
| query | String | 否 | - | 查询语句;未配置table_path和table_list时必填 |
| connection_check_timeout_sec | Int | 否 | 30 | 等待连接验证数据库操作完成的时间(秒) |
| partition_column | String | 否 | - | 用于并行分片的列名,仅支持数值类型主键,只能配置一列 |
| partition_lower_bound | BigDecimal | 否 | - | partition_column的最小值;未设置时 SeaTunnel 自动查询数据库获取 |
| partition_upper_bound | BigDecimal | 否 | - | partition_column的最大值;未设置时 SeaTunnel 自动查询数据库获取 |
| partition_num | Int | 否 | job parallelism | 分割数量,仅支持正整数,默认等于任务并行度 |
| fetch_size | Int | 否 | 0 | 查询行提取大小(fetch size);0 表示使用 JDBC 默认值。注意 Oracle 方言在 fetch_size <= 0 时实际使用默认 128(见下文源码说明) |
| properties | Map | 否 | - | 其他连接配置参数;当 properties 与 URL 中参数相同时,优先级由驱动实现决定,Oracle 下 properties 优先于 URL |
| use_regex | Boolean | 否 | false | 控制table_path是否按正则表达式匹配;true时 table_path 视为正则,否则视为精确路径 |
| table_path | String | 否 | - | 表的完整路径,可替代query,如"test_schema.table1" |
| table_list | Array | 否 | - | 要读取的表列表,可替代table_path;元素支持table_path与可选query字段 |
| where_condition | String | 否 | - | 所有表/查询的公共行过滤条件,必须以where开头 |
| split.size | Int | 否 | 8096 | 一个分片包含多少行 |
| split.even-distribution.factor.lower-bound | Double | 否 | 0.05 | 分片键分布因子下限,用于判断数据是否均匀分布 |
| split.even-distribution.factor.upper-bound | Double | 否 | 100 | 分片键分布因子上限,用于判断数据是否均匀分布 |
| split.sample-sharding.threshold | Int | 否 | 1000 | 触发采样分片策略的估算分片数阈值 |
| split.inverse-sampling.rate | Int | 否 | 1000 | 采样分片策略的采样率倒数,1000表示 1/1000 采样 |
| split.allow-sampling | Boolean | 否 | true | 是否允许对分布不均匀的分片键使用采样分片策略;false时回退到迭代式不均匀分片 |
| use_select_count | Boolean | 否 | false | 分片前是否使用select count(*)估算表行数 |
| skip_analyze | Boolean | 否 | false | 是否跳过分片前的表行数分析 |
| decimal_type_narrowing | Boolean | 否 | true | Decimal 类型收窄;为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时:
| Oracle | SeaTunnel |
|---|---|
| NUMBER(1, 0) | Boolean |
| NUMBER(6, 0) | INT |
| NUMBER(10, 0) | BIGINT |
decimal_type_narrowing = false时:
| Oracle | SeaTunnel |
|---|---|
| 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 及其子类(DynamicChunkSplitter、FixedChunkSplitter),分片枚举与下发由 JdbcSourceSplitEnumerator.java 负责,实际读取由 JdbcSourceReader.java 完成。
分片键(Split Key)选择规则
- 如果设置了
partition_column,则用该列计算分片。该列必须在支持的分片数据类型中。 - 如果
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 中的选项描述与源码可以还原如下决策链:
- 计算分片键分布因子
(MAX(id) - MIN(id) + 1) / 行数; - 若分布因子落在
[lower-bound, upper-bound]区间内 → 判定数据均匀,使用均匀分片优化(直接按范围等分); - 若分布因子超出区间且估算分片数超过
split.sample-sharding.threshold→ 触发采样分片策略:按1 / split.inverse-sampling.rate采样率抽样分片键值,基于采样分布切分; - 若不允许采样(
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_tables的NUM_ROWS; - 其余情况(query 含 WHERE、或只有 query 无 table_path)使用
select count(*); - 显式设置
use_select_count = true时强制走count(*)路径。
这也解释了skip_analyze与use_select_count的用途:前者跳过 ANALYZE 语句避免对大表造成统计开销,后者强制使用精确 count 而非统计值。
提示与限制
- 如果表无法拆分(例如表没有主键或唯一索引,且未设置
partition_column),将以单并发运行。- 单表读取可使用
table_path代替query。需要读取多张表时,请使用table_list。table_list与table_path/query是互斥的(由JdbcSourceFactory的TableListExclusiveValidator在提交期强制校验),同一作业只能选择一种表选择模式。
任务示例
以下示例均来自 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_list与table_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_list或query中的所有条目应用一个公共过滤条件。字符串必须以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 } }说明:提交时
JdbcSourceFactory的WhereConditionPrefixValidator会校验该字符串必须以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),仅供参考