☰
SeaTunnel JDBC Sink 使用指南:基于 XA 事务的精确一次写入与 CDC 事件落库实践
2026/10/9 1:33:41 网站建设 项目流程
  • 数据工程
  • 大数据
  • 批处理
  • 流处理

【免费下载链接】seatunnel

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

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

本文基于 Apache SeaTunnel 仓库中 JDBC 数据接收器文档 编写,深入讲解 JDBC Sink 插件的全部配置参数、写入原理与实战示例。JDBC Sink 是 SeaTunnel 连接任意 JDBC 兼容数据库(MySQL、PostgreSQL、Oracle、SQL Server、达梦、OceanBase 等)的核心输出插件,支持批处理与流处理、并发写入、基于 XA 事务的 exactly-once 语义,以及 CDC(Change Data Capture)事件的 INSERT/UPDATE/DELETE 落库。读完本文,你将掌握 JDBC Sink 的驱动部署、参数调优、精确一次配置、表结构/数据保存策略,以及 MySQL、PostgreSQL 等主流数据库的完整接入方案。

功能概述

JDBC Sink 通过 JDBC 协议将上游数据写入目标数据库,具备以下核心能力:

  • 批处理与流处理双模式:既能承接离线批量同步,也能在流式任务中持续写入;
  • 并发写入:多并行度 Writer 同时写入,配合连接池复用连接资源;
  • 精确一次语义(exactly-once):基于 XA 事务(两阶段提交)保证,仅对支持 XA 事务的数据库生效,可通过is_exactly_once=true开启;
  • CDC 事件支持:可消费上游 CDC 产生的 INSERT、UPDATE、DELETE 事件并转化为对应的 SQL 操作,详见 Connector-V2 特性说明。

从源码看,插件的工厂标识符为Jdbc(见 JdbcSinkFactory.java),Sink 实现在 JdbcSink.java 中根据is_exactly_once选择不同的 Writer:开启时使用JdbcExactlyOnceSinkWriter(XA 两阶段提交),未开启时使用普通JdbcSinkWriter(批量提交)。

使用依赖:驱动 JAR 部署

用于 Spark/Flink 引擎

需要确保 JDBC 驱动 JAR 包已放入目录${SEATUNNEL_HOME}/plugins/下。

适用于 SeaTunnel Zeta 引擎

需要确保 JDBC 驱动 JAR 包已放入${SEATUNNEL_HOME}/lib/目录下。

驱动类加载是 Sink 正常工作的前提:在 JdbcSink.java 的getSaveModeHandler()中会先执行Class.forName(driverName)校验驱动是否存在,若类找不到将直接抛出ClassNotFoundException。

Options 参数总览

名称类型是否必须默认值
urlString是-
driverString是-
userString否-
passwordString否-
queryString否-
compatible_modeString否-
databaseString否-
tableString否-
primary_keysArray否-
support_upsert_by_query_primary_key_existBoolean否false
connection_check_timeout_secInt否30
max_retriesInt否0
batch_sizeInt否1000
is_exactly_onceBoolean否false
generate_sink_sqlBoolean否false
xa_data_source_class_nameString否-
max_commit_attemptsInt否3
transaction_timeout_secInt否-1
auto_commitBoolean否true
field_ideString否-
propertiesMap否-
common-options否-
schema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST
data_save_modeEnum否APPEND_DATA
custom_sqlString否-
enable_upsertBoolean否true
use_copy_statementBoolean否false

以上默认值均可在 JdbcOptions.java 的 Option 定义中找到对应实现。此外,JdbcSinkFactory.java 的optionRule()定义了参数之间的条件依赖关系,编写配置时需要遵守:

  • 必填:url、driver、schema_save_mode、data_save_mode;
  • is_exactly_once=true时,必须提供xa_data_source_class_name,可选max_commit_attempts、transaction_timeout_sec;
  • is_exactly_once=false时,可使用max_retries;
  • generate_sink_sql=true时需提供database;generate_sink_sql=false时需提供query;
  • data_save_mode=CUSTOM_PROCESSING时必须提供custom_sql。

