尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

SeaTunnel MySQL JDBC Sink 连接器实战指南:批量写入、XA 精确一次与多表同步

SeaTunnel MySQL JDBC Sink 连接器实战指南:批量写入、XA 精确一次与多表同步 SeaTunnel MySQL JDBC Sink 连接器实战指南批量写入、XA 精确一次与多表同步【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel本篇指南系统讲解 SeaTunnel 中基于 JDBC 的 MySQL Sink 连接器的完整使用方法涵盖支持版本、依赖部署、数据类型映射、全部 Sink 参数、XA 事务精确一次语义以及 CDC / 多表同步等生产场景配置。读者读完可掌握从单表写入到多表动态映射的完整配置能力并能理解其底层实现原理连接池、批量执行器、XA 事务提交链路以便排障与调优。连接器概述MySQL Sink 是 SeaTunnel 通过 JDBC 协议将上游数据写入 MySQL 的 Sink 插件其工厂标识为Jdbc对应的方言实现为MysqlDialect源码见 MysqlDialect.java。该连接器具备以下核心能力批处理与流模式既可用于 BATCH 任务的批量落库也可用于 STREAMING 任务的实时写入并发写入配合partition_column等并行切分能力实现多并发写入精确一次Exactly-Once通过XA 事务保证仅适用于支持 XA 的数据库多表写入支持 MySQL CDC、JDBC Source 场景下的多表同步与占位符动态表名映射自动生成 Sink SQL配置databasetable后自动生成 insert/update/delete 语句无需手写 SQL。支持的 MySQL 版本连接器官方支持以下 MySQL 版本5.5 / 5.6 / 5.7 / 8.0 / 8.1 / 8.2 / 8.3 / 8.4注意不同的 MySQL 版本与驱动版本组合可能需要不同的驱动类与 URL 参数请以实际部署的驱动版本为准。引擎支持该连接器同时支持三种运行引擎Spark / Flink / SeaTunnel Zeta无论运行在哪种引擎上Sink 插件本身的配置语法与参数语义保持一致。主要功能清单功能说明精确一次依赖 XA 事务设置is_exactly_oncetrue且max_retries0启用CDC 支持可接收 CDC 变更事件Insert / Update / Delete并按主键生成对应 SQL多表写入支持 MySQL CDC、JDBC Source 的多表同步见 Connector V2 功能特性定时刷新通过batch_interval_ms实现基于时间间隔的批量刷新依赖项安装JDBC 连接器本身位于 connector-jdbc 模块但MySQL 的 JDBC 驱动 jar 包需要单独放置。对于 Spark / Flink 引擎将 mysql-connector-java 驱动 jar 包 放置到目录${SEATUNNEL_HOME}/plugins/中。对于 SeaTunnel Zeta 引擎将 MySQL 驱动 jar 包放置到目录${SEATUNNEL_HOME}/lib/中。支持的数据源信息数据源支持的版本驱动器网址Maven 下载链接Mysql不同的依赖版本具有不同的驱动程序类com.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/test下载从源码看MySQL 方言在 MysqlDialect.java 中通过defaultParameter()默认注入rewriteBatchedStatementstrue这会使 JDBC 驱动将多条 INSERT 合并为一条多值语句批量执行是批量写入性能的关键参数。数据类型映射Sink 写入时上游 SeaTunnel 数据类型会被转换reconvert为 MySQL 类型反之读取时 MySQL 类型会转换为 SeaTunnel 类型。完整映射关系见下表实现位于 MySqlTypeConverter.javaMySQL 数据类型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 UNSIGNEDDOUBLECHARVARCHARTINYTEXTMEDIUMTEXTTEXTLONGTEXTJSONSTRINGDATEDATETIMETIMEDATETIMETIMESTAMPTIMESTAMPTINYBLOBMEDIUMBLOBBLOBLONGBLOBBINARYVARBINARYBIT(n)BYTESGEOMETRYUNKNOWNNot supported yet值得注意的几点实现细节可从 MySqlTypeConverter.java 源码确认TINYINT(1)在启用int_type_narrowing时映射为 BOOLEAN否则映射为 BYTEMYSQL_TINYINT分支BIGINT UNSIGNED因超出 Java long 范围被提升为DECIMAL(20,0)TIMESTAMP被识别为 LTZ 类型存储为 UTC、按会话时区展示映射到OFFSET_DATE_TIME_TYPE回写时再转换回 MySQL TIMESTAMP反向转换SeaTunnel STRING → MySQL时根据列长度自动选择 VARCHAR / TEXT / MEDIUMTEXT / LONGTEXT字符串长度超过上限时默认使用 LONGTEXT。Sink 参数详解以下参数定义均可在源码 JdbcSinkOptions.java 与 JdbcSinkFactory.java 的optionRule()中得到印证。名称类型是否必填默认值描述urlString是-JDBC 连接的 URL。示例jdbc:mysql://localhost:3306/testdriverString是-用于连接远程数据源的 JDBC 类名MySQL 为com.mysql.cj.jdbc.DriverusernameString否-连接实例用户名passwordString否-连接实例密码queryString否-使用此 SQL 将上游输入数据写入数据库例如INSERT ...query具有更高优先级databaseString否-使用此database与table自动生成 SQL 并写入数据库此选项与query互斥优先级更高tableString否-使用 database 与此表名自动生成 SQL 并写入数据库此选项与query互斥优先级更高primary_keysArray否-用于自动生成 SQL 时支持insert、delete、update等操作connection_check_timeout_secInt否30等待用于验证连接的数据库操作完成的时间秒max_retriesInt否0提交失败的重试次数executeBatchbatch_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 数据源类名MySQL 为com.mysql.cj.jdbc.MysqlXADataSource其他数据源参考对应驱动文档max_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务打开后的超时秒默认 -1 表示永不超时注意设置超时可能影响精确一次语义auto_commitBoolean否true默认启用自动事务提交field_ideString否-字段名大小写转换ORIGINAL不转换UPPERCASE转大写LOWERCASE转小写propertiesMap否-其他连接配置参数当属性与 URL 参数重复时优先级由驱动实现决定例如 MySQL 中属性优先于 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存在primary_keys时启用 upsert若任务只有插入设为false可加快数据导入multi_table_sink_replicaInt否1多表写入时使用的 Sink Writer 副本数量参数联动与条件约束源码级从 JdbcSinkFactory.java 的optionRule()可以看到以下条件规则配置时必须遵守is_exactly_oncetrue时必须提供xa_data_source_class_name可选max_commit_attempts、transaction_timeout_sec此时max_retries会被强制置为 0见 JdbcConnectionConfig.javais_exactly_oncefalse时才会读取max_retriesgenerate_sink_sqltrue时必须提供databasegenerate_sink_sqlfalse时必须提供querydata_save_modeCUSTOM_PROCESSING时必须提供custom_sql。此外ExactlyOnceMaxRetriesValidator会在任务提交阶段校验启用精确一次时若max_retries ! 0会直接抛出OptionValidationException原因是 XA Sink 不支持重试否则可能造成重复数据。提示如果未设置partition_columnSink 将以单并发运行如果设置了partition_column将根据任务的并发度并行执行。任务示例示例一最简单的批量写入此示例通过 FakeSource 自动生成 16 行数据每行含name字符串与age整型两个字段写入 MySQL 的test_table表。运行前需先在 MySQL 中创建好测试库与test_table表若尚未安装部署 SeaTunnel请参照 安装 SeaTunnel运行方式参照 快速启动 SeaTunnel 引擎。# 定义运行时环境 env { parallelism 1 job.mode BATCH } source { # 仅供测试与演示的示例 Source 插件 FakeSource { parallelism 1 plugin_output fake row.num 16 schema { fields { name string age int } } } } transform { } sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username root password 123456 query insert into test_table(name,age) values(?,?) } }执行后目标表test_table中将出现 16 行数据。示例二自动生成 Sink SQL不需要手写复杂 SQL配置库名表名即可自动生成插入语句。query与database/table互斥配置generate_sink_sqltrue后走自动生成分支sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username root password 123456 # 根据数据库表名自动生成 SQL 语句 generate_sink_sql true database test table test_table } }从源码看自动生成的 INSERT 语句由方言的getInsertIntoStatement生成当配置了primary_keys且enable_upserttrue时MySQL 方言会生成INSERT ... ON DUPLICATE KEY UPDATE ...形式的 upsert 语句见 MysqlDialect.java。示例三精确一次XA 事务启用精确一次语义需要满足数据库支持 XA 事务 is_exactly_oncetruemax_retries0 正确配置xa_data_source_class_namesink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver max_retries 0 username root password 123456 query insert into test_table(name,age) values(?,?) is_exactly_once true xa_data_source_class_name com.mysql.cj.jdbc.MysqlXADataSource } }精确一次的底层实现启用is_exactly_once后Sink 会使用 JdbcExactlyOnceSinkWriter 替代普通 Writer其工作链路为开启 XA 事务beginTx()通过XidGenerator.semanticXidGenerator()生成全局唯一Xid并调用xaFacade.start(currentXid)开启事务源码 JdbcExactlyOnceSinkWriter.java写入与缓冲write()将记录克隆后写入JdbcOutputFormat数据先进入 XA 事务内的批缓冲Checkpoint 时 PrepareprepareCommit(checkpointId)先outputFormat.flush()刷空缓冲再调用xaFacade.endAndPrepare(currentXid)将事务置于 prepared 状态并把Xid存入状态快照snapshotState返回JdbcSinkState(prepareXid)提交/回滚由JdbcSinkCommitter/JdbcSinkAggregatedCommitter在 Checkpoint 完成后统一 commit任务失败时abortPrepare()会回滚已 prepare 的事务重启后还会通过recoverAndRollback清理悬挂的待恢复事务见tryOpen()中的恢复逻辑。正是这种写入→prepare→commit的两阶段提交流程配合 XA 数据源MysqlXADataSource保证了端到端的精确一次。这也是为什么开启精确一次时max_retries必须为 0——重试会破坏 XA 提交的幂等性可能产生重复数据。示例四CDC 变更数据事件Sink 支持接收 CDC 变更事件INSERT / UPDATE / DELETE。此时需要配置database、table与primary_keys由连接器按主键自动生成对应的增删改 SQLsink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username root password 123456 generate_sink_sql true # 需要同时配置 database 与 table database test table sink_table primary_keys [id,name] field_ide UPPERCASE schema_save_mode CREATE_SCHEMA_WHEN_NOT_EXIST data_save_mode APPEND_DATA } }示例五多表同步示例 1MySQL CDC 多表同步通过 MySQL CDC 将多张表同步到目标 MySQL使用${database_name}、${table_name}、${primary_key}占位符实现动态表名与主键映射env { parallelism 1 job.mode STREAMING checkpoint.interval 5000 } source { Mysql-CDC { url jdbc:mysql://127.0.0.1:3306/seatunnel username root password ****** table-names [seatunnel.role,seatunnel.user,galileo.Bucket] } } transform { } sink { Mysql { url jdbc:mysql://localhost:3306?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username root password 123456 generate_sink_sql true database ${database_name}_test table ${table_name}_test primary_keys [${primary_key}] } }示例 2JDBC Source 多表同步到 MySQL从 MySQL 使用 JDBC Source 以批量方式将多张表同步到另一个 MySQL 数据库env { parallelism 1 job.mode BATCH } source { Jdbc { driver com.mysql.cj.jdbc.Driver url jdbc:mysql://localhost:3306/source_db username root password 123456 table_list [ { table_path source_db.table_1 }, { table_path source_db.table_2 } ] } } transform { } sink { Mysql { url jdbc:mysql://localhost:3306?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver username root password 123456 generate_sink_sql true database ${database_name}_target table ${table_name}_copy primary_keys [${primary_key}] } }多表场景下可用multi_table_sink_replica参数控制 Sink Writer 副本数量默认 1以平衡并发写入能力与数据库压力。实践建议与注意事项URL 参数官方示例统一携带useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue。其中rewriteBatchedStatementstrue对批量写入性能至关重要且 MySQL 方言已将其作为默认参数注入MysqlDialect.java查询与自动生成二选一query与database/table互斥只有generate_sink_sqltrue时才会解析database与table精确一次的代价启用 XA 精确一次会带来额外的事务协调开销若业务可容忍 at-least-once 语义保持默认is_exactly_oncefalse可获得更高吞吐字段大小写field_ide支持ORIGINAL/UPPERCASE/LOWERCASE三种模式MySQL 方言在生成标识符时会统一加上反引号并通过getFieldIde做大小写转换MysqlDialect.java表结构处理schema_save_mode与data_save_mode支持建表、删表、追加等组合策略更多细节可参考 Sink Write Modes 文档连接校验JdbcSinkFactory实现了SupportSinkDryRunValidation可通过 dry-run 模式在任务启动前校验连通性、目标表存在性与字段兼容性建议 CI 阶段使用。变更日志连接器的功能迭代与修复记录见 connector-jdbc 变更日志。【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表