- 数据工程
- 大数据
- 批处理
- 流处理
【免费下载链接】seatunnel
SeaTunnel is a next-generation super high-performance, distributed, massive data integration tool.
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)。
支持的数据源信息
| Datasource | Supported Versions | Driver | Url | Maven |
|---|---|---|---|---|
| Mysql | 不同依赖版本对应不同驱动类 | com.mysql.cj.jdbc.Driver | jdbc:mysql://localhost:3306/test | mysql-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 Type | SeaTunnel Data Type |
|---|---|
BIT(1)、INT UNSIGNED | BOOLEAN |
TINYINT、TINYINT UNSIGNED、SMALLINT、SMALLINT UNSIGNED、MEDIUMINT、MEDIUMINT UNSIGNED、INT、INTEGER、YEAR | INT |
INT UNSIGNED、INTEGER UNSIGNED、BIGINT | BIGINT |
BIGINT UNSIGNED | DECIMAL(20,0) |
DECIMAL(x,y)(列精度 < 38) | DECIMAL(x,y) |
DECIMAL(x,y)(列精度 > 38) | DECIMAL(38,18) |
DECIMAL UNSIGNED | DECIMAL(精度+1, 小数位) |
FLOAT、FLOAT UNSIGNED | FLOAT |
DOUBLE、DOUBLE UNSIGNED | DOUBLE |
CHAR、VARCHAR、TINYTEXT、MEDIUMTEXT、TEXT、LONGTEXT、JSON | STRING |
DATE | DATE |
TIME | TIME |
DATETIME、TIMESTAMP | TIMESTAMP |
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 的全部配置项(继承自官方文档表格):
| Name | Type | Required | Default | Description |
|---|---|---|---|---|
url | String | Yes | - | JDBC 连接 URL,示例:jdbc:mysql://localhost:3306/test |
driver | String | Yes | - | 连接远程数据源使用的 JDBC 驱动类名,MySQL 填com.mysql.cj.jdbc.Driver |
user | String | No | - | 连接实例用户名 |
password | String | No | - | 连接实例密码 |
query | String | No | - | 自定义写入 SQL,如INSERT ...,优先级最高 |
database | String | No | - | 配合table自动生成写入 SQL;与query互斥且优先级更高 |
table | String | No | - | 配合database自动生成写入 SQL;与query互斥且优先级更高 |
primary_keys | Array | No | - | 自动生成 SQL 时用于支持insert、delete、update操作 |
support_upsert_by_query_primary_key_exist | Boolean | No | false | 数据库不支持 upsert 语法时,通过查询主键是否存在来选择 INSERT 或 UPDATE SQL 处理更新事件(INSERT、UPDATE_AFTER)。注意:该方式性能较低 |
connection_check_timeout_sec | Int | No | 30 | 等待数据库连接校验操作完成的超时时间(秒) |
max_retries | Int | No | 0 | 提交失败(executeBatch)时的重试次数 |
batch_size | Int | No | 1000 | 批量写入时,当缓冲记录数达到batch_size或时间达到checkpoint.interval时,将数据刷入数据库 |
is_exactly_once | Boolean | No | false | 是否开启 Exactly-Once 语义(使用 XA 事务)。开启后需设置xa_data_source_class_name |
generate_sink_sql | Boolean | No | false | 根据目标数据库表自动生成 SQL 语句 |
xa_data_source_class_name | String | No | - | 数据库驱动的 XA 数据源类名,MySQL 为com.mysql.cj.jdbc.MysqlXADataSource,其他数据源见附录 |
max_commit_attempts | Int | No | 3 | 事务提交失败的重试次数 |
transaction_timeout_sec | Int | No | -1 | 事务开启后的超时时间,默认 -1(永不超时)。注意:设置超时可能影响 Exactly-Once 语义 |
auto_commit | Boolean | No | true | 默认开启自动事务提交 |
field_ide | String | No | - | 源到 Sink 同步时字段是否需要转换:ORIGINAL不转换;UPPERCASE转大写;LOWERCASE转小写 |
properties | Map | No | - | 额外的连接配置参数。当properties与 URL 中存在相同参数时,优先级由驱动具体实现决定,例如 MySQL 中properties优先于 URL |
common-options | - | No | - | Sink 插件通用参数,详见 Sink Common Options |
schema_save_mode | Enum | No | CREATE_SCHEMA_WHEN_NOT_EXIST | 同步任务开启前,对目标端表结构存在情况的不同处理方案 |
data_save_mode | Enum | No | APPEND_DATA | 同步任务开启前,对目标端已存在数据的不同处理方案 |
custom_sql | String | No | - | 当data_save_mode选择CUSTOM_PROCESSING时,填写可执行的 SQL,该 SQL 在同步任务开始前执行 |
enable_upsert | Boolean | No | true | 基于主键存在与否启用 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)
另外两点值得注意的源码行为:
- XA 模式强制
max_retries = 0:JdbcExactlyOnceSinkWriter构造时校验maxRetries必须为 0,否则会因重试导致数据重复(见 JdbcExactlyOnceSinkWriter.java)。 - 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 事件同步"四种主流场景的成熟连接器。实践要点可归纳为:
- 按运行引擎正确放置 MySQL 驱动 JAR;
- 数据量小、结构简单时用
query手动 SQL;希望免写 SQL 时开启generate_sink_sql并配置database+table; - 对准确性有硬性要求时开启
is_exactly_once = true并保持max_retries = 0; - CDC 场景务必配置
primary_keys,必要时通过field_ide统一字段大小写; - 需要自动建表或清空存量数据时,配合使用
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.
相关推荐
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大数据批处理流处理变更数据捕获
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考