核心参数详解

driver [string]

用于连接远程数据源的 JDBC 驱动类名。例如使用 MySQL 时值为com.mysql.cj.jdbc.Driver。

user [string]

连接数据库的用户名。

password [string]

连接数据库的密码。

url [string]

JDBC 连接的 URL,参考案例:jdbc:postgresql://localhost/test。当开启精确一次语义时,URL 会被驱动加载器解析并与properties合并(见 JdbcSinkFactory.java 中的dialect.connectionUrlParse逻辑)。

query [string]

使用 SQL 语句将上游输入数据写入数据库,如INSERT ...。使用占位符?即可按上游字段顺序绑定参数。注意:配置了query时无法使用 Save Mode 相关的表结构/数据处理策略(源码中 JdbcSink.java 明确在simpleSql非空时返回空的 SaveModeHandler)。

compatible_mode [string]

数据库的兼容模式,当数据库支持多种兼容模式时需要设置。例如使用 OceanBase 数据库时,需要设置为'mysql'或'oracle'。Postgres 9.5 及以下版本,请设置为postgresLow以支持 CDC。该参数会传递给JdbcDialectLoader决定加载哪种方言实现(仓库中存在独立的 PostgresLowDialect.java)。

database [string]

配合table自动生成写入 SQL,并接收上游输入的数据写入数据库。此选项与query选项互斥,且优先级更高。

table [string]

配合database自动生成 SQL。table可以填入任意表名,该名字最终用作创建/写入的目标表名,并且支持变量(${table_name}、${schema_name})。替换规则:${schema_name}替换为传递给目标端的 SCHEMA 名称,${table_name}替换为传递给目标端的表名。

MySQL 接收器示例:

  1. test_${schema_name}_${table_name}_test
  2. sink_sinktable
  3. ss_${table_name}

PostgreSQL(Oracle、SQL Server 等)接收器示例:

  1. ${schema_name}.${table_name}_test
  2. dbo.tt_${table_name}_sink
  3. public.sink_table

Tip:如果目标数据库有 SCHEMA 概念,则表参数必须写成xxx.xxx形式。

在源码中,table参数的解析与变量替换位于 JdbcSinkFactory.java:table按.拆分出 schema 与表名,并依次支持table_prefix/table_suffix前后缀追加以及${table_name}、${schema_name}、${database_name}的替换。

primary_keys [array]

该选项用于辅助生成 insert、delete、update 等 SQL 语句。设置后,Sink 会根据主键列自动生成对应的 upsert/update SQL。若未显式配置,JdbcSinkFactory.java 会尝试从上游 CatalogTable 的主键或唯一键(UNIQUE_KEY)约束中自动推导主键列。

support_upsert_by_query_primary_key_exist [boolean]

根据查询主键是否存在来选择使用 INSERT SQL 或 UPDATE SQL 处理变更事件(INSERT、UPDATE_AFTER)。仅当数据库不支持 upsert 语法时才使用此配置。

注意:该方法性能较低,因为每条记录都要先执行一次主键存在性查询。其实现对应源码中的 InsertOrUpdateBatchStatementExecutor.java,会为每条数据执行SELECT判断存在性,再决定走 insert 还是 update 批。

connection_check_timeout_sec [int]

用于验证数据库连接有效性时等待数据库操作完成所需的时间,单位秒,默认 30。写入失败重试前会通过该机制判断连接是否需要重建(见 JdbcOutputFormat.java 的重连逻辑)。

max_retries [int]

重试提交失败的最大次数(针对executeBatch),默认 0 即不重试。注意:开启is_exactly_once=true时该参数会被强制置为 0(见 JdbcConnectionConfig.java),因为 XA sink 的重试可能造成数据重复;JdbcExactlyOnceSinkWriter.java 构造函数同样用checkArgument强制maxRetries == 0。

batch_size [int]

