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=true与max_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 | 支持版本 >= 2008 | com.microsoft.sqlserver.jdbc.SQLServerDriver | jdbc:sqlserver://localhost:1433 | com.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 数据类型 |
|---|---|
| BIT | BOOLEAN |
| TINYINT、SMALLINT | SHORT |
| INTEGER | INT |
| BIGINT | LONG |
| DECIMAL、NUMERIC、MONEY、SMALLMONEY | DECIMAL((获取指定列的列大小)+1, (获取指定列的小数点右侧的位数)) |
| REAL | FLOAT |
| FLOAT | DOUBLE |
| CHAR、NCHAR、VARCHAR、NTEXT、NVARCHAR、TEXT | STRING |
| DATE | LOCAL_DATE |
| TIME | LOCAL_TIME |
| DATETIME、DATETIME2、SMALLDATETIME、DATETIMEOFFSET | LOCAL_DATE_TIME |
| TIMESTAMP、BINARY、VARBINARY、IMAGE、UNKNOWN | 尚未支持 |
从源码看映射细节
结合 SqlServerTypeConverter.java 的实现,可以补充几个值得注意的细节:
- FLOAT 精度分支:
FLOAT在精度<= 24时按REAL(FLOAT_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 = 38、MAX_SCALE = 37,转换时若超出上限会自动截断并打印告警日志,避免生成非法 DDL。 - 二进制类型:文档表格将
TIMESTAMP(注意:SQL Server 的TIMESTAMP是行版本号,非时间类型)、BINARY、VARBINARY、IMAGE标注为"尚未支持",但从源码看它们已被映射为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 源码核实):
| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
|---|---|---|---|---|
| url | String | 是 | - | JDBC 连接 URL。示例:jdbc:sqlserver://localhost:1433;databaseName=mydatabase |
| driver | String | 是 | - | JDBC 驱动类名,SQL Server 固定为com.microsoft.sqlserver.jdbc.SQLServerDriver |
| username | String | 否 | - | 连接实例的用户名 |
| password | String | 否 | - | 连接实例的密码 |
| query | String | 否 | - | 使用此 SQL 将上游输入数据写入数据库,如INSERT ...。query优先级更高 |
| database | String | 否 | - | 使用此database和table自动生成 SQL 并写入数据。与query互斥,且优先级更高 |
| table | String | 否 | - | 配合database自动生成 SQL 的目标表名。与query互斥,且优先级更高 |
| primary_keys | Array | 否 | - | 自动生成 SQL 时支持insert、delete、update(及 upsert)操作的键列 |
| connection_check_timeout_sec | Int | 否 | 30 | 用于验证连接完成的数据库操作等待时间(秒) |
| max_retries | Int | 否 | 0 | 提交失败(executeBatch)的重试次数 |
| batch_size | Int | 否 | 1000 | 批量写入的缓冲记录数,达到后刷新到数据库;若batch_interval_ms大于 0,超时也会触发刷新 |
| batch_interval_ms | Long | 否 | 0 | 定时刷新间隔(毫秒)。0表示关闭;大于 0 时每条记录写入时检查间隔,达到后同步刷新 |
| is_exactly_once | Boolean | 否 | false | 是否启用精确一次语义(使用 XA 事务)。启用时需设置xa_data_source_class_name |
| generate_sink_sql | Boolean | 否 | false | 根据目标数据库表自动生成 SQL 语句 |
| xa_data_source_class_name | String | 否 | - | 数据库驱动的 XA 数据源类名,SQL Server 为com.microsoft.sqlserver.jdbc.SQLServerXADataSource |
| max_commit_attempts | Int | 否 | 3 | 事务提交失败的重试次数 |
| transaction_timeout_sec | Int | 否 | -1 | 事务打开后的超时时间,-1 表示永不超时。注意:设置超时可能影响精确一次语义 |
| auto_commit | Boolean | 否 | true | 默认启用自动事务提交 |
| properties | Map | 否 | - | 额外的 JDBC 连接参数;与 URL 中同参数冲突时,优先级由 SQL Server JDBC 驱动决定 |
| common-options | - | 否 | - | Sink 插件通用参数,详见 Sink Common Options |
| schema_save_mode | Enum | 否 | CREATE_SCHEMA_WHEN_NOT_EXIST | 任务启动前控制目标表结构(Schema)的处理方式 |
| data_save_mode | Enum | 否 | APPEND_DATA | 任务启动前控制目标表已有数据的处理方式 |
| custom_sql | String | 否 | - | 当data_save_mode为CUSTOM_PROCESSING时,任务启动前需要执行的 SQL |
| enable_upsert | Boolean | 否 | true | 通过主键启用 upsert;若任务中没有键重复数据,设为false可加快导入速度 |
| multi_table_sink_replica | Int | 否 | 1 | 多表写入时使用的 Sink Writer 副本数量 |
参数分组解读
连接类参数(必填):url、driver是唯二必填项;username、password按目标库配置提供;connection_check_timeout_sec控制连接校验超时;properties可附加驱动级参数(如字符集、加密选项)。
写入 SQL 控制类:三选一的写入策略——使用query完全自定义 SQL;或使用database+table(可配合generate_sink_sql=true)自动生成 SQL;CDC 场景还需配置primary_keys。注意database/table与query互斥,且前者优先级更高。
批处理与性能类:batch_size(默认 1000 条)、batch_interval_ms(默认 0,关闭定时刷新)共同决定刷新节奏;max_retries控制executeBatch失败重试;enable_upsert在无键重复时可关闭以加速。
精确一次(XA)类:is_exactly_once=true、xa_data_source_class_name=com.microsoft.sqlserver.jdbc.SQLServerXADataSource、max_retries=0三者组合启用 XA 事务;max_commit_attempts控制提交阶段重试;transaction_timeout_sec默认为 -1(永不超时),设置超时需评估对精确一次的影响。
写入模式与 Save Mode
generate_sink_sql、query、schema_save_mode、data_save_mode、custom_sql、primary_keys、enable_upsert等参数的选择,本质上是在回答两个问题:每一行数据如何写(写入模式),以及写入前如何处理目标端的表和数据(Save Mode)。官方在 Sink 写入模式与 Save Mode 中给出了快速决策表,JDBC 系列 Sink 的规则可归纳如下:
| 目标 | 优先选择 |
|---|---|
| 让 SeaTunnel 自动生成 INSERT / UPSERT / UPDATE / DELETE SQL | generate_sink_sql = true+database+table(通常还需primary_keys) |
| 完全控制写入 SQL | query = "INSERT ... VALUES (?, ...)"(不要与generate_sink_sql=true同时配置) |
| 目标表不存在时自动创建 | 配置schema_save_mode(仅自动生成 SQL 模式生效) |
| 写入前保留 / 清空 / 校验目标数据 | 配置data_save_mode |
| 写入前执行自定义 SQL | data_save_mode = "CUSTOM_PROCESSING"+custom_sql(写入前钩子,非逐行 SQL) |
| 使用数据库原生 Upsert | generate_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_mode、data_save_mode或custom_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.table与database.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 变更数据时,需要配置database、table和primary_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 仅需driver、url、username、password、database、table与generate_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),仅供参考