ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

SeaTunnel JDBC SQL Server Sink 连接器实战指南:配置详解、数据类型映射与精确一次写入 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或更高版本文档标注仅供参考实际以目标环境验证为准支持的引擎SparkFlinkSeaTunnel Zeta连接器特性精确一次Exactly-Once见 connector-v2-featuresCDC变更数据捕获见 connector-v2-features支持多表写入见 connector-v2-features定时刷新基于batch_interval_ms的时间触发写入精确一次语义通过XA 事务保证因此仅支持启用 XA 事务的数据库。可通过设置is_exactly_oncetrue与max_retries0组合启用。底层工作方式从源码结构看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.javaSqlServerDialect负责 upsert 语句生成、标识符引用、表行数估算、分片查询与 schema 变更 DDLSqlServerTypeConverter/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-jdbcURL 中可以通过分号追加连接属性例如 e2e 测试配置 jdbc_sqlserver_source_to_sink.conf 中使用了jdbc:sqlserver://sqlserver;databaseNamemaster;encryptfalse;其中databaseNamemaster指定默认数据库encryptfalse关闭 TLS 加密测试环境常用生产环境建议按安全要求显式配置加密策略。数据类型映射JDBC Sink 在进行字段写入、自动生成建表 SQL 时需要完成 SQL Server 原生类型与 SeaTunnel 内部类型SeaTunnelDataType的双向转换。官方文档给出的映射关系如下SQL Server 数据类型SeaTunnel 数据类型BITBOOLEANTINYINT、SMALLINTSHORTINTEGERINTBIGINTLONGDECIMAL、NUMERIC、MONEY、SMALLMONEYDECIMAL((获取指定列的列大小)1, (获取指定列的小数点右侧的位数))REALFLOATFLOATDOUBLECHAR、NCHAR、VARCHAR、NTEXT、NVARCHAR、TEXTSTRINGDATELOCAL_DATETIMELOCAL_TIMEDATETIME、DATETIME2、SMALLDATETIME、DATETIMEOFFSETLOCAL_DATE_TIMETIMESTAMP、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 38、MAX_SCALE 37转换时若超出上限会自动截断并打印告警日志避免生成非法 DDL。二进制类型文档表格将TIMESTAMP注意SQL Server 的TIMESTAMP是行版本号非时间类型、BINARY、VARBINARY、IMAGE标注为尚未支持但从源码看它们已被映射为PrimitiveByteArrayTypeBYTES——写入方向的支持情况以实际版本运行结果为准。反向转换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;databaseNamemydatabasedriverString是-JDBC 驱动类名SQL Server 固定为com.microsoft.sqlserver.jdbc.SQLServerDriverusernameString否-连接实例的用户名passwordString否-连接实例的密码queryString否-使用此 SQL 将上游输入数据写入数据库如INSERT ...。query优先级更高databaseString否-使用此database和table自动生成 SQL 并写入数据。与query互斥且优先级更高tableString否-配合database自动生成 SQL 的目标表名。与query互斥且优先级更高primary_keysArray否-自动生成 SQL 时支持insert、delete、update及 upsert操作的键列connection_check_timeout_secInt否30用于验证连接完成的数据库操作等待时间秒max_retriesInt否0提交失败executeBatch的重试次数batch_sizeInt否1000批量写入的缓冲记录数达到后刷新到数据库若batch_interval_ms大于 0超时也会触发刷新batch_interval_msLong否0定时刷新间隔毫秒。0表示关闭大于 0 时每条记录写入时检查间隔达到后同步刷新is_exactly_onceBoolean否false是否启用精确一次语义使用 XA 事务。启用时需设置xa_data_source_class_namegenerate_sink_sqlBoolean否false根据目标数据库表自动生成 SQL 语句xa_data_source_class_nameString否-数据库驱动的 XA 数据源类名SQL Server 为com.microsoft.sqlserver.jdbc.SQLServerXADataSourcemax_commit_attemptsInt否3事务提交失败的重试次数transaction_timeout_secInt否-1事务打开后的超时时间-1 表示永不超时。注意设置超时可能影响精确一次语义auto_commitBoolean否true默认启用自动事务提交propertiesMap否-额外的 JDBC 连接参数与 URL 中同参数冲突时优先级由 SQL Server JDBC 驱动决定common-options-否-Sink 插件通用参数详见 Sink Common Optionsschema_save_modeEnum否CREATE_SCHEMA_WHEN_NOT_EXIST任务启动前控制目标表结构Schema的处理方式data_save_modeEnum否APPEND_DATA任务启动前控制目标表已有数据的处理方式custom_sqlString否-当data_save_mode为CUSTOM_PROCESSING时任务启动前需要执行的 SQLenable_upsertBoolean否true通过主键启用 upsert若任务中没有键重复数据设为false可加快导入速度multi_table_sink_replicaInt否1多表写入时使用的 Sink Writer 副本数量参数分组解读连接类参数必填url、driver是唯二必填项username、password按目标库配置提供connection_check_timeout_sec控制连接校验超时properties可附加驱动级参数如字符集、加密选项。写入 SQL 控制类三选一的写入策略——使用query完全自定义 SQL或使用databasetable可配合generate_sink_sqltrue自动生成 SQLCDC 场景还需配置primary_keys。注意database/table与query互斥且前者优先级更高。批处理与性能类batch_size默认 1000 条、batch_interval_ms默认 0关闭定时刷新共同决定刷新节奏max_retries控制executeBatch失败重试enable_upsert在无键重复时可关闭以加速。精确一次XA类is_exactly_oncetrue、xa_data_source_class_namecom.microsoft.sqlserver.jdbc.SQLServerXADataSource、max_retries0三者组合启用 XA 事务max_commit_attempts控制提交阶段重试transaction_timeout_sec默认为 -1永不超时设置超时需评估对精确一次的影响。写入模式与 Save Modegenerate_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 SQLgenerate_sink_sql truedatabasetable通常还需primary_keys完全控制写入 SQLquery INSERT ... VALUES (?, ...)不要与generate_sink_sqltrue同时配置目标表不存在时自动创建配置schema_save_mode仅自动生成 SQL 模式生效写入前保留 / 清空 / 校验目标数据配置data_save_mode写入前执行自定义 SQLdata_save_mode CUSTOM_PROCESSINGcustom_sql写入前钩子非逐行 SQL使用数据库原生 Upsertgenerate_sink_sql trueprimary_keysenable_upsert trueschema_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_keysenable_upserttrue时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全字段 INSERTenv { # 可以在此设置引擎配置 parallelism 10 } source { # 这是一个示例源插件**仅用于测试和演示功能** Jdbc { driver com.microsoft.sqlserver.jdbc.SQLServerDriver url jdbc:sqlserver://localhost:1433;databaseNamecolumn_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;databaseNamecolumn_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类型自动生成对应 SQLJdbc { plugin_input customers driver com.microsoft.sqlserver.jdbc.SQLServerDriver url jdbc:sqlserver://localhost:1433;databaseNamecolumn_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_oncetruexa_data_source_class_namemax_retries0启用 XA 事务Jdbc { driver com.microsoft.sqlserver.jdbc.SQLServerDriver url jdbc:sqlserver://localhost:1433;databaseNamecolumn_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;databaseNamemaster;encryptfalse; 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 idpartition_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),仅供参考
返回列表