SeaTunnel OceanBase JDBC Sink 连接器实战指南:配置、类型映射与精确一次写入
SeaTunnel OceanBase JDBC Sink 连接器实战指南配置、类型映射与精确一次写入【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本文基于 docs/zh/connectors/sink/OceanBase.md 展开围绕 SeaTunnel 官方 JDBC Sink 对 OceanBase 的接入能力系统讲解其支持引擎、关键特性、数据类型映射、Sink 选项、CDC 多表写入与 XA 精确一次语义并结合 connector-jdbc 模块源码佐证底层实现原理。读完本文你将能够独立编写从 FakeSource 或 CDC 上游到 OceanBase 的批/流同步任务正确配置 MySQL/Oracle 兼容模式并在需要时启用基于 XA 事务的精确一次写入。一、连接器概述与支持引擎OceanBase Sink 是 SeaTunnel JDBC 连接器家族中的一个方言实现dialect通过 JDBC 协议向 OceanBase 数据库写入数据。它支持批处理任务和流处理任务两类作业形态能够处理CDC 行类型INSERT / UPDATE / DELETE 事件、自动生成 SQL、保存模式schema_save_mode / data_save_mode、多表写入并且在配置 XA 事务后支持精确一次exactly-once语义。该 Sink 可运行于以下引擎SparkFlinkSeaTunnel Zeta从源码结构看OceanBase 的接入位于 seatunnel-connectors-v2/connector-jdbc 模块内部并非独立连接器其核心类集中在 internal/dialect/oceanbase 与 catalog/oceanbase 两个包下。二、关键特性连接器官方文档列出的关键能力如下精确一次Exactly-Once依赖 XA 事务实现需要配置is_exactly_once true、max_retries 0并填写 OceanBase JDBC 驱动提供的有效 XA 数据源类名。CDC 支持可消费上游 CDC 变更数据INSERT / UPDATE / DELETE。多表写入支持在table中通过${schema_name}、${table_name}等占位符实现多表分发。定时刷新通过batch_interval_ms实现基于写入触发的定时刷新。精确一次的硬性前置条件启用时需同时满足is_exactly_once true、max_retries 0并填写 OceanBase JDBC 驱动提供的有效 XA 数据源类名。源码 JdbcExactlyOnceSinkWriter 中通过checkArgument强制校验maxRetries 0否则直接抛出异常——因为非零重试与 XA 提交组合可能导致数据重复。更详细的特性说明可参考 Connector V2 特性。三、支持的数据源信息与依赖部署官方给出的连接信息如下数据源支持版本DriverUrlMavenOceanBase所有 OceanBase 服务版本com.oceanbase.jdbc.Driverjdbc:oceanbase://localhost:2883/testoceanbase-client数据库相关依赖请下载 Maven 中央仓库中与com.oceanbase/oceanbase-client对应的驱动 JAR并将其复制到$SEATUNNEL_HOME/plugins/jdbc/lib/目录# 示例将驱动复制到 JDBC 插件依赖目录 cp oceanbase-client-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/从源码 OceanBaseDialectFactory 可以看出连接器的方言识别规则是URL 以jdbc:oceanbase:前缀开头即被识别为 OceanBase随后依据compatible_mode参数决定使用哪套方言实现compatible_mode oracle→ 复用OracleDialect其他取值如mysql→ 使用OceanBaseMysqlDialect这也解释了为什么compatible_mode是必填项缺少它时方言工厂的create()会直接抛出UnsupportedOperationException(Cant create JdbcDialect without compatible mode for OceanBase)。四、数据类型映射OceanBase 同时兼容 MySQL 与 Oracle 两套 SQL 方言因此类型映射也分为两张表。映射逻辑的底层实现在 OceanBaseMySqlTypeConverterMySQL 模式以及 Oracle 对应的转换器中。4.1 MySQL 模式MySQL 数据类型SeaTunnel 数据类型BIT(1)INT UNSIGNEDBOOLEANTINYINTTINYINT UNSIGNEDSMALLINTSMALLINT UNSIGNEDMEDIUMINTMEDIUMINT UNSIGNEDINTINTEGERYEARINTINT UNSIGNEDINTEGER UNSIGNEDBIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20,0)DECIMAL(x,y)指定列大小 38DECIMAL(x,y)DECIMAL(x,y)指定列大小 38DECIMAL(38,18)DECIMAL UNSIGNEDDECIMAL(精度1, 小数位)FLOATFLOAT UNSIGNEDFLOATDOUBLEDOUBLE UNSIGNEDDOUBLECHARVARCHARTINYTEXTMEDIUMTEXTTEXTLONGTEXTJSONSTRINGDATEDATETIMETIMEDATETIMETIMESTAMPTIMESTAMPTINYBLOBMEDIUMBLOBBLOBLONGBLOBBINARYVARBINARBIT(n)BYTESGEOMETRYUNKNOWN暂不支持从 OceanBaseMySqlTypeConverter 的convert方法可以看到几个值得注意的实现细节BIT 类型的双态处理BIT(1)或未指定长度时映射为 BOOLEANBIT(n)n1映射为 BYTES且长度按n/8向上取整。BIGINT UNSIGNED → DECIMAL(20,0)因为无符号 BIGINT 超出 SeaTunnel BIGINT 的表示范围源码中通过new DecimalType(20, 0)处理。DECIMAL 超精度截断当精度超过默认值 38 时映射为DECIMAL(38,18)并输出will probably cause value overflow警告日志DECIMAL UNSIGNED则按“精度 1”扩容以容纳符号位。TEXT 系列长度上限TINYTEXT/TEXT/MEDIUMTEXT/LONGTEXT 分别按2^8-1、2^16-1、2^24-1、2^32-1截断。4.2 Oracle 模式Oracle 数据类型SeaTunnel 数据类型Number(p), p 9INTNumber(p), p 18BIGINTNumber(p), p 18DECIMAL(38,18)REALBINARY_FLOATFLOATBINARY_DOUBLEDOUBLECHARNCHARNVARCHAR2NCLOBCLOBROWIDSTRINGDATEDATETIMESTAMPTIMESTAMP WITH LOCAL TIME ZONETIMESTAMPBLOBRAWLONG RAWBFILEBYTESUNKNOWN暂不支持提示在 Oracle 模式下 SeaTunnel 直接复用 OracleDialect 及其行转换器因此类型映射、标识符引用如双引号与 SQL 生成均遵循 Oracle 语法风格。五、Sink 选项详解以下为 OceanBase Sink 的全部配置项。其中多数参数由 JdbcSinkOptions 定义并经 JdbcSinkConfig.of() 解析进运行时配置对象。参数名类型是否必填默认值描述urlString是-JDBC 连接 URL示例jdbc:oceanbase://localhost:2883/testdriverString是-连接远程数据源的 JDBC 类名应为com.oceanbase.jdbc.DriverusernameString否-连接实例用户名passwordString否-连接实例密码queryString否-自定义写入 SQL如INSERT ...。当generate_sink_sql false时必须配置query自定义 query 模式下不会执行保存模式相关配置compatible_modeString是-OceanBase 兼容模式取值mysql或oracledatabaseString否-generate_sink_sql true时使用的数据库此时必须配置tableString否-generate_sink_sql true时使用的目标表多表写入支持${schema_name}、${table_name}等占位符primary_keysArray否-自动生成 SQL 时支持 insert/delete/update 等操作connection_check_timeout_secInt否30等待用于验证连接的数据库操作完成的时间秒max_retriesInt否0提交失败的重试次数executeBatchbatch_sizeInt否1000批量写入缓冲上限达到batch_size条记录时刷新到 OceanBase若batch_interval_ms 0达到间隔也会触发刷新batch_interval_msLong否0写入触发的定时刷新间隔毫秒。0表示关闭大于 0 时写入器在每条记录写入时检查间隔达到间隔后同步刷新is_exactly_onceBoolean否false是否通过 XA 事务启用精确一次。启用后需配置xa_data_source_class_name并保持max_retries 0generate_sink_sqlBoolean否false根据目标数据库表自动生成 SQL 语句xa_data_source_class_nameString否-OceanBase JDBC 驱动提供的 XA 数据源类名is_exactly_once true时必填max_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务打开后的超时时间-1 表示永不超时设置超时可能影响精确一次语义auto_commitBoolean否true默认启用自动事务提交field_ideString否-控制字段名大小写转换可选ORIGINAL、UPPERCASE、LOWERCASEpropertiesMap否-其他连接配置参数当属性和 URL 携带相同参数时优先级由驱动实现决定如 MySQL 中属性优先于 URLschema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST同步任务启动前对目标表结构的处理方式data_save_modeEnum否APPEND_DATA同步任务启动前对目标表已有数据的处理方式custom_sqlString否-data_save_mode CUSTOM_PROCESSING时在同步前执行的 SQL自定义 query 模式下不执行common-options-否-Sink 插件通用参数详见 Sink 通用选项enable_upsertBoolean否true通过 primary_keys 存在启用 upsert若任务无主键重复数据设为false可加快数据导入is_primary_key_updatedBoolean否true自动生成更新语句时是否把主键字段放入更新字段中support_upsert_by_insert_onlyBoolean否false兼容方言下是否通过仅 INSERT 语句实现 upsertmulti_table_sink_replicaInt否1多表写入时使用的 Sink Writer 副本数量5.1 关键参数源码解读compatible_mode决定方言选择如第三节所述OceanBaseDialectFactory.create(compatibleMode, fieldIde) 中oracle返回OracleDialect其余返回OceanBaseMysqlDialect。MySQL 模式下标识符用反引号包裹见 OceanBaseMysqlDialect.quoteIdentifier并且会应用field_ide对字段名做大小写转换。generate_sink_sql与 upsert 自动生成当启用自动生成 SQL 时OceanBaseMysqlDialect.getUpsertStatement() 会基于INSERT INTO ... ON DUPLICATE KEY UPDATE colVALUES(col), ...生成 upsert 语句——这正是 MySQL 兼容模式下enable_upsert/primary_keys的底层支撑。默认 JDBC 连接参数MySQL 模式下的方言默认会向连接附加rewriteBatchedStatementstrue与allowMultiQueriestrue见 OceanBaseMysqlDialect.defaultParameter()前者对executeBatch批量写入性能有明显提升后者则服务于多语句执行场景。field_ide字段命名策略枚举定义在 FieldIdeEnum支持ORIGINAL保持原样、UPPERCASE转大写、LOWERCASE转小写在自动生成建表/写入 SQL 时统一生效。batch_interval_ms的写入触发机制从 JdbcSinkOptions.BATCH_INTERVAL_MS 的描述可以看出该参数是**写入触发write-triggered**的每条记录进入写入路径时都会检查距上次刷新的耗时达到间隔即同步刷新并没有后台调度线程。因此在空闲无新记录时段缓冲行会一直保留到下一条记录到达或下一个 checkpoint 完成——batch_interval_ms自身并不能保证低吞吐流上的精确 wall-clock 时延边界不应把它当作严格的实时定时器。5.2 配置提示OceanBase MySQL 模式请配置compatible_mode mysqlOracle 模式请配置compatible_mode oracle。想自己写完整写入 SQL 时使用query想让 SeaTunnel 自动生成 INSERT/UPSERT SQL 并执行保存模式时使用generate_sink_sql true并配置database和table。消费 CDC 数据时建议使用自动生成 SQL 并配置primary_keys否则 UPDATE 和 DELETE 事件无法安全映射。六、任务示例以下示例均以 config 目录下的 HOCON 配置格式编写。运行前需先按 安装 SeaTunnel 完成部署再参考 快速启动 SeaTunnel 引擎 运行作业。6.1 简单示例自定义 SQL此示例通过 FakeSource 自动生成 16 行数据row.num 16字段为namestring和ageint并写入 JDBC Sink最终目标表test_table中也将有 16 行数据。运行前需在数据库中创建test库和test_table表。# 定义运行环境 env { parallelism 1 job.mode BATCH } source { # 这是一个示例源插件仅用于测试和演示功能 FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age int } } } } transform { # 如需了解 transform 插件完整列表请参考 transforms 文档 } sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username root password 123456 compatible_mode mysql query insert into test_table(name,age) values(?,?) } }6.2 自动生成 Sink SQL无需手写复杂 SQL仅配置数据库名与表名即可自动生成插入语句sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username root password 123456 compatible_mode mysql # 根据数据库表名自动生成 sql 语句 generate_sink_sql true database test table test_table } }6.3 自动生成 SQL 并设置保存模式配合primary_keys与保存模式在启动同步任务前对目标表结构与存量数据做统一处理sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username roottest password compatible_mode mysql generate_sink_sql true database test table sink_table primary_keys [id] schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }保存模式的底层执行由 JdbcSaveModeHandler 与 catalog/oceanbase 下的 Catalog 实现协作完成。例如 MySQL 模式通过INFORMATION_SCHEMA.COLUMNS读取表结构见 OceanBaseMySqlCatalog.SELECT_COLUMNS_SQL_TEMPLATE并由 OceanBaseMysqlCreateTableSqlBuilder 生成建表语句。6.4 CDCChange Data Capture数据变更事件OceanBase Sink 支持 CDC 变更数据。此场景下需要同时配置database、table与primary_keys以便将 INSERT / UPDATE / DELETE 事件安全映射为对应的写操作sink { jdbc { url jdbc:oceanbase://localhost:3306/test driver com.oceanbase.jdbc.Driver username root password 123456 compatible_mode mysql generate_sink_sql true # 您需要同时配置数据库和表 database test table sink_table primary_keys [id,name] } }说明文档示例中该配置的 URL 端口为 3306实际部署时应替换为 OceanBase 实例的真实连接端口默认 2883。6.5 Oracle 兼容模式当目标 OceanBase 租户运行在 Oracle 兼容模式时compatible_mode必须设为oracle并使用 Oracle 风格的 SQL 语法sink { jdbc { url jdbc:oceanbase://localhost:2883/TESTUSER driver com.oceanbase.jdbc.Driver username TESTUSERtest password compatible_mode oracle query INSERT INTO SINK_TABLE (ID, NAME, CREATE_TIME) VALUES (?, ?, ?) } }6.6 多表写入当上游数据携带表身份信息时可在table中使用占位符实现多表分发。${table_name}会被替换为上游记录对应的实际表名sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username roottest password compatible_mode mysql generate_sink_sql true database test table ${table_name}_sink primary_keys [id] multi_table_sink_replica 2 } }multi_table_sink_replica用于控制多表写入时使用的 Sink Writer 副本数量多表场景下的资源管理由 JdbcMultiTableResourceManager 承担。6.7 基于 XA 事务的精确一次写入启用 XA 精确一次语义需要三要素齐备is_exactly_once true、提供 OceanBase JDBC 驱动中的xa_data_source_class_name、max_retries 0。写入器会把每个 checkpoint 批次包装在 XA 事务里要么与源 checkpoint 一起提交要么失败回滚env { parallelism 2 job.mode STREAMING checkpoint.interval 10000 } sink { Jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username root password 123456 compatible_mode mysql generate_sink_sql true database test table sink_table primary_keys [id] is_exactly_once true xa_data_source_class_name com.oceanbase.jdbc.OceanBaseXADataSource max_retries 0 batch_size 1000 } }从 JdbcExactlyOnceSinkWriter 的实现可以看到完整的 XA 生命周期开启事务beginTx()通过xidGenerator.generateXid()生成带语义的全局事务 ID并调用xaFacade.start(currentXid)启动 XA 事务批量写入write()将每条记录克隆后经outputFormat.writeRecord()写入数据在prepareCurrentTx()阶段通过outputFormat.flush()落盘预提交与快照prepareCommit()调用xaFacade.endAndPrepare(currentXid)完成 end prepare 两阶段snapshotState()仅保存prepareXid若事务为空EmptyXaTransactionException会跳过预提交并返回空状态恢复与回滚writer 打开时执行xaGroupOps.recoverAndRollback()回滚残留的未决事务close()或abortPrepare()时也会尽力回滚当前与已预提交的事务配合 JdbcSinkAggregatedCommitter 完成最终两阶段提交。6.8 批量 定时刷新组合流式作业可以同时设置batch_size与batch_interval_ms基于距离上次刷新的耗时来刷新缓冲行sink { Jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver username root password 123456 compatible_mode mysql generate_sink_sql true database test table sink_table primary_keys [id] batch_size 2000 batch_interval_ms 5000 } }再次强调该模式的语义边界刷新是写入触发的每条记录进入写入路径时都会检查耗时达到间隔才同步刷新不存在后台调度线程。因此在空闲时段缓冲行会一直保留到下一条记录到达或下一个 checkpoint 完成。配合batch_size使用可以兼顾吞吐与单条记录时延但请不要把它当作严格的实时定时器。七、变更日志OceanBase Sink 随 JDBC 连接器一同演进具体变更历史请查阅 JDBC 连接器变更日志。八、相关源码索引如需深入阅读 OceanBase 接入的实现细节可重点关注以下文件方言工厂与识别OceanBaseDialectFactory.javaMySQL 模式方言upsert 生成、默认参数、标识符引用OceanBaseMysqlDialect.javaMySQL 模式类型转换OceanBaseMySqlTypeConverter.java行数据转换OceanBaseMysqlJdbcRowConverter.javaCatalog 与建表语句OceanBaseMySqlCatalog.java、OceanBaseOracleCatalog.java精确一次写入器JdbcExactlyOnceSinkWriter.java参数定义JdbcSinkOptions.java此外OceanBaseMySqlTypeConverterTest.java 与 OceanBaseMysqlDialectTest.java 提供了类型映射与 SQL 生成的单元测试可作为验证与二次开发的参考。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考