SeaTunnel MySQL JDBC Sink 连接器完全指南:配置、数据类型映射与 Exactly-Once 实战
数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载SeaTunnel 的 MySQL JDBC Sink 连接器通过标准 JDBC 接口将上游数据写入 MySQL支持批量Batch与流式Streaming两种模式、并发写入以及基于 XA 事务的 Exactly-Once 语义。本文以官方文档 Mysql.md 为主体结合仓库源码逐一讲解支持的版本、驱动依赖、数据类型映射、全部 Sink 选项、SQL 自动生成、CDC 事件处理与四种可落地的任务配置示例帮助你快速在 Spark、Flink 与 SeaTunnel Zeta 引擎中完成 MySQL 数据写入任务。支持的 MySQL 版本与运行引擎支持 MySQL 版本5.5 / 5.6 / 5.7 / 8.0 / 8.4支持运行引擎SparkFlinkSeaTunnel Zeta该连接器本质上是通用 JDBC Sink插件名称为Jdbc见 JdbcSink.java在 MySQL 方言Dialect下的具体实现。连接器通过MysqlDialect提供 MySQL 专属的 SQL 生成、批量写入参数与类型映射能力并通过MySqlCatalog支持建表等 Save Mode 操作。使用依赖驱动 JAR 的放置位置使用前需要确保 MySQL JDBC 驱动 JAR 已放置到正确目录Spark / Flink 引擎将驱动 JAR 放入${SEATUNNEL_HOME}/plugins/SeaTunnel Zeta 引擎将驱动 JAR 放入${SEATUNNEL_HOME}/lib/从源码看连接器在初始化多表资源管理器时会执行Class.forName(driver)显式加载驱动类见 JdbcSinkWriter.java因此驱动类必须存在于运行引擎的类路径中。关键特性Exactly-OnceCDCChange Data CaptureExactly-Once 的语义说明连接器使用XA 事务来保证 Exactly-Once因此该能力仅对支持 XA 事务的数据库生效。通过设置is_exactly_once true即可开启。源码层面的支撑JdbcSink在is_exactly_once开启时创建JdbcExactlyOnceSinkWriter与JdbcSinkAggregatedCommitter见 JdbcSink.java写入器通过XaFacade完成 XA 事务的 begin / end / preparecheckpoint 时记录Xid状态见 JdbcExactlyOnceSinkWriter.java。支持的数据源信息DatasourceSupported VersionsDriverUrlMavenMysql不同依赖版本对应不同驱动类com.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testmysql-connector-java需要留意的是不同版本的 MySQL Connector 驱动类名可能不同例如旧版com.mysql.jdbc.Driver与新版com.mysql.cj.jdbc.Driver务必根据实际引入的驱动版本填写driver配置项。MySQL 与 SeaTunnel 数据类型映射下表完整列出 MySQL 数据列到 SeaTunnel 数据类型的映射关系官方文档连接器通过MySqlTypeConverter与MySqlTypeMapper在源码中实现这一映射见 MySqlTypeConverter.java 与 MySqlTypeMapper.javaMySQL Data TypeSeaTunnel Data TypeBIT(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暂不支持几个源码级的补充细节可作为理解映射的参考源码中BIT(1)与TINYINT(1)均映射为BOOLEANBIT(n)n 1映射为BYTES字节长度按n/8向上取整计算见 MySqlTypeConverter.java。DECIMAL默认精度常量DEFAULT_PRECISION 38、默认小数位DEFAULT_SCALE 18超过 38 位精度时会被截断为DECIMAL(38,18)并输出可能溢出的告警日志同上文件#L196-L213。字符串类型在反向建表时按长度自动选择长度 2^8 用VARCHAR 2^16 用TEXT 2^24 用MEDIUMTEXT否则用LONGTEXT同上文件#L444-L467。MySqlTypeMapper对CHAR/VARCHAR/ENUM会按 4 字节字符集计算实际列长避免 UTF-8/UTF-8MB4 场景下精度失真见 MySqlTypeMapper.java。Sink 选项Sink Options详解以下为 MySQL JDBC Sink 的全部配置项继承自官方文档表格NameTypeRequiredDefaultDescriptionurlStringYes-JDBC 连接 URL示例jdbc:mysql://localhost:3306/testdriverStringYes-连接远程数据源使用的 JDBC 驱动类名MySQL 填com.mysql.cj.jdbc.DriveruserStringNo-连接实例用户名passwordStringNo-连接实例密码queryStringNo-自定义写入 SQL如INSERT ...优先级最高databaseStringNo-配合table自动生成写入 SQL与query互斥且优先级更高tableStringNo-配合database自动生成写入 SQL与query互斥且优先级更高primary_keysArrayNo-自动生成 SQL 时用于支持insert、delete、update操作support_upsert_by_query_primary_key_existBooleanNofalse数据库不支持 upsert 语法时通过查询主键是否存在来选择 INSERT 或 UPDATE SQL 处理更新事件INSERT、UPDATE_AFTER。注意该方式性能较低connection_check_timeout_secIntNo30等待数据库连接校验操作完成的超时时间秒max_retriesIntNo0提交失败executeBatch时的重试次数batch_sizeIntNo1000批量写入时当缓冲记录数达到batch_size或时间达到checkpoint.interval时将数据刷入数据库is_exactly_onceBooleanNofalse是否开启 Exactly-Once 语义使用 XA 事务。开启后需设置xa_data_source_class_namegenerate_sink_sqlBooleanNofalse根据目标数据库表自动生成 SQL 语句xa_data_source_class_nameStringNo-数据库驱动的 XA 数据源类名MySQL 为com.mysql.cj.jdbc.MysqlXADataSource其他数据源见附录max_commit_attemptsIntNo3事务提交失败的重试次数transaction_timeout_secIntNo-1事务开启后的超时时间默认 -1永不超时。注意设置超时可能影响 Exactly-Once 语义auto_commitBooleanNotrue默认开启自动事务提交field_ideStringNo-源到 Sink 同步时字段是否需要转换ORIGINAL不转换UPPERCASE转大写LOWERCASE转小写propertiesMapNo-额外的连接配置参数。当properties与 URL 中存在相同参数时优先级由驱动具体实现决定例如 MySQL 中properties优先于 URLcommon-options-No-Sink 插件通用参数详见 Sink Common Optionsschema_save_modeEnumNoCREATE_SCHEMA_WHEN_NOT_EXIST同步任务开启前对目标端表结构存在情况的不同处理方案data_save_modeEnumNoAPPEND_DATA同步任务开启前对目标端已存在数据的不同处理方案custom_sqlStringNo-当data_save_mode选择CUSTOM_PROCESSING时填写可执行的 SQL该 SQL 在同步任务开始前执行enable_upsertBooleanNotrue基于主键存在与否启用 upsert。若任务只有insert将其设为false可加快数据导入源码中的默认值与解析实现上述选项的默认值、类型与解析逻辑均可在 JdbcOptions.java 与 JdbcSinkConfig.java 中逐一印证例如connection_check_timeout_sec默认 30#L38-L42max_retries默认 0#L50-L51batch_size默认 1000#L81-L82is_exactly_once默认 false#L92-L96max_commit_attempts默认 3#L110-L114transaction_timeout_sec默认 -1#L116-L120auto_commit默认 true#L75-L79schema_save_mode默认CREATE_SCHEMA_WHEN_NOT_EXIST、data_save_mode默认APPEND_DATA#L61-L70enable_upsert默认 true#L137-L141support_upsert_by_query_primary_key_exist默认 false#L131-L135field_ide可选值来自FieldIdeEnumORIGINAL/UPPERCASE/LOWERCASE见#L183-L187另外两点值得注意的源码行为XA 模式强制max_retries 0JdbcExactlyOnceSinkWriter构造时校验maxRetries必须为 0否则会因重试导致数据重复见 JdbcExactlyOnceSinkWriter.java。MySQL 默认开启批量重写MysqlDialect.defaultParameter()会默认注入rewriteBatchedStatementstrue以提升批量写入性能见 MysqlDialect.java这也是官方示例 URL 中显式带上该参数的原因。Tips如果未设置partition_column任务将以单并发运行设置partition_column后将按任务并发数并行执行。任务示例Task Example以下示例均假定运行任务前已在 MySQL 中创建好数据库与目标表若尚未安装部署 SeaTunnel请先参考 安装 SeaTunnel再按 SeaTunnel Engine 快速上手 运行任务。示例一简单写入手动 SQL本示例通过 FakeSource 自动生成 16 行数据row.num16每行包含name字符串与age整数两个字段由 JDBC Sink 写入 MySQL 的test_table表最终表中应有 16 行数据。运行前需在 MySQL 中创建test库与test_table表。# Defining the runtime environment env { parallelism 1 job.mode BATCH } source { # 演示用 FakeSource 源插件 FakeSource { parallelism 1 result_table_name fake row.num 16 schema { fields { name string age int } } } } transform { # 如需了解 transform 插件配置请参考项目 transform-v2 文档 } sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver user root password 123456 query insert into test_table(name,age) values(?,?) } }要点说明URL 中rewriteBatchedStatementstrue与源码中MysqlDialect的默认参数一致用于优化批量写入使用query自定义 SQL 时占位符?的个数与顺序必须与上游 schema 字段一致本方式下generate_sink_sql、database、table均不需要配置。示例二自动生成 Sink SQL无需手写复杂 SQL只需配置数据库名与表名连接器即可自动生成插入语句。sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver user root password 123456 # 根据数据库表名自动生成 SQL 语句 generate_sink_sql true database test table test_table } }要点说明generate_sink_sql true时连接器基于上游 schema 与database、table生成 INSERT 语句相关配置解析见 JdbcSinkConfig.java该模式是后续 CDC 事件处理与 upsert 能力的基础。示例三Exactly-Once 精确一次写入适用于对数据准确性要求严格的场景通过 XA 事务保证每条数据仅写入一次。sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver max_retries 0 user 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 true开启 XA 事务xa_data_source_class_name必须填写 MySQL 的 XA 数据源类名com.mysql.cj.jdbc.MysqlXADataSource如前面源码所述XA 模式下max_retries必须保持为 0否则会导致重复数据见 JdbcExactlyOnceSinkWriter.java写入流程为生成 Xid → 开启 XA 事务xaFacade.start→ 批刷数据 → 事务 end/prepare → checkpoint 记录 Xid → 提交阶段由 AggregatedCommitter 完成 commit/rollback见 JdbcExactlyOnceSinkWriter.java。示例四CDCChange Data Capture事件处理当上游为 CDC 数据源如 MySQL CDC时Sink 可识别 INSERT / UPDATE / DELETE 等变更事件需要配置database、table与primary_keys。sink { jdbc { url jdbc:mysql://localhost:3306/test?useUnicodetruecharacterEncodingUTF-8rewriteBatchedStatementstrue driver com.mysql.cj.jdbc.Driver user 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_modeAPPEND_DATA } }要点说明generate_sink_sql true配合database、table自动生成 SQLprimary_keys声明主键字段Sink 据此生成 upsert / update / delete 语句处理 CDC 事件field_ide UPPERCASE表示字段名统一转换为大写可选ORIGINAL、LOWERCASEschema_save_mode与data_save_mode控制任务开启前对目标表结构与存量数据的处理策略。源码层面MySQL 的 upsert 通过INSERT ... ON DUPLICATE KEY UPDATE实现MysqlDialect.getUpsertStatement()会基于全部字段生成ON DUPLICATE KEY UPDATE 字段VALUES(字段)子句见 MysqlDialect.java对应执行器为 InsertOrUpdateBatchStatementExecutor.java。当任务仅包含 INSERT 且不需要 upsert 时可将enable_upsert false以提升导入速度。补充Save Mode 与建表能力当配置了database、table且未使用query自定义 SQL 时连接器通过DefaultSaveModeHandler执行schema_save_mode/data_save_mode策略见 JdbcSink.javaschema_save_mode可选值如CREATE_SCHEMA_WHEN_NOT_EXIST默认表不存在时自动建表、RECREATE_SCHEMA等用于处理目标表结构data_save_mode可选值如APPEND_DATA默认直接追加、TRUNCATE_TABLE、CUSTOM_PROCESSING配合custom_sql在任务启动前执行自定义 SQL等用于处理目标端存量数据MySQL 建表 SQL 由 MysqlCreateTableSqlBuilder.java 基于上游 schema 生成数据类型转换遵循本文前面给出的映射表。总结SeaTunnel 的 MySQL JDBC Sink 是一个覆盖手动 SQL 写入、自动生成 SQL、Exactly-Once 精确写入、CDC 事件同步四种主流场景的成熟连接器。实践要点可归纳为按运行引擎正确放置 MySQL 驱动 JAR数据量小、结构简单时用query手动 SQL希望免写 SQL 时开启generate_sink_sql并配置databasetable对准确性有硬性要求时开启is_exactly_once true并保持max_retries 0CDC 场景务必配置primary_keys必要时通过field_ide统一字段大小写需要自动建表或清空存量数据时配合使用schema_save_mode/data_save_mode/custom_sql。更深层的实现细节可继续阅读仓库源码JdbcOptions.java、MySqlTypeConverter.java、MysqlDialect.java 以及 Sink 写入器 JdbcSinkWriter.java。赞分享数据工程大数据批处理流处理【免费下载链接】seatunnelSeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.项目地址https://gitcode.com/gh_mirrors/sea/seatunnel点击查看免费下载相关推荐SeaTunnel JDBC Oracle Sink 连接器完整实战指南配置参数、数据类型映射与 XA 事务 Exactly-Once 写入SeaTunnel JDBC Oracle Sink 连接器完整实战指南配置参数、数据类型映射与 XA 事务 Exactly Once 写入 本文以 Apac数据工程大数据批处理流处理SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入SeaTunnel Kingbase Sink 连接器完全指南JDBC 配置、类型映射与实战写入 本文围绕 Kingbase Sink 连接器文档 https数据集成ETL大数据批处理流处理变更数据捕获SeaTunnel JDBC Snowflake Sink 连接器配置、CDC 写入与数据类型映射实战指南SeaTunnel JDBC Snowflake Sink 连接器配置、CDC 写入与数据类型映射实战指南 本文面向使用 Apache SeaTunnel h数据集成ETL大数据批处理流处理变更数据捕获上一篇10分钟上手TileStache从安装到启动地图瓦片服务的完整教程下一篇从Demo到实战Godot Card Game Framework卡牌扩展与定制教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考