对于批量写入,当缓冲的记录数达到batch_size数量或者时间达到checkpoint.interval时,数据将被刷新到数据库,默认 1000。批提交逻辑在 JdbcOutputFormat.java 中实现:batchCount达到阈值即触发flush()执行executeBatch()。

is_exactly_once [boolean]

是否启用通过 XA 事务实现的精确一次语义,默认 false。开启后还需设置xa_data_source_class_name。开启后 Sink 会启用两阶段提交:Writer 端beginTx → write → prepare,再由JdbcSinkAggregatedCommitter统一commit(详见下文"精确一次写入原理")。

generate_sink_sql [boolean]

根据要写入的数据库表结构生成 SQL 语句,默认 false。开启时需配合database使用。

xa_data_source_class_name [string]

数据库驱动的 XA 数据源类名。以 MySQL 为例,其类名为com.mysql.cj.jdbc.MysqlXADataSource。其他数据库的 XA 数据源类名可参考文末附录。

max_commit_attempts [int]

事务提交失败的最大重试次数,默认 3。该值会在 JdbcSinkAggregatedCommitter.java 的commit()中被传递给 XA 分组提交逻辑,提交失败的 Xid 会被放入待重试列表,由下一次 checkpoint 继续提交。

transaction_timeout_sec [int]

事务开启后的超时时间,默认 -1(即永不超时)。注意:设置超时时间可能会影响 exactly-once 语义。在 JdbcConnectionConfig.java 中,小于 0 的值被解析为Optional.empty(),即不设置事务超时。

auto_commit [boolean]

默认启用自动事务提交(true)。在普通写入模式下,若关闭自动提交,Writer 会在prepareCommit()与close()时显式执行commit()(见 JdbcSinkWriter.java)。

field_ide [String]

字段field_ide用于在从 source 同步到 sink 时,确定字段是否需要转换大小写:ORIGINAL表示不转换,UPPERCASE表示转换为大写,LOWERCASE表示转换为小写。该值同时会写入 CatalogTable options,影响建表时字段名的引用方式(见 JdbcSinkFactory.java)。

properties [Map]

附加连接配置参数,当属性和 URL 具有相同参数时,优先级由驱动程序的实现决定。例如在 MySQL 中,属性配置优先于 URL。示例:properties { rewriteBatchedStatements = true }。

common options

Sink 插件常用参数,请参考 Sink 常用选项 了解详情。

schema_save_mode [Enum]

在启动同步任务之前,针对目标侧已有的表结构选择不同的处理方案:

  • RECREATE_SCHEMA:当表不存在时会创建,当表已存在时会删除并重建;
  • CREATE_SCHEMA_WHEN_NOT_EXIST:当表不存在时会创建,当表已存在时则跳过创建(默认值);
  • ERROR_WHEN_SCHEMA_NOT_EXIST:当表不存在时将抛出错误。

data_save_mode [Enum]

在启动同步任务之前,针对目标侧已存在的数据选择不同的处理方案:

  • DROP_DATA:保留数据库结构,删除数据;
  • APPEND_DATA:保留数据库结构,保留数据(默认值);
  • CUSTOM_PROCESSING:允许用户自定义数据处理方式;
  • ERROR_WHEN_DATA_EXISTS:当有数据时抛出错误。

custom_sql [String]

当data_save_mode选择CUSTOM_PROCESSING时必须填写。该参数通常填写一条可执行的 SQL,将在同步任务之前执行(例如先清空目标表数据)。SaveMode 处理器会将其作为自定义 SQL 传入(见 JdbcSink.java)。

enable_upsert [boolean]

启用通过主键更新插入(upsert),默认 true。如果任务没有 key 重复数据,设置该参数为 false 可以加快数据导入速度(跳过 upsert 语句的生成与判断开销)。该开关在 JdbcSinkConfig.java 中被读取并决定 SQL 生成策略。

use_copy_statement [boolean]

使用COPY ${table} FROM STDIN语句导入数据,默认 false。仅支持具有getCopyAPI()方法连接的驱动程序,例如 PostgreSQL 驱动org.postgresql.Driver。

