☰
SeaTunnel MySQL JDBC Sink 连接器完全指南:配置、数据类型映射与 Exactly-Once 实战
2026/9/28 3:03:10 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel 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
  • 支持运行引擎:
    • Spark
    • Flink
    • SeaTunnel 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-Once
  • CDC(Change Data Capture)

Exactly-Once 的语义说明:连接器使用XA 事务来保证 Exactly-Once,因此该能力仅对支持 XA 事务的数据库生效。通过设置is_exactly_once = true即可开启。

源码层面的支撑:JdbcSink在is_exactly_once开启时创建JdbcExactlyOnceSinkWriter与JdbcSinkAggregatedCommitter(见 JdbcSink.java),写入器通过XaFacade完成 XA 事务的 begin / end / prepare,checkpoint 时记录Xid状态(见 JdbcExactlyOnceSinkWriter.java)。

支持的数据源信息

DatasourceSupported VersionsDriverUrlMaven
Mysql不同依赖版本对应不同驱动类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.java):

MySQL Data TypeSeaTunnel Data Type
BIT(1)、INT UNSIGNEDBOOLEAN
TINYINT、TINYINT UNSIGNED、SMALLINT、SMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEARINT
INT UNSIGNED、INTEGER UNSIGNED、BIGINTBIGINT
BIGINT UNSIGNEDDECIMAL(20,0)
DECIMAL(x,y)(列精度 < 38)DECIMAL(x,y)
DECIMAL(x,y)(列精度 > 38)DECIMAL(38,18)
DECIMAL UNSIGNEDDECIMAL(精度+1, 小数位)
FLOAT、FLOAT UNSIGNEDFLOAT
DOUBLE、DOUBLE UNSIGNEDDOUBLE
CHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、JSONSTRING
DATEDATE
TIMETIME
DATETIME、TIMESTAMPTIMESTAMP
TINYBLOB、MEDIUMBLOB、BLOB、LONGBLOB、BINARY、VARBINARY、BIT(n)BYTES
GEOMETRY、UNKNOWN暂不支持

几个源码级的补充细节(可作为理解映射的参考):

  • 源码中BIT(1)与TINYINT(1)均映射为BOOLEAN,BIT(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 的全部配置项(继承自官方文档表格):

NameTypeRequiredDefaultDescription
urlStringYes-JDBC 连接 URL,示例:jdbc:mysql://localhost:3306/test
driverStringYes-连接远程数据源使用的 JDBC 驱动类名,MySQL 填com.mysql.cj.jdbc.Driver
userStringNo-连接实例用户名
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_name
generate_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优先于 URL
common-options-No-Sink 插件通用参数,详见 Sink Common Options
schema_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-L42)
  • max_retries默认 0(#L50-L51)
  • batch_size默认 1000(#L81-L82)
  • is_exactly_once默认 false(#L92-L96)
  • max_commit_attempts默认 3(#L110-L114)
  • transaction_timeout_sec默认 -1(#L116-L120)
  • auto_commit默认 true(#L75-L79)
  • schema_save_mode默认CREATE_SCHEMA_WHEN_NOT_EXIST、data_save_mode默认APPEND_DATA(#L61-L70)
  • enable_upsert默认 true(#L137-L141)
  • support_upsert_by_query_primary_key_exist默认 false(#L131-L135)
  • field_ide可选值来自FieldIdeEnum(ORIGINAL/UPPERCASE/LOWERCASE,见#L183-L187)

另外两点值得注意的源码行为:

  1. XA 模式强制max_retries = 0:JdbcExactlyOnceSinkWriter构造时校验maxRetries必须为 0,否则会因重试导致数据重复(见 JdbcExactlyOnceSinkWriter.java)。
  2. MySQL 默认开启批量重写:MysqlDialect.defaultParameter()会默认注入rewriteBatchedStatements=true,以提升批量写入性能(见 MysqlDialect.java),这也是官方示例 URL 中显式带上该参数的原因。

Tips

如果未设置partition_column,任务将以单并发运行;设置partition_column后,将按任务并发数并行执行。

任务示例(Task Example)

以下示例均假定:运行任务前已在 MySQL 中创建好数据库与目标表;若尚未安装部署 SeaTunnel,请先参考 安装 SeaTunnel,再按 SeaTunnel Engine 快速上手 运行任务。

示例一:简单写入(手动 SQL)

本示例通过 FakeSource 自动生成 16 行数据(row.num=16,每行包含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?useUnicode=true&characterEncoding=UTF-8&rewriteBatchedStatements=true" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" query = "insert into test_table(name,age) values(?,?)" } }

要点说明:

  • URL 中rewriteBatchedStatements=true与源码中MysqlDialect的默认参数一致,用于优化批量写入;
  • 使用query自定义 SQL 时,占位符?的个数与顺序必须与上游 schema 字段一致;
  • 本方式下generate_sink_sql、database、table均不需要配置。

示例二:自动生成 Sink SQL

无需手写复杂 SQL,只需配置数据库名与表名,连接器即可自动生成插入语句。

sink { jdbc { url = "jdbc:mysql://localhost:3306/test?useUnicode=true&characterEncoding=UTF-8&rewriteBatchedStatements=true" 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?useUnicode=true&characterEncoding=UTF-8&rewriteBatchedStatements=true" 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)。

示例四:CDC(Change Data Capture)事件处理

当上游为 CDC 数据源(如 MySQL CDC)时,Sink 可识别 INSERT / UPDATE / DELETE 等变更事件,需要配置database、table与primary_keys。

sink { jdbc { url = "jdbc:mysql://localhost:3306/test?useUnicode=true&characterEncoding=UTF-8&rewriteBatchedStatements=true" 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_mode="APPEND_DATA" } }

要点说明:

  • generate_sink_sql = true配合database、table自动生成 SQL;
  • primary_keys声明主键字段,Sink 据此生成 upsert / update / delete 语句处理 CDC 事件;
  • field_ide = UPPERCASE表示字段名统一转换为大写(可选ORIGINAL、LOWERCASE);
  • schema_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.java):

  • schema_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 事件同步"四种主流场景的成熟连接器。实践要点可归纳为:

  1. 按运行引擎正确放置 MySQL 驱动 JAR;
  2. 数据量小、结构简单时用query手动 SQL;希望免写 SQL 时开启generate_sink_sql并配置database+table;
  3. 对准确性有硬性要求时开启is_exactly_once = true并保持max_retries = 0;
  4. CDC 场景务必配置primary_keys,必要时通过field_ide统一字段大小写;
  5. 需要自动建表或清空存量数据时,配合使用schema_save_mode/data_save_mode/custom_sql。

更深层的实现细节可继续阅读仓库源码:JdbcOptions.java、MySqlTypeConverter.java、MysqlDialect.java 以及 Sink 写入器 JdbcSinkWriter.java。

  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.

项目地址:https://gitcode.com/gh_mirrors/sea/seatunnel
点击查看免费下载

相关推荐

上一篇:10分钟上手TileStache:从安装到启动地图瓦片服务的完整教程
下一篇:从Demo到实战:Godot Card Game Framework卡牌扩展与定制教程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询