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 引擎将 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 数据类型INTEGERDECIMAL(38,0)FLOATDECIMAL(38, 18)NUMBER(precision 9, scale 0)INTNUMBER(9 precision 18, scale 0)BIGINTNUMBER(18 precision, scale 0)DECIMAL(38, 0)NUMBER(scale ! 0)DECIMAL(38, 18)BINARY_DOUBLEDOUBLEBINARY_FLOAT、REALFLOATCHAR、NCHAR、VARCHAR、NVARCHAR2、VARCHAR2、LONG、ROWID、NCLOB、CLOB、XML、INTERVALSTRINGDATETIMESTAMPTIMESTAMP、TIMESTAMP WITH LOCAL TIME ZONETIMESTAMPBLOB、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_string为true时BLOB 会被映射为 STRING而不是 BYTES该开关当前仅对 Oracle 生效。字符串列长度上限CHAR 最大 2000 字节、VARCHAR2 最大 4000 字节、CLOB/NCLOB 最大 4GB-1、LONG 最大 2GB-1这些常量在转换器中均有定义。反转换reconvert说明转换器同时实现了reconvertSeaTunnel 类型 → 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开头。参数名类型必须默认值描述urlString是-JDBC 连接 URL如jdbc:oracle:thin:datasource01:1523:xedriverString是-JDBC 驱动类名Oracle 为oracle.jdbc.OracleDriverusernameString否-连接实例用户名passwordString否-连接实例密码queryString否-查询语句未配置table_path和table_list时必填connection_check_timeout_secInt否30等待连接验证数据库操作完成的时间秒partition_columnString否-用于并行分片的列名仅支持数值类型主键只能配置一列partition_lower_boundBigDecimal否-partition_column的最小值未设置时 SeaTunnel 自动查询数据库获取partition_upper_boundBigDecimal否-partition_column的最大值未设置时 SeaTunnel 自动查询数据库获取partition_numInt否job parallelism分割数量仅支持正整数默认等于任务并行度fetch_sizeInt否0查询行提取大小fetch size0 表示使用 JDBC 默认值。注意 Oracle 方言在 fetch_size 0 时实际使用默认 128见下文源码说明propertiesMap否-其他连接配置参数当 properties 与 URL 中参数相同时优先级由驱动实现决定Oracle 下 properties 优先于 URLuse_regexBoolean否false控制table_path是否按正则表达式匹配true时 table_path 视为正则否则视为精确路径table_pathString否-表的完整路径可替代query如test_schema.table1table_listArray否-要读取的表列表可替代table_path元素支持table_path与可选query字段where_conditionString否-所有表/查询的公共行过滤条件必须以where开头split.sizeInt否8096一个分片包含多少行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采样分片策略的采样率倒数1000表示 1/1000 采样split.allow-samplingBoolean否true是否允许对分布不均匀的分片键使用采样分片策略false时回退到迭代式不均匀分片use_select_countBoolean否false分片前是否使用select count(*)估算表行数skip_analyzeBoolean否false是否跳过分片前的表行数分析decimal_type_narrowingBoolean否trueDecimal 类型收窄为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时OracleSeaTunnelNUMBER(1, 0)BooleanNUMBER(6, 0)INTNUMBER(10, 0)BIGINTdecimal_type_narrowing false时OracleSeaTunnelNUMBER(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列进行分片。支持的分片数据类型StringNumberint、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_boundpartition_column的扫描最大值。未设置时 SeaTunnel 会查询数据库获取。partition_lower_boundpartition_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→ 回退到迭代式不均匀分片通过循环执行queryNextChunkMaxOracle 实现见 OracleDialect.queryNextChunkMax内部使用ORDER BY ... ASCROWNUM 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),仅供参考