注意:不支持MAP、ARRAY、ROW类型。实现上,CopyManagerBatchStatementExecutor.java 通过反射调用连接的getCopyAPI()获取 CopyManager,若驱动不支持会提示关闭use_copy_statement。

写入流程与底层实现

普通模式(批量写入 + 自动/手动提交)

未开启精确一次时,写入链路为:JdbcSinkWriter.write()→JdbcOutputFormat.writeRecord()。核心机制见 JdbcOutputFormat.java:

  1. 记录先被加入PreparedStatement批中(addToBatch),batchCount累加;
  2. 当batchCount >= batch_size时触发flush(),执行executeBatch();
  3. flush()失败时按max_retries进行重试,每次重试前会校验连接有效性,必要时重建连接;
  4. 若关闭了auto_commit,prepareCommit()或close()阶段会显式提交事务;
  5. 每个 checkpoint 周期也会驱动一次 flush,因此流处理模式下数据刷新时机为"达到 batch_size 或到达 checkpoint.interval"。

JdbcSinkWriter内部通过方言工厂构建连接 Provider 与输出格式(见 JdbcSinkWriter.java),多表场景下还会基于 HikariCP 构建共享连接池(initMultiTableResourceManager)。

精确一次写入原理(XA 两阶段提交)

开启is_exactly_once=true后,写入链路切换为 XA 事务模式,整个过程分布在 Writer 与 AggregatedCommitter 两个角色上:

  1. Writer 端(JdbcExactlyOnceSinkWriter.java):
    • 初始化时通过XaFacade.fromJdbcConnectionOptions创建 XA 连接门面,并用XidGenerator.semanticXidGenerator()生成语义化 Xid;
    • 每轮 checkpoint 先beginTx()(xaFacade.start(xid)),写入数据;
    • prepareCommit()时执行outputFormat.flush()并调用xaFacade.endAndPrepare(xid)进入 prepare 阶段,将XidInfo作为提交信息返回;空事务会被跳过(EmptyXaTransactionException);
    • 故障恢复时,Writer 会调用recoverAndRollback回滚不属于当前恢复状态的悬空事务,避免重复数据;
  2. AggregatedCommitter 端(JdbcSinkAggregatedCommitter.java):
    • 聚合所有 Writer 上报的XidInfo(combine),统一执行 XAcommit;
    • 提交失败的事务按max_commit_attempts重试,仍失败者返回待重试列表,等待下一次 checkpoint 再提交;
    • 任务失败触发回滚时,通过abort()对未提交的 Xid 执行 rollback。

注意:XA 模式下max_retries被强制为 0,请勿手动配置重试,否则可能产生重复数据。

CDC 事件(INSERT/UPDATE/DELETE)处理

Sink 对 CDC 变更事件的支持依赖两个执行器:

  • 存在性查询模式(support_upsert_by_query_primary_key_exist=true):每条记录先按主键查询目标表,存在则走 UPDATE,不存在则走 INSERT,见 InsertOrUpdateBatchStatementExecutor.java;
  • 缓冲归并模式(默认,无查询开销):BufferReducedBatchStatementExecutor.java 按主键将记录缓冲在LinkedHashMap中,同一主键的多条变更被合并为最终状态,checkpoint 时统一执行 upsert 或 delete 批,既保证语义正确又显著降低写库次数。

上游事件对应的 RowKind 转换规则为:INSERT / UPDATE_AFTER 走 upsert,DELETE / UPDATE_BEFORE 走 delete,UPDATE_BEFORE 本身会被忽略(只作为 UPDATE_AFTER 的前置占位)。

表结构与数据处理策略(Save Mode)

当配置了database+table(且未配置query)时,Sink 会在任务启动前通过 Catalog 对目标侧执行 SaveMode 处理(源码入口见 JdbcSink.java):

  • 驱动类加载校验 → 按 URL 与兼容模式查找对应方言的 Catalog 实现;
  • 结合schema_save_mode决定建表/重建/报错;
  • 结合data_save_mode决定保留数据/清空数据/自定义 SQL 预执行/数据存在即报错;
  • field_ide会作为fieldIde写入 CatalogTable options,影响建表 SQL 中字段引用的大小写。

query与 Save Mode 互斥:一旦配置了query,任务将直接按用户 SQL 写入,不再进行表结构管理。

Tips:启用精确一次的数据库要求

在is_exactly_once = "true"的情况下使用 XA 事务,需要数据库支持。部分数据库需要额外设置:

  1. PostgreSQL:需要设置max_prepared_transactions > 1,例如ALTER SYSTEM set max_prepared_transactions to 10;
  2. MySQL:版本需要>= 8.0.29,且非 root 用户需要授予XA_RECOVER_ADMIN权限。例如:将test_db.*上的XA_RECOVER_ADMIN授予'user1'@'%';
  3. MySQL 性能优化:可以尝试在 url 中添加rewriteBatchedStatements=true参数以获得更好的批量写入性能。

附录:常见数据源接入参考

附录参数仅提供参考,Maven 列给出的是可在 Maven 中央仓库检索到的驱动 artifact 名称:

数据源driverurlxa_data_source_class_nameMaven artifact
MySQLcom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/testcom.mysql.cj.jdbc.MysqlXADataSourcemysql:mysql-connector-java
PostgreSQLorg.postgresql.Driverjdbc:postgresql://localhost:5432/postgresorg.postgresql.xa.PGXADataSourceorg.postgresql:postgresql
DM(达梦)dm.jdbc.driver.DmDriverjdbc:dm://localhost:5236dm.jdbc.driver.DmdbXADataSourcecom.dameng:DmJdbcDriver18
Phoenixorg.apache.phoenix.queryserver.client.Driverjdbc:phoenix:thin:url=http://localhost:8765;serialization=PROTOBUF/com.aliyun.phoenix:ali-phoenix-shaded-thin-client
SQL Servercom.microsoft.sqlserver.jdbc.SQLServerDriverjdbc:sqlserver://localhost:1433com.microsoft.sqlserver.jdbc.SQLServerXADataSourcecom.microsoft.sqlserver:mssql-jdbc
Oracleoracle.jdbc.OracleDriverjdbc:oracle:thin:@localhost:1521/xepdb1oracle.jdbc.xa.OracleXADataSourcecom.oracle.database.jdbc:ojdbc8
SQLiteorg.sqlite.JDBCjdbc:sqlite:test.db/org.xerial:sqlite-jdbc
GBase8acom.gbase.jdbc.Driverjdbc:gbase://e2e_gbase8aDb:5258/test/gbase-connector-java
StarRockscom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/test/mysql:mysql-connector-java
DB2com.ibm.db2.jcc.DB2Driverjdbc:db2://localhost:50000/testdbcom.ibm.db2.jcc.DB2XADataSourcecom.ibm.db2.jcc:db2jcc4
SAP HANAcom.sap.db.jdbc.Driverjdbc:sap://localhost:39015/com.sap.cloud.db.jdbc:ngdbc
Doriscom.mysql.cj.jdbc.Driverjdbc:mysql://localhost:3306/test/mysql:mysql-connector-java
Teradatacom.teradata.jdbc.TeraDriverjdbc:teradata://localhost/DBS_PORT=1025,DATABASE=test/com.teradata.jdbc:terajdbc
Redshiftcom.amazon.redshift.jdbc42.Driverjdbc:redshift://localhost:5439/testdbcom.amazon.redshift.xa.RedshiftXADataSourcecom.amazon.redshift:redshift-jdbc42
Snowflakenet.snowflake.client.jdbc.SnowflakeDriverjdbc//<account_name>.snowflakecomputing.com/net.snowflake:snowflake-jdbc
Verticacom.vertica.jdbc.Driverjdbc:vertica://localhost:5433/com.vertica:vertica-jdbc
Kingbasecom.kingbase8.Driverjdbc:kingbase8://localhost:54321/db_test/cn.com.kingbase:kingbase8
OceanBasecom.oceanbase.jdbc.Driverjdbc:oceanbase://localhost:2881/com.oceanbase:oceanbase-client

各数据库对应的方言实现(含类型映射与 SQL 生成)位于仓库 internal/dialect 目录下,可按需阅读。

配置示例

1. 简单示例(自定义 query 写入)

jdbc { url = "jdbc:mysql://localhost:3306/test" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" query = "insert into test_table(name,age) values(?,?)" }

2. 精确一次(Exactly-once)

通过设置is_exactly_once = "true"并指定 XA 数据源类名开启精确一次语义:

jdbc { url = "jdbc:mysql://localhost:3306/test" 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" }

3. 变更数据捕获(CDC)事件写入

JDBC Sink 消费并落库 CDC 事件的示例(database/table自动生成 SQL,primary_keys辅助生成 upsert 语句):

sink { jdbc { url = "jdbc:mysql://localhost:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" database = "sink_database" table = "sink_table" primary_keys = ["key1", "key2", ...] } }

4. 配置表生成策略(schema_save_mode)

通过将schema_save_mode配置为CREATE_SCHEMA_WHEN_NOT_EXIST,在目标表不存在时自动建表;data_save_mode="APPEND_DATA"表示保留已有数据:

sink { jdbc { url = "jdbc:mysql://localhost:3306" driver = "com.mysql.cj.jdbc.Driver" user = "root" password = "123456" database = "sink_database" table = "sink_table" primary_keys = ["key1", "key2", ...] schema_save_mode = "CREATE_SCHEMA_WHEN_NOT_EXIST" data_save_mode = "APPEND_DATA" } }

5. Postgres 9.5 及以下版本 CDC 示例

Postgres 9.5 及以下版本需将compatible_mode配置为postgresLow来支持 Postgres CDC 操作,配合support_upsert_by_query_primary_key_exist=true与generate_sink_sql=true使用:

sink { jdbc { url = "jdbc:postgresql://localhost:5432" driver = "org.postgresql.Driver" user = "root" password = "123456" compatible_mode = "postgresLow" database = "sink_database" table = "sink_table" support_upsert_by_query_primary_key_exist = true generate_sink_sql = true primary_keys = ["key1", "key2", ...] } }

变更日志

2.3.0-beta 2022-10-20

  • [BugFix] 修复 JDBC split 异常;
  • [Feature] 支持 Phoenix JDBC Sink;
  • [Feature] 支持 SQL Server JDBC Sink;
  • [Feature] 支持 Oracle JDBC Sink;
  • [Feature] 支持 StarRocks JDBC Sink;
  • [Feature] 支持 DB2 JDBC Sink。

next version

  • [Feature] 支持 CDC 写入 DELETE/UPDATE/INSERT 事件;
  • [Feature] 支持 Teradata JDBC Sink;
  • [Feature] 支持 SQLite JDBC Sink;
  • [Feature] 支持 Doris JDBC Sink;
  • [Feature] 支持 Redshift JDBC Sink;
  • [Improve] 新增按查询启用 upsert 的配置项;
  • [Improve] Sink 配置中新增 database 字段;
  • [Improve] 新增 Vertica connector。

小结

JDBC Sink 是 SeaTunnel 生态中覆盖数据库最广的输出插件之一。结合本仓库源码可以看到,它的可靠性来自三套设计:普通模式下基于batch_size与 checkpoint 驱动的批量提交、精确一次模式下基于 XA 两阶段提交的端到端一致性、以及 CDC 场景下按主键缓冲归并的 upsert/delete 执行器。实际使用时,只需根据引擎类型放好驱动 JAR,再按本文的参数矩阵与示例组合出适合自己业务的配置即可。

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

【免费下载链接】seatunnel

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

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

相关推荐

上一篇:Picasso请求生命周期全解析:cancelRequest、tag批量暂停恢复与Priority优先级调度原理
下一篇:Android RecyclerView高级用法:基于Sample项目的自定义LayoutManager实现

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

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

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

立即咨询