SeaTunnel JDBC SQL Server Sink 连接器实战指南:配置详解、数据类型映射与精确一次写入
2026/9/19 12:57:40 网站建设 项目流程

SeaTunnel JDBC SQL Server Sink 连接器实战指南:配置详解、数据类型映射与精确一次写入

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

本文以 Apache SeaTunnel 仓库中的 docs/zh/connectors/sink/SqlServer.md 为主体,结合connector-jdbc模块中 SQL Server 方言与类型转换的真实源码,系统讲解如何通过 SeaTunnel 将数据写入 SQL Server:从驱动安装、连接参数、数据类型映射,到自动生成 SQL、CDC 事件处理、XA 事务精确一次语义与 Save Mode 写入策略,并提供可直接运行的 HOCON 任务示例。

连接器概述

Jdbc SQLServer Sink是 SeaTunnel 基于 JDBC 协议实现的 SQL Server 写入端连接器,属于connector-jdbc插件体系。它的核心能力是:通过 JDBC 将上游任意 Source(如 JDBC、Kafka、CDC 等)产生的数据写入 SQL Server,同时支持批处理(Batch)与流处理(Streaming)两种运行模式,支持并发写入,并可通过 XA 事务实现精确一次(Exactly-Once)语义。

支持的 SQL Server 版本

  • SQL Server2008或更高版本(文档标注"仅供参考",实际以目标环境验证为准)

支持的引擎

  • Spark
  • Flink
  • SeaTunnel Zeta

连接器特性

  • 精确一次(Exactly-Once),见 connector-v2-features
  • CDC(变更数据捕获),见 connector-v2-features
  • 支持多表写入,见 connector-v2-features
  • 定时刷新(基于batch_interval_ms的时间触发写入)

精确一次语义通过XA 事务保证,因此仅支持启用 XA 事务的数据库。可通过设置is_exactly_once=truemax_retries=0组合启用。

底层工作方式

从源码结构看,connector-jdbc采用"方言(Dialect)+ 类型转换器(TypeConverter)"的可插拔架构:SQL Server 相关实现集中在 seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/dialect/sqlserver/ 目录,包含:

  • SqlServerDialectFactory:通过url.startsWith("jdbc:sqlserver:")识别 SQL Server 连接,注册对应方言实例(见 SqlServerDialectFactory.java);
  • SqlServerDialect:负责 upsert 语句生成、标识符引用、表行数估算、分片查询与 schema 变更 DDL;
  • SqlServerTypeConverter/SqlserverTypeMapper:负责 SQL Server 类型与 SeaTunnel 内部类型之间的双向转换。

环境准备与驱动依赖

SQL Server Sink 依赖 Microsoft 官方 JDBC 驱动mssql-jdbc,驱动 JAR 的放置位置随引擎不同而不同:

引擎驱动 JAR 放置目录
Spark / Flink${SEATUNNEL_HOME}/plugins/
SeaTunnel Zeta${SEATUNNEL_HOME}/lib/

同时在"数据库依赖"一节中,官方文档还给出了另一种放置方式:将 Maven 依赖中的驱动 JAR 复制到$SEATUNNEL_HOME/plugins/jdbc/lib/工作目录,例如:

cp mssql-jdbc-xxx.jar $SEATUNNEL_HOME/plugins/jdbc/lib/

提示:mssql-jdbc驱动请从 Maven 中央仓库获取(com.microsoft.sqlserver:mssql-jdbc),并按上述路径放置后重启任务,否则运行时会抛出找不到驱动类com.microsoft.sqlserver.jdbc.SQLServerDriver的错误。

支持的数据源信息

数据源支持的版本驱动类名URL 格式Maven 依赖
SQL Server支持版本 >= 2008com.microsoft.sqlserver.jdbc.SQLServerDriverjdbc:sqlserver://localhost:1433com.microsoft.sqlserver:mssql-jdbc

URL 中可以通过分号追加连接属性,例如 e2e 测试配置 jdbc_sqlserver_source_to_sink.conf 中使用了:

jdbc:sqlserver://sqlserver;databaseName=master;encrypt=false;

其中databaseName=master指定默认数据库,encrypt=false关闭 TLS 加密(测试环境常用,生产环境建议按安全要求显式配置加密策略)。

数据类型映射

JDBC Sink 在进行字段写入、自动生成建表 SQL 时,需要完成 SQL Server 原生类型与 SeaTunnel 内部类型(SeaTunnelDataType)的双向转换。官方文档给出的映射关系如下:

SQL Server 数据类型SeaTunnel 数据类型
BITBOOLEAN
TINYINT、SMALLINTSHORT
INTEGERINT
BIGINTLONG
DECIMAL、NUMERIC、MONEY、SMALLMONEYDECIMAL((获取指定列的列大小)+1, (获取指定列的小数点右侧的位数))
REALFLOAT
FLOATDOUBLE
CHAR、NCHAR、VARCHAR、NTEXT、NVARCHAR、TEXTSTRING
DATELOCAL_DATE
TIMELOCAL_TIME
DATETIME、DATETIME2、SMALLDATETIME、DATETIMEOFFSETLOCAL_DATE_TIME
TIMESTAMP、BINARY、VARBINARY、IMAGE、UNKNOWN尚未支持

从源码看映射细节

结合 SqlServerTypeConverter.java 的实现,可以补充几个值得注意的细节:

  • FLOAT 精度分支FLOAT在精度<= 24时按REALFLOAT_TYPE)处理,否则按DOUBLE_TYPE处理;SqlserverTypeMapper.java 还会在读取元数据时对float精度 15 统一规整为 53,并对nchar/nvarchar的长度做双字节修正(char 长度 × 2)。
  • DATETIMEOFFSET 的时区语义:文档表格将其归入LOCAL_DATE_TIME,但从源码看,DATETIMEOFFSET实际被映射为OFFSET_DATE_TIME_TYPE(带时区偏移的 LTZ 类型),DATETIME/DATETIME2才映射为LOCAL_DATE_TIME_TYPE
  • DECIMAL 精度上限:源码中定义了MAX_PRECISION = 38MAX_SCALE = 37,转换时若超出上限会自动截断并打印告警日志,避免生成非法 DDL。
  • 二进制类型:文档表格将TIMESTAMP(注意:SQL Server 的TIMESTAMP是行版本号,非时间类型)、BINARYVARBINARYIMAGE标注为"尚未支持",但从源码看它们已被映射为PrimitiveByteArrayType(BYTES)——写入方向的支持情况以实际版本运行结果为准。
  • 反向转换(reconvert,用于自动建表):SeaTunnel 的STRING会转换为NVARCHAR(n)(长度 ≤ 4000)或NVARCHAR(MAX)BYTES转换为VARBINARY(n)VARBINARY(MAX)TIMESTAMP_TZ转换为DATETIMEOFFSET。这解释了为什么在自动生成建表 SQL 时,字符串列在 SQL Server 中通常体现为NVARCHAR

Sink 选项详解

以下为Jdbc SQLServer Sink的完整参数表(默认值与类型定义可对照 JdbcSinkOptions.java 源码核实):

名称类型是否必填默认值描述
urlString-JDBC 连接 URL。示例:jdbc:sqlserver://localhost:1433;databaseName=mydatabase
driverString-JDBC 驱动类名,SQL Server 固定为com.microsoft.sqlserver.jdbc.SQLServerDriver
usernameString-连接实例的用户名
passwordString-连接实例的密码
queryString-使用此 SQL 将上游输入数据写入数据库,如INSERT ...query优先级更高
databaseString-使用此databasetable自动生成 SQL 并写入数据。与query互斥,且优先级更高
tableString-配合database自动生成 SQL 的目标表名。与query互斥,且优先级更高
primary_keysArray-自动生成 SQL 时支持insertdeleteupdate(及 upsert)操作的键列
connection_check_timeout_secInt30用于验证连接完成的数据库操作等待时间(秒)
max_retriesInt0提交失败(executeBatch)的重试次数
batch_sizeInt1000批量写入的缓冲记录数,达到后刷新到数据库;若batch_interval_ms大于 0,超时也会触发刷新
batch_interval_msLong0定时刷新间隔(毫秒)。0表示关闭;大于 0 时每条记录写入时检查间隔,达到后同步刷新
is_exactly_onceBooleanfalse是否启用精确一次语义(使用 XA 事务)。启用时需设置xa_data_source_class_name
generate_sink_sqlBooleanfalse根据目标数据库表自动生成 SQL 语句
xa_data_source_class_nameString-数据库驱动的 XA 数据源类名,SQL Server 为com.microsoft.sqlserver.jdbc.SQLServerXADataSource
max_commit_attemptsInt3事务提交失败的重试次数
transaction_timeout_secInt-1事务打开后的超时时间,-1 表示永不超时。注意:设置超时可能影响精确一次语义
auto_commitBooleantrue默认启用自动事务提交
propertiesMap-额外的 JDBC 连接参数;与 URL 中同参数冲突时,优先级由 SQL Server JDBC 驱动决定
common-options--Sink 插件通用参数,详见 Sink Common Options
schema_save_modeEnumCREATE_SCHEMA_WHEN_NOT_EXIST任务启动前控制目标表结构(Schema)的处理方式
data_save_modeEnumAPPEND_DATA任务启动前控制目标表已有数据的处理方式
custom_sqlString-data_save_modeCUSTOM_PROCESSING时,任务启动前需要执行的 SQL
enable_upsertBooleantrue通过主键启用 upsert;若任务中没有键重复数据,设为false可加快导入速度
multi_table_sink_replicaInt1多表写入时使用的 Sink Writer 副本数量

参数分组解读

连接类参数(必填)urldriver是唯二必填项;usernamepassword按目标库配置提供;connection_check_timeout_sec控制连接校验超时;properties可附加驱动级参数(如字符集、加密选项)。

写入 SQL 控制类:三选一的写入策略——使用query完全自定义 SQL;或使用database+table(可配合generate_sink_sql=true)自动生成 SQL;CDC 场景还需配置primary_keys。注意database/tablequery互斥,且前者优先级更高。

批处理与性能类batch_size(默认 1000 条)、batch_interval_ms(默认 0,关闭定时刷新)共同决定刷新节奏;max_retries控制executeBatch失败重试;enable_upsert在无键重复时可关闭以加速。

精确一次(XA)类is_exactly_once=truexa_data_source_class_name=com.microsoft.sqlserver.jdbc.SQLServerXADataSourcemax_retries=0三者组合启用 XA 事务;max_commit_attempts控制提交阶段重试;transaction_timeout_sec默认为 -1(永不超时),设置超时需评估对精确一次的影响。

写入模式与 Save Mode

generate_sink_sqlqueryschema_save_modedata_save_modecustom_sqlprimary_keysenable_upsert等参数的选择,本质上是在回答两个问题:每一行数据如何写(写入模式),以及写入前如何处理目标端的表和数据(Save Mode)。官方在 Sink 写入模式与 Save Mode 中给出了快速决策表,JDBC 系列 Sink 的规则可归纳如下:

目标优先选择
让 SeaTunnel 自动生成 INSERT / UPSERT / UPDATE / DELETE SQLgenerate_sink_sql = true+database+table(通常还需primary_keys
完全控制写入 SQLquery = "INSERT ... VALUES (?, ...)"(不要与generate_sink_sql=true同时配置)
目标表不存在时自动创建配置schema_save_mode(仅自动生成 SQL 模式生效)
写入前保留 / 清空 / 校验目标数据配置data_save_mode
写入前执行自定义 SQLdata_save_mode = "CUSTOM_PROCESSING"+custom_sql(写入前钩子,非逐行 SQL)
使用数据库原生 Upsertgenerate_sink_sql = true+primary_keys+enable_upsert = true

schema_save_mode取值RECREATE_SCHEMA(存在即删后重建)、CREATE_SCHEMA_WHEN_NOT_EXIST(不存在才创建,默认)、ERROR_WHEN_SCHEMA_NOT_EXIST(不存在则报错)、IGNORE(跳过结构处理)。

data_save_mode取值DROP_DATA(保留结构清空数据)、APPEND_DATA(追加写入,默认)、CUSTOM_PROCESSING(先执行custom_sql)、ERROR_WHEN_DATA_EXISTS(已有数据则报错)。

注意:使用query自定义 SQL 模式时,JDBC Sink 不会执行schema_save_modedata_save_modecustom_sql;这些 Save Mode 仅在自动生成 SQL 且能解析目标 Catalog 表时才生效。

源码视角:SQL Server 方言实现要点

Upsert:基于 MERGE 语句

在自动生成 SQL + 配置primary_keys+enable_upsert=true时,SqlServerDialect.getUpsertStatement() 会生成 SQL Server 原生的MERGE INTO ... USING ... WHEN MATCHED THEN UPDATE ... WHEN NOT MATCHED THEN INSERT语句,以主键作为匹配条件实现 upsert。这也解释了enable_upsert只有拿到可用主键或唯一键后才有意义:没有键时,自动生成的 SQL 会退化为普通 INSERT。

标识符引用与 URL 识别

  • SQL Server 方言使用方括号[]引用标识符(quoteIdentifier),支持schema.tabledatabase.schema.table全限定名;
  • 行数估算与分片查询使用sys.dm_db_partition_stats视图与SELECT TOP (n) ...语法,体现了 SQL Server 特有的 T-SQL 方言(见 SqlServerDialect.java)。

Schema 变更(CDC 场景)

对于 CDC 数据同步中的表结构演进,SQL Server 方言实现了ALTER TABLE ADD / ALTER COLUMN / DROP COLUMN及列重命名(sp_rename)、默认约束管理、列注释(sp_updateextendedproperty)等 DDL 能力(见 SqlServerDialect.java),并针对"向非空表添加无默认值的 NOT NULL 列"这一 SQL Server 限制做了降级处理(先按 NULL 添加,由后续 CDC 事件回填)。e2e 测试 SqlServerSchemaChangeIT.java 对mysqlcdc_to_sqlserver_with_schema_change.conf场景做了覆盖验证。

任务示例

简单示例:SQL Server 到 SQL Server

以下示例从 SQL Server 读取full_types_jdbc表(按id分 10 片并行读取),写入另一张表full_types_jdbc_sink(基于query全字段 INSERT):

env { # 可以在此设置引擎配置 parallelism = 10 } source { # 这是一个示例源插件,**仅用于测试和演示功能** Jdbc { driver = com.microsoft.sqlserver.jdbc.SQLServerDriver url = "jdbc:sqlserver://localhost:1433;databaseName=column_type_test" username = SA password = "Y.sa123456" query = "select * from column_type_test.dbo.full_types_jdbc" # 并行分片读取字段 partition_column = "id" # 分片数量 partition_num = 10 } # 完整源插件列表请参阅仓库 docs/zh/connectors/source 下的 Jdbc 文档 } transform { # 转换插件示例请参阅仓库 docs/zh/transforms 下的相关文档 } sink { Jdbc { driver = com.microsoft.sqlserver.jdbc.SQLServerDriver url = "jdbc:sqlserver://localhost:1433;databaseName=column_type_test" username = SA password = "Y.sa123456" query = "insert into full_types_jdbc_sink( id, val_char, val_varchar, val_text, val_nchar, val_nvarchar, val_ntext, val_decimal, val_numeric, val_float, val_real, val_smallmoney, val_money, val_bit, val_tinyint, val_smallint, val_int, val_bigint, val_date, val_time, val_datetime2, val_datetime, val_smalldatetime ) values( ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? )" } # 完整接收器插件列表请参阅仓库 docs/zh/connectors/sink 下的 Jdbc 文档 }

CDC(变更数据捕获)事件

写入 CDC 变更数据时,需要配置databasetableprimary_keys,由 SeaTunnel 根据 CDC 事件的INSERT / UPDATE / DELETE类型自动生成对应 SQL:

Jdbc { plugin_input = "customers" driver = com.microsoft.sqlserver.jdbc.SQLServerDriver url = "jdbc:sqlserver://localhost:1433;databaseName=column_type_test" username = SA password = "Y.sa123456" generate_sink_sql = true database = "column_type_test" table = "dbo.full_types_sink" batch_size = 100 primary_keys = ["id"] }

plugin_input指定上游数据集名称(当作业存在多个 source / transform / sink 时必须显式指定,详见 Sink Common Options)。

精确一次接收器

事务性写入可能较慢,但数据更准确。通过is_exactly_once=true+xa_data_source_class_name+max_retries=0启用 XA 事务:

Jdbc { driver = com.microsoft.sqlserver.jdbc.SQLServerDriver url = "jdbc:sqlserver://localhost:1433;databaseName=column_type_test" username = SA password = "Y.sa123456" max_retries = 0 query = "insert into full_types_jdbc_sink( id, val_char, val_varchar, val_text, val_nchar, val_nvarchar, val_ntext, val_decimal, val_numeric, val_float, val_real, val_smallmoney, val_money, val_bit, val_tinyint, val_smallint, val_int, val_bigint, val_date, val_time, val_datetime2, val_datetime, val_smalldatetime ) values( ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ? )" is_exactly_once = "true" xa_data_source_class_name = "com.microsoft.sqlserver.jdbc.SQLServerXADataSource" }

基于自动生成 SQL 的最小化配置

仓库 e2e 测试 jdbc_sqlserver_source_to_sink.conf 展示了最简写法:Sink 仅需driverurlusernamepassworddatabasetablegenerate_sink_sql = true,即可完成整表写入:

sink { Jdbc { driver = com.microsoft.sqlserver.jdbc.SQLServerDriver url = "jdbc:sqlserver://sqlserver;databaseName=master;encrypt=false;" username = SA password = "A_Str0ng_Required_Password" database = "master" table = "dbo.sink" generate_sink_sql = true } }

并发与分片提示

  • 如果未设置partition_column,任务将以单并发运行;
  • 如果设置了partition_column,将根据任务的并发度(env.parallelism或插件级parallelism并行执行

该提示同样适用于 Source 侧:简单示例中通过partition_column = "id"+partition_num = 10将读取切分为 10 个分片并行拉取,充分利用 SQL Server 的SELECT TOP ... ORDER BY分片查询能力(见前文方言实现)。

变更日志

本连接器的版本变更记录维护在 connector-jdbc 变更日志 中,升级前建议对照确认各版本行为差异。

【免费下载链接】seatunnelSeaTunnel is a multimodal, high-performance, distributed, massive data integration tool.项目地址: https://gitcode.com/GitHub_Trending/se/seatunnel

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

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

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

立即咨询