
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载本文围绕 Apache SeaTunnel 的 JDBC OceanBase Sink 连接器展开系统讲解其引擎支持范围、双兼容模式MySQL/Oracle下的数据类型映射、全部 Sink 配置参数以及基于 FakeSource 的同步、自动生成 SQL、CDC 变更事件三类任务配置示例。读完本文你将掌握如何在 SeaTunnel 中通过 JDBC 将数据稳定写入 OceanBase并理解底层方言分发、参数默认值与事务语义的实现原理。概述OceanBase Sink 连接器基于 SeaTunnel 的 JDBC 连接器实现源码位于 connector-jdbc 模块负责通过 JDBC 将上游数据写入 OceanBase 数据库。它支持**批处理Batch与流处理Streaming**两种任务模式支持并发写入并支持 exactly-once 精确一次语义同时可对接 CDCChange Data Capture变更数据实现 INSERT / DELETE / UPDATE 等行级操作的落库。特性支持情况exactly-once支持通过事务/两阶段提交实现cdc支持配合primary_keys自动生成增删改 SQL支持的引擎SparkFlinkSeaTunnel ZetaSeaTunnel 自研引擎支持的数据库信息与依赖数据源支持版本DriverURL 示例Maven 坐标OceanBase全部 OceanBase 服务端版本com.oceanbase.jdbc.Driverjdbc:oceanbase://localhost:2883/testcom.oceanbase:oceanbase-client数据库驱动依赖安装需要下载上述 Maven 坐标对应的驱动 jar并复制到$SEATUNNEL_HOME/plugins/jdbc/lib/工作目录例如cp oceanbase-client-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/在 SeaTunnel 的插件目录中JDBC 连接器通过统一的插件加载机制在 config/plugin_config 中注册plugins/jdbc/lib/是该连接器的驱动存放目录详见 plugins/README.md。兼容模式与底层方言分发OceanBase 同时兼容 MySQL 与 Oracle 两种租户模式因此 Sink 连接器在建立连接和执行 SQL 时必须明确compatible_modemysql或oracle。从源码看这一分发逻辑位于 OceanBaseDialectFactory.javaacceptsURL通过url.startsWith(jdbc:oceanbase:)识别 OceanBase JDBC URLcreate()无参会直接抛出UnsupportedOperationException提示无法在未指定兼容模式的情况下创建 OceanBase 方言——这说明compatible_mode是必填项create(compatibleMode, fieldIde)在兼容模式为oracle时返回OracleDialect否则返回MysqlDialect即底层 SQL 方言生成分别复用 Oracle 与 MySQL 的实现。对应地SeaTunnel 也提供了 OceanBase 的 Catalog 工厂 OceanBaseCatalogFactory.java在oracle模式创建OceanBaseOracleCatalog否则创建OceanBaseMySqlCatalog用于元数据schema/表结构的读取与校验。数据类型映射数据类型映射由OceanBase 兼容模式下的原生类型 → SeaTunnel 内部类型决定写库时 SeaTunnel 会将上游数据类型按映射关系转换为目标列类型。MySQL 模式MySQL 数据类型SeaTunnel 数据类型BIT(1)、INT UNSIGNEDBOOLEANTINYINT、TINYINT UNSIGNED、SMALLINT、SMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINTINT UNSIGNED、INTEGER UNSIGNED、BIGINTBIGINTBIGINT UNSIGNEDDECIMAL(20,0)DECIMAL(x,y)精度 38DECIMAL(x,y)DECIMAL(x,y)精度 38DECIMAL(38,18)DECIMAL UNSIGNEDDECIMAL(列精度1, 小数位数)FLOAT、FLOAT UNSIGNEDFLOATDOUBLE、DOUBLE UNSIGNEDDOUBLECHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、JSONSTRINGDATEDATETIMETIMEDATETIME、TIMESTAMPTIMESTAMPTINYBLOB、MEDIUMBLOB、BLOB、LONGBLOB、BINARY、VARBINARY、BIT(n)BYTESGEOMETRY、UNKNOWN暂不支持Oracle 模式Oracle 数据类型SeaTunnel 数据类型Number(p)p 9INTNumber(p)p 18BIGINTNumber(p)p 18DECIMAL(38,18)REAL、BINARY_FLOATFLOATBINARY_DOUBLEDOUBLECHAR、NCHAR、NVARCHAR2、NCLOB、CLOB、ROWIDSTRINGDATEDATETIMESTAMP、TIMESTAMP WITH LOCAL TIME ZONETIMESTAMPBLOB、RAW、LONG RAW、BFILEBYTESUNKNOWN暂不支持Sink 配置参数详解下表完整列出 OceanBase JDBC Sink 支持的全部配置项其中默认值与参数类型与 JdbcOptions.java 中的定义一一对应。参数名类型必填默认值说明urlString是-JDBC 连接 URL例如jdbc:oceanbase://localhost:2883/testdriverString是-连接远端数据源的 JDBC 驱动类名应为com.oceanbase.jdbc.DriveruserString否-连接实例的用户名passwordString否-连接实例的密码queryString否-使用该 SQL 将上游数据写入数据库例如INSERT ...query优先级更高compatible_modeString是-OceanBase 兼容模式可选mysql或oracledatabaseString否-配合table自动生成 SQL 并接收上游数据写入与query互斥且优先级更高tableString否-配合database自动生成 SQL与query互斥且优先级更高primary_keysArray否-自动生成 SQL 时用于支持insert、delete、update等操作support_upsert_by_query_primary_key_existBoolean否false当数据库不支持 upsert 语法时通过查询主键是否存在来选择 INSERT 或 UPDATE SQL 处理更新事件INSERT、UPDATE_AFTER。注意该方式性能较低connection_check_timeout_secInt否30等待用于校验连接的数据库操作完成的秒数max_retriesInt否0提交失败executeBatch时的重试次数batch_sizeInt否1000批量写入时缓冲记录数达到batch_size或时间达到checkpoint.interval时刷入数据库generate_sink_sqlBoolean否false根据目标数据库表自动生成 SQL 语句max_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务开启后的超时时间默认 -1 表示永不超时。注意设置超时可能影响 exactly-once 语义auto_commitBoolean否true默认开启自动事务提交propertiesMap否-额外的连接配置参数当 properties 与 URL 中存在同名参数时优先级由驱动实现决定例如 MySQL 中 properties 优先于 URLcommon-options-否-Sink 插件公共参数详见 Sink Common Optionsenable_upsertBoolean否true依据主键是否存在启用 upsert若任务数据无主键重复设置为false可加速数据导入参数规则的源码印证从 JdbcSinkFactory.java 的optionRule()可以看到参数的依赖关系必填项为url、driver、schema_save_mode、data_save_modeis_exactly_oncetrue时才会启用xa_data_source_class_name、max_commit_attempts、transaction_timeout_secis_exactly_oncefalse时启用max_retriesgenerate_sink_sqltrue时必须配置databasegenerate_sink_sqlfalse时配置query。Tips若不设置分区列任务将以单并发运行若设置了分区列则按任务的并发度并行执行该配置对应 JDBC Source 侧的partition_column系列参数请参考 JdbcOptions.java 中的PARTITION_COLUMN、PARTITION_UPPER_BOUND、PARTITION_LOWER_BOUND、PARTITION_NUM。任务配置示例以下示例均使用 SeaTunnel 配置文件格式HOCON。运行前需完成 SeaTunnel 的安装部署参见 Install SeaTunnel首次运行任务请参考 Quick Start With SeaTunnel Engine。示例一简单写入FakeSource → JDBC Sink该示例通过 FakeSource 自动生成 16 行数据row.num16每行包含namestring和ageint两个字段最终写入目标表test_table运行前需在数据库中预先创建test库和test_table表任务结束后表中应有 16 行数据# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { # 仅用于测试和演示的示例 Source 插件 FakeSource { parallelism 1 result_table_name fake row.num 16 schema { fields { name string age int } } } } transform { # 如需查看完整的 Transform 插件列表请查阅仓库 docs/en/transform-v2 目录 } sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver user root password 123456 compatible_mode mysql query insert into test_table(name,age) values(?,?) } }要点说明query使用占位符?接收上游字段字段顺序与占位符一一对应compatible_mode mysql表示目标为 MySQL 兼容模式的 OceanBase 租户。示例二自动生成 Sink SQL当不想手写复杂 SQL 时配置generate_sink_sql true并给出database与tableSeaTunnel 会根据表结构自动生成写入语句sink { jdbc { url jdbc:oceanbase://localhost:2883/test driver com.oceanbase.jdbc.Driver user root password 123456 compatible_mode mysql # 根据数据库表名自动生成 SQL 语句 generate_sink_sql true database test table test_table } }示例三CDCChange Data Capture事件写入SeaTunnel 也支持 CDC 变更数据写入 OceanBase。此时需要同时配置database、table与primary_keysSeaTunnel 会根据主键自动生成 INSERT / DELETE / UPDATE 语句来应用上游的增删改事件sink { jdbc { url jdbc:oceanbase://localhost:3306/test driver com.oceanbase.jdbc.Driver user root password 123456 compatible_mode mysql generate_sink_sql true # 需要同时配置 database 和 table database test table sink_table primary_keys [id,name] } }关于 upsert 与性能的调优建议默认enable_upsert true连接器会依据主键判断记录是否存在以决定 INSERT 还是 UPDATE如果数据本身无主键重复将enable_upsert设为false可以显著提升导入速度当数据库不支持原生 upsert 语法时可开启support_upsert_by_query_primary_key_exist通过先查后写实现等价语义但该方案性能较低仅在必要时启用批量写入大小由batch_size默认 1000控制同时受checkpoint.interval触发刷写事务提交失败重试次数由max_commit_attempts默认 3控制。小结OceanBase JDBC Sink 是 SeaTunnel 连接生态中面向 OceanBase 的标准落库通道通过compatible_mode无缝适配 MySQL / Oracle 两种租户模式底层由 OceanBaseDialectFactory 分发到MysqlDialect/OracleDialect支持手工 SQL、自动生成 SQL 与 CDC 事件三类写入方式并可通过enable_upsert、batch_size、max_commit_attempts等参数在 exactly-once 语义与吞吐性能之间取得平衡。更多公共参数与相邻能力可进一步参考 Sink Common Options 与 connector-v2-features。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel JDBC Snowflake Sink 连接器配置、CDC 写入与数据类型映射实战指南SeaTunnel JDBC Snowflake Sink 连接器配置、CDC 写入与数据类型映射实战指南 本文面向使用 Apache SeaTunnel h数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel OceanBase JDBC Sink 连接器实战指南配置、类型映射与精确一次写入SeaTunnel OceanBase JDBC Sink 连接器实战指南配置、类型映射与精确一次写入 本文基于 docs/zh/connectors/sin数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入 本文围绕 Kingbase Sink 连接器文档 https数据集成ETL大数据批处理流处理变更数据捕获上一篇使用 AWS CLI 的 codeguru-reviewer put-recommendation-feedback 提交代码评审建议反馈下一篇显卡驱动冲突的终极解决方案Display Driver Uninstaller完全指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考