十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

SeaTunnel Oracle Sink Connector 完整指南:JDBC 写入、数据类型映射与精确一次语义

SeaTunnel Oracle Sink Connector 完整指南:JDBC 写入、数据类型映射与精确一次语义 SeaTunnel Oracle Sink Connector 完整指南JDBC 写入、数据类型映射与精确一次语义【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 SeaTunnel 官方文档 Oracle JDBC Sink Connector并结合仓库内 JDBC 连接器源码seatunnel-connectors-v2/connector-jdbc撰写。Oracle Sink 是 SeaTunnel 基于 JDBC 协议实现的数据库写入端插件支持批处理与流处理两种模式、并发写入、基于 XA 事务的 exactly-once精确一次语义、CDC 事件写入与多表写入。读完本文你将掌握 Oracle Sink 的驱动部署、全部配置参数、数据类型映射规则以及 Simple、自动生成 SQL、精确一次、CDC、定时刷写、多表占位符等 6 种实战配置写法并理解其底层实现原理。支持引擎Oracle Sink 连接器可以运行在以下引擎之上SparkFlinkSeaTunnel Zeta其中 SeaTunnel Zeta 是项目自带的分布式引擎也是日常使用中最常用的运行方式。三种引擎下该 Sink 的配置写法完全一致差异仅体现在 JDBC 驱动 jar 的放置位置见下文使用依赖一节。功能特性Oracle Sink 具备以下核心能力均已在官方 Connector V2 Features 中声明支持exactly-once精确一次通过Xa transactionsXA 分布式事务保证。因此仅对支持Xa transactions的数据库生效启用方式为设置is_exactly_oncetrue且max_retries0。CDCChange Data Capture事件可接收上游 CDC 插件产出的变更事件INSERT / UPDATE / DELETE配合主键自动生成对应的写 SQL。多表写入multiple table write借助${schema_name}、${table_name}占位符可将数据按元数据路由到不同目标表。定时刷写timer flush通过batch_interval_ms参数在缓冲未满时按时间间隔触发刷写降低流式任务的写入延迟。使用依赖驱动 jar 的准备连接 Oracle 数据库需要 ojdbc8Oracle JDBC 驱动根据运行引擎不同放置位置有所区别对于 Spark / Flink 引擎需要将 JDBC 驱动 jar 放入${SEATUNNEL_HOME}/plugins/目录。对于 SeaTunnel Zeta 引擎需要将 JDBC 驱动 jar 放入${SEATUNNEL_HOME}/lib/目录。字符集与 i18n 支持Oracle 的国际化字符集支持由单独的orai18n.jar提供。为了让中文等多语言字符集能够正确读写需要将orai18n.jar一并复制到${SEATUNNEL_HOME}/lib/目录。提示官方文档示例给出的复制命令为cp ojdbc8-xxxxxx.jar $SEATUNNEL_HOME/lib/。不同 Oracle 版本对应的驱动类名可能不同请以实际下载的驱动版本为准。支持的数据库信息Datasource支持的版本DriverUrlMavenOracle不同依赖版本对应不同驱动类oracle.jdbc.OracleDriverjdbc:oracle:thin:datasource01:1523:xecom.oracle.database.jdbc:ojdbc8URL 使用 Oracle Thin 驱动的标准格式jdbc:oracle:thin:host:port:serviceName示例中的xe为 Oracle XE 实例的服务名/ SID。连接器的方言标识为Oracle对应源码OracleDialect.dialectName()返回DatabaseIdentifier.ORACLE见 OracleDialect.java。数据类型映射Oracle Sink 在将 SeaTunnel 数据类型写入 Oracle 时遵循 OracleTypeConverter.java 中定义的映射规则。官方文档给出的 Oracle → SeaTunnel 映射表如下Oracle Data TypeSeaTunnel Data TypeINTEGERINTFLOATDECIMAL(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、NVARCHAR2、VARCHAR2、LONG、ROWID、NCLOB、CLOBSTRINGDATEDATETIMESTAMP、TIMESTAMP WITH LOCAL TIME ZONETIMESTAMPBLOB、RAW、LONG RAW、BFILEBYTES结合源码可以补充几个值得注意的细节NUMBER 的精度分级映射器按精度 ≤ 9、9 精度 ≤ 18、精度 18三档分别映射为INT、BIGINT、DECIMAL(38, 0)而带小数部分scale ! 0的 NUMBER 统一映射为DECIMAL(38, 18)。这与OracleTypeConverter中MAX_PRECISION 38、DEFAULT_SCALE 18的常量定义一致。FLOAT 特殊处理在 OracleTypeMapper.java 中当 JDBC 元数据返回的NUMBER类型scale -127时会被识别为FLOAT同时NVARCHAR2、NCHAR等双字节字符类型的长度会按charToDoubleByteLength换算为字节长度避免长度判断偏差。时间精度TIMESTAMP默认精度为 6微秒最大支持 9纳秒DATE映射为 SeaTunnel 的DATE类型。配置参数详解以下参数表来自官方文档并结合 JdbcSinkOptions.java 中的 Option 定义默认值均与源码一致整理参数名类型必填默认值说明urlString是-JDBC 连接 URL示例jdbc:oracle:thin:datasource01:1523:xedriverString是-JDBC 驱动类名Oracle 使用oracle.jdbc.OracleDriverusernameString否-数据库用户名passwordString否-数据库密码queryString否-自定义写入 SQL如INSERT ...query优先级更高databaseString否-与table搭配自动生成 SQL与query互斥优先级更高tableString否-与database搭配自动生成 SQL与query互斥优先级更高primary_keysArray否-自动生成 SQL 时用于支持 insert / delete / update 操作的主键列表connection_check_timeout_secInt否30连接校验操作的等待超时时间秒max_retriesInt否0提交失败executeBatch时的重试次数batch_sizeInt否1000批写模式下缓冲记录数达到batch_size即刷写若batch_interval_ms 0时间也可触发刷写batch_interval_msLong否0写入触发的刷写间隔毫秒。0表示关闭基于时间的刷写大于 0 时每条记录写入都会检查耗时达到间隔则同步刷写is_exactly_onceBoolean否false是否启用精确一次语义使用 XA 事务开启时需设置xa_data_source_class_namegenerate_sink_sqlBoolean否false根据目标表自动生成 SQL 语句xa_data_source_class_nameString否-数据库驱动的 XA 数据源类名Oracle 为oracle.jdbc.xa.client.OracleXADataSourcemax_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务开启后的超时时间默认 -1 表示永不超时注意设置超时可能影响精确一次语义auto_commitBoolean否true默认开启自动提交事务propertiesMap否-额外的连接配置参数当 properties 与 URL 含相同参数时优先级取决于驱动具体实现如 MySQL 中 properties 优先于 URLcommon-options-否-Sink 插件通用参数详见 Sink Common Optionsschema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST同步任务开启前对目标端表结构存在与否采取的建表策略data_save_modeEnum否APPEND_DATA同步任务开启前对目标端已有数据采取的处理策略custom_sqlString否-当data_save_mode选择CUSTOM_PROCESSING时填写通常是同步任务开始前执行的 SQLenable_upsertBoolean否true存在主键时启用 upsert若任务没有主键重复数据设为false可加速导入multi_table_sink_replicaInt否1多表写入时使用的 sink writer 副本数参数要点补充源码视角query与database/table的互斥关系在 JdbcSinkConfig.java 中JdbcSinkConfig.of()会分别解析QUERY、DATABASE、TABLE三个选项query提供了最直接的 SQL 控制而databasetable则用于自动生成 SQL。两者只能选其一。schema_save_mode/data_save_mode两者在 JdbcSinkOptions 中定义为枚举默认分别为CREATE_SCHEMA_WHEN_NOT_EXIST表不存在时自动建表与APPEND_DATA追加写入。搭配custom_sql可实现在任务启动前执行自定义初始化 SQL。enable_upsert仅当配置了primary_keys时生效。Oracle 方言下的 upsert 最终会生成MERGE INTO语句详见下文Upsert 实现。其他 JDBC 通用选项连接器还支持field_ide字段大小写转换、use_copy_statement、create_index、oracle_insert_mode等选项见 JdbcSinkOptions.java。其中oracle_insert_mode可设置为CONVENTIONAL普通 INSERT或APPEND_VALUES为仅插入场景添加 OracleAPPEND_VALUEShint提升插入性能。Tips若未设置partition_column分区列任务将以单并发运行设置partition_column后将按任务的并发度并行执行。任务示例以下示例均以 FakeSource 生成数据、Oracle Sink 写入为例。运行前需要在 Oracle 中预先创建数据库test和表test_tableSimple 示例尚未安装部署 SeaTunnel 时先按 Install SeaTunnel 完成部署再按 Quick Start With SeaTunnel Engine 学习如何提交任务。1. Simple基础写入本例定义了一个 SeaTunnel 批同步任务FakeSource 自动生成 16 行数据row.num16每行包含namestring和ageint两个字段最终写入 Oracle 的test_table表16 行。使用query指定插入 SQL# 定义运行时环境 env { parallelism 1 job.mode BATCH } source { FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age int } } } } transform { } sink { jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver username root password 123456 query INSERT INTO TEST.TEST_TABLE(NAME,AGE) VALUES(?,?) } }要点query中的?为占位符按字段顺序绑定上游数据表名使用TEST.TEST_TABLEschema.table写法与 Oracle 的命名习惯一致。2. Generate Sink SQL自动生成 SQL如果不想手写复杂 SQL可以配置database和table让连接器自动生成 INSERT 语句sink { Jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver username root password 123456 generate_sink_sql true database XE table TEST.TEST_TABLE } }3. Exactly-once精确一次针对需要精确写入的场景配合 XA 事务实现。注意两个关键点is_exactly_once true与max_retries 0并指定 Oracle 的 XA 数据源类sink { jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver max_retries 0 username root password 123456 query INSERT INTO TEST.TEST_TABLE(NAME,AGE) VALUES(?,?) is_exactly_once true xa_data_source_class_name oracle.jdbc.xa.client.OracleXADataSource } }原理说明精确一次要求max_retries 0是因为一旦启用 XA 事务后任何写入失败都必须在事务层面回滚而不是靠重试单个 batch 来兜底否则可能破坏事务的一致性边界。4. CDCChange Data Capture事件连接器支持接收上游 CDC 插件产生的变更数据INSERT / UPDATE / DELETE。此时需要同时配置database、table和primary_keys并配合自动生成 SQLsink { jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver username root password 123456 generate_sink_sql true # database 和 table 必须同时配置 database XE table TEST.TEST_TABLE primary_keys [ID] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }5. Timer Flush Batch 组合定时刷写 批量对于流式任务如果不想等缓冲区填满才刷写可同时设置batch_size和batch_interval_ms。需要理解其刷写机制是**写入触发式write-triggered**的每条写入都会检查缓冲行数与已耗时任一阈值达到即同步刷写。没有后台调度器因此在空闲时段无新数据到达缓冲行会一直保留直到下一条数据到达或 checkpoint 完成——batch_interval_ms本身并不保证严格的墙钟延迟上界。sink { Jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver username root password 123456 generate_sink_sql true database XE table TEST.TEST_TABLE primary_keys [ID] batch_size 2000 batch_interval_ms 5000 } }6. 多表写入Multi-Table Write With Placeholder在table中使用${schema_name}和${table_name}占位符可根据上游行携带的元数据将每条记录路由到匹配的目标表。通过multi_table_sink_replica控制多表写入时使用的 writer 副本数sink { Jdbc { url jdbc:oracle:thin:datasource01:1523:xe driver oracle.jdbc.OracleDriver username root password 123456 generate_sink_sql true database XE table ${schema_name}.${table_name}_SINK primary_keys [ID] multi_table_sink_replica 2 } }注tablePrefix/tableSuffix参数已废弃官方建议统一改用table占位符写法见 JdbcSinkOptions.java 中的Deprecated标注。底层实现原理源码级解读Upsert 实现Oracle 方言生成 MERGE INTO当配置primary_keys且启用enable_upsert时Oracle 方言在 OracleDialect.getUpsertStatement() 中生成标准的 OracleMERGE INTO语句USING (SELECT :field field, ... FROM DUAL) SOURCE利用 Oracle 的DUAL虚拟表构造单行数据源ON (TARGET.pk SOURCE.pk)以主键作为匹配条件WHEN MATCHED THEN UPDATE SET ...主键命中时更新非主键字段WHEN NOT MATCHED THEN INSERT (...) VALUES (...):未命中时执行插入。这意味着 Oracle Sink 的 upsert 完全借助 Oracle 原生MERGE能力实现无需依赖额外的方言插件。精确一次XA 事务与两阶段提交is_exactly_once开启后写入路径切换到 JdbcExactlyOnceSinkWriter.java。其核心机制包括通过xa_data_source_class_name指定的OracleXADataSource创建 XA 连接事务以 XID事务分支标识为单位管理写入阶段将数据缓存在 XA 事务内并prepare提交阶段由 JdbcSinkAggregatedCommitter.java 统一协调提交源码中针对EmptyXaTransactionException空 XA 事务做了特殊处理事务内没有任何写入时不会生成提交记录同时实现了失败时的rollback与 prepared 事务恢复逻辑rollbackPrepareXidOrThrow/tryRecoverPreparedTransactionsAfterRollbackFailure用于清理异常残留的预提交事务。因此精确一次是写入端的保证其有效性取决于目标数据库对 XA 事务的支持程度Oracle 原生支持 XA。类型映射与数据转换写入方向SeaTunnel → Oracle由 OracleJdbcRowConverter.java 负责将 SeaTunnel 类型转换为 JDBC 参数读取/建表方向由OracleTypeConverter负责。建表场景下SeaTunnel 的 DECIMAL 精度上限为 38、默认小数位 18与 Oracle NUMBER 的最大精度一致可以无损对齐。连接与连接池写入端通过 ConnectionPoolManager.java 管理 JDBC 连接复用connection_check_timeout_sec用于控制连接校验操作的超时批量写入路径 JdbcSinkWriter.java 负责executeBatch提交与max_retries重试。总结Oracle Sink 连接器是 SeaTunnel 对接 Oracle 数据库的标准写入通道具备三个层次的用法基础写入query自定义 SQL 或databasetable自动生成 SQL增强写入primary_keysenable_upsert实现 upsertbatch_sizebatch_interval_ms控制刷写节奏schema_save_mode/data_save_mode管理目标表与存量数据高级语义is_exactly_once XA 事务实现精确一次CDC 事件写入以及${schema_name}/${table_name}占位符的多表路由。在动手配置前请务必确认驱动 jar 已按引擎类型放入正确的目录Zeta 引擎为${SEATUNNEL_HOME}/lib/并视字符集需求补充orai18n.jar。更详细的插件通用参数可参考 Sink Common Options连接器的完整变更记录见 connector-jdbc Changelog。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表