Flink CDC SQLServer Connector 实战指南:从数据库 CDC 启用、Flink SQL 建表到增量快照与无主键表采集 📅 发布时间:2026/9/17 7:42:28 👁 浏览次数: Flink CDC SQLServer Connector 实战指南从数据库 CDC 启用、Flink SQL 建表到增量快照与无主键表采集【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdcFlink CDC 的 SQLServer CDC Connector 支持从 SQL Server 数据库同时读取全量快照数据与增量变更数据基于 SQL Server 自身的变更数据捕获 Change Data Capture 机制。本文以当前仓库Flink CDC中的官方文档与源码实现为依据完整讲解如何在 SQL Server 端开启 CDC、在 Flink SQL 中注册 sqlserver-cdc 表、逐项解析全部 Connector 配置参数并深入说明精确一次语义、启动位点、DataStream 编程接口、无主键表支持以及数据类型映射等核心能力。读完本文你可以独立完成SQL Server → Flink的实时数据接入方案并在遇到快照阶段 Checkpoint 超时、无主键表数据一致性等问题时具备排查与规避能力。依赖引入Maven 依赖与 SQL Client JAR接入 SQLServer CDC Connector 有两种方式使用构建工具Maven/SBT 等引入依赖或在 SQL Client 中使用打包好的 SQL JAR。Maven 依赖在基于构建工具的项目中引入flink-connector-sqlserver-cdc依赖可借助仓库docs目录下artifactshortcode 自动生成对应版本坐标flink-connector-sqlserver-cdcSQL Client JAR下载flink-sql-connector-sqlserver-cdc连接器 JAR放入FLINK_HOME/lib/目录后即可在 SQL Client 中直接使用。该连接器的打包实现位于仓库 flink-sql-connector-sqlserver-cdc 模块内含一个仅用于注册连接器标识的入口类用于把依赖重打包成可直接投放的 SQL 连接器 JAR。更多发布版本可到 Maven 中央仓库检索。准备 SQL Server 数据库开启变更数据捕获CDCSQL Server 的 CDC 能力是连接器读取增量变更的基础。需要由SQL Server 管理员在想要捕获的源表上开启 CDC且数据库本身必须已经启用 CDC。前置条件目标 SQL Server 数据库已启用 CDCSQL Server Agent 服务正在运行执行操作的用户是目标数据库db_owner固定数据库角色的成员。开启表级 CDC通过数据库管理工具连接到 SQL Server 数据库后执行以下 SQL 语句为表开启 CDCUSE MyDB GO EXEC sys.sp_cdc_enable_table source_schema Ndbo, -- 指定源表所属的 schema source_name NMyTable, -- 指定要捕获的源表名称 role_name NMyRole, -- 指定角色 MyRole可向该角色添加用户并授予其对源表捕获列的 SELECT 权限 -- sysadmin 或 db_owner 角色的用户也可访问指定的更改表 -- 将该值设为 NULL 则仅允许 sysadmin 或 db_owner 成员完全访问捕获信息 filegroup_name NMyDB_CT,-- 指定 SQL Server 放置更改表change table的文件组 -- 该文件组必须已存在建议不要将更改表放在源表所在的文件组中 supports_net_changes 0 GO各参数含义如下参数说明source_schema源表所属的 schema 名如dbosource_name要捕获的源表名role_name授予捕获列 SELECT 权限的角色NULL表示仅 sysadmin/db_owner 可完全访问filegroup_name更改表存放的文件组必须已存在建议与源表文件组分离supports_net_changes是否支持净更改net changes查询验证用户对 CDC 表的访问权限开启后可通过存储过程sys.sp_cdc_help_change_data_capture验证-- 以下示例在数据库 MyDB 上运行存储过程 sys.sp_cdc_help_change_data_capture USE MyDB; GO EXEC sys.sp_cdc_help_change_data_capture GO该查询返回数据库中每个已启用 CDC、且调用者有权访问其更改数据的表的配置信息。如果结果为空请确认用户对捕获实例capture instance和 CDC 表均具备访问权限。在 Flink SQL 中创建 SQLServer CDC 表开启数据库 CDC 后即可在 Flink SQL 中按如下方式注册一张orders表-- register a SqlServer table orders in Flink SQL CREATE TABLE orders ( id INT, order_date DATE, purchaser INT, quantity INT, product_id INT, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector sqlserver-cdc, hostname localhost, port 1433, username sa, password Password!, database-name inventory, table-name dob.orders ); -- read snapshot and binlogs from orders table SELECT * FROM orders;执行SELECT后Flink 会先读取orders表的全量快照随后持续消费增量变更。连接器标识sqlserver-cdc在源码 SqlServerTableFactory.java 中通过factoryIdentifier()注册hostname、username、password、database-name、table-name在 requiredOptions() 中被声明为必填项。Connector 配置参数详解下表完整列出 SQLServer CDC Connector 的配置项来自官方文档并给出默认值与含义OptionRequiredDefaultTypeDescriptionconnectorrequired(none)String指定使用的连接器此处应为sqlserver-cdchostnamerequired(none)StringSQL Server 数据库的 IP 地址或主机名usernamerequired(none)String连接 SQL Server 数据库的用户名passwordrequired(none)String连接 SQL Server 数据库的密码database-namerequired(none)String要监控的 SQL Server 数据库名table-namerequired(none)String要监控的 SQL Server 表名格式如db1.table1portoptional1433IntegerSQL Server 数据库的端口号server-time-zoneoptionalUTCString数据库服务器会话时区如Asia/Shanghaiscan.incremental.snapshot.enabledoptionaltrueBoolean是否启用并行快照incremental snapshotchunk-meta.group.sizeoptional1000Integerchunk 元数据的分组大小元数据规模超过该值时会被分成多组chunk-key.even-distribution.factor.lower-boundoptional0.05dDoublechunk key 分布因子下界。分布因子用于判断表数据是否均匀分布数据均匀时走均匀切分优化不均匀时走查询式切分。分布因子计算方式为(MAX(id) - MIN(id) 1) / rowCountchunk-key.even-distribution.factor.upper-boundoptional1000.0dDoublechunk key 分布因子上界含义同上debezium.*optional(none)String透传给 Debezium Embedded Engine 的属性用于捕获 SQL Server 变更例如debezium.snapshot.mode initial_only详见 Debezium SQLServer Connector 配置文档scan.incremental.close-idle-reader.enabledoptionalfalseBoolean是否在快照阶段结束时关闭空闲 reader。当execution.checkpointing.checkpoints-after-tasks-finish.enabled设为 true 时要求 Flink 版本 ≥ 1.14Flink ≥ 1.15 时该配置默认即为 true无需显式设置scan.incremental.snapshot.chunk.key-columnoptional(none)String表快照的 chunk key。读取快照时按 chunk key 将表拆分为多个 chunk默认取主键第一列也可使用非主键列但可能导致查询性能下降。警告使用非主键列作为 chunk key 可能造成数据不一致详见下文无主键表的支持scan.incremental.snapshot.unbounded-chunk-first.enabledoptionaltrueBoolean快照读取阶段是否优先分配无界 chunkunbounded chunk。优先分配有助于降低对最大无界 chunk 做快照时 TaskManager 发生 OOM 的风险scan.incremental.snapshot.backfill.skipoptionalfalseBoolean快照读取阶段是否跳过 backfill。跳过时快照期间的变更会在后续变更日志读取阶段消费而非合入快照。警告跳过 backfill 可能造成数据不一致因为快照阶段内的部分变更日志事件会被重放仅保证 at-least-once 语义例如对快照中已更新的值再次更新、对已删除的记录再次删除这些重放事件需特殊处理参数解析与校验的源码印证以上参数并非只是文档描述均可在 SqlServerTableFactory.java 中找到对应ConfigOption定义port默认值 1433、server-time-zone默认值UTC与文档一致工厂类还会做如下运行时校验开启并行快照scan.incremental.snapshot.enabledtrue时scan.incremental.snapshot.chunk-size、scan.snapshot.fetch.size、chunk-meta.group.size、connection.pool.size必须大于 1connect.max-retries必须大于 0分布因子下界chunk-key.even-distribution.factor.lower-bound必须落在[0, 1]区间上界必须 ≥ 1.0。除文档表格中的参数外工厂的 optionalOptions() 还声明了一批继承自 JDBC 基础连接器的常用项包括scan.incremental.snapshot.chunk-size快照单 chunk 大小、scan.snapshot.fetch.size快照查询抓取行数、connect.timeout、connect.max-retries、connection.pool.size以及scan.startup.mode/scan.startup.timestamp-millis见下文启动读取位点。可用的元数据字段SQLServer CDC 连接器可以把如下元数据以只读VIRTUAL列形式暴露在表定义中KeyDataTypeDescriptiontable_nameSTRING NOT NULL包含该行数据的表名schema_nameSTRING NOT NULL包含该行数据的 schema 名database_nameSTRING NOT NULL包含该行数据的数据库名op_tsTIMESTAMP_LTZ(3) NOT NULL该变更在数据库中被执行的时间若记录来自表快照而非变更流则该值恒为 0这些元数据的实现定义在 SqlServerReadableMetadata.java其取值直接解析 DebeziumSourceRecord的 source 结构体如table、schema、db、ts_ms字段从源码可以确认op_ts通过TimestampData.fromEpochMillis(...)把数据库时间戳毫秒值转为TIMESTAMP_LTZ(3)。暴露元数据字段的扩展 CREATE TABLE 示例CREATE TABLE products ( table_name STRING METADATA FROM table_name VIRTUAL, schema_name STRING METADATA FROM schema_name VIRTUAL, db_name STRING METADATA FROM database_name VIRTUAL, operation_ts TIMESTAMP_LTZ(3) METADATA FROM op_ts VIRTUAL, id INT NOT NULL, name STRING, description STRING, weight DECIMAL(10,3) ) WITH ( connector sqlserver-cdc, hostname localhost, port 1433, username sa, password Password!, database-name inventory, table-name dbo.products );限制扫描表快照期间无法进行 Checkpoint注意该限制仅适用于未启用增量快照框架即scan.incremental.snapshot.enabled设为false的情况。在扫描数据库表快照期间由于没有可恢复的位点recoverable position无法执行 Checkpoint。此时 SQLServer CDC Source 会让 Checkpoint 一直等到超时超时后的 Checkpoint 会被判定为失败而默认情况下失败 Checkpoint 会触发 Flink 作业的 failover。因此如果数据库表很大建议添加如下 Flink 配置以避免因 Checkpoint 超时导致的作业重启execution.checkpointing.interval: 10min execution.checkpointing.tolerable-failed-checkpoints: 100 restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 2147483647而当scan.incremental.snapshot.enabledtrue默认开启时连接器走增量快照框架快照被切成多个带位点的 chunk可以正常支持 Checkpoint 与并行读取。核心特性精确一次处理Exactly-Once ProcessingSQLServer CDC 连接器是一个 Flink Source先读取数据库快照之后持续读取变更事件即使发生故障也保证精确一次exactly-once处理语义。其工作方式基于 Debezium SQLServer Connector 的底层机制快照 日志抓取 LSN 位点追踪。启动读取位点scan.startup.mode配置项scan.startup.mode指定 SQLServer CDC 消费者的启动模式有效枚举值如下initial默认对捕获表的结构和数据做全量快照适合需要完整初始化数据的场景latest-offset仅对捕获表的结构做快照适合只关心从现在起发生的新变更的场景。注意scan.startup.mode的机制依赖 Debezium 的snapshot.mode配置因此两者不要同时使用。如果在表 DDL 中同时指定了scan.startup.mode与debezium.snapshot.mode可能导致scan.startup.mode不生效。从源码 SqlServerSourceConfigFactory.java 可以看到两者的映射关系INITIAL → snapshot.modeinitial、LATEST_OFFSET → snapshot.modeschema_only、SNAPSHOT → snapshot.modeinitial_only并且该工厂会把database.history设为 Flink 的EmbeddedFlinkDatabaseHistory状态内嵌的数据库历史支持 Flink Checkpoint 恢复使用 Microsoft JDBC 驱动com.microsoft.sqlserver.jdbc.SQLServerDriver建立连接。此外工厂还支持一种文档未在参数表中列出的timestamp模式scan.startup.mode timestamp配合scan.startup.timestamp-millis从指定时间戳映射出的 LSN 开始消费。需要注意该模式仅在启用增量快照时可用——validateStartupOptions() 会在timestamp模式且scan.incremental.snapshot.enabledfalse时抛出校验异常因为传统非并行 Source 不支持时间戳启动。单线程读取SQL Server CDC Source不支持并行读取变更事件因为同一时刻只有一个 task 能接收变更事件SQL Server CDC 表结构的变更流读取本质上是单通道的。在 DataStream 编程中如需保证消息顺序应将下游 Sink 的并行度设为 1。DataStream Source 编程接口SQLServer CDC 连接器同样可以作为 DataStream Source 使用。2.4.0 之前使用传统的SourceFunction风格import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.streaming.api.functions.source.SourceFunction; import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema; import org.apache.flink.cdc.connectors.sqlserver.SqlServerSource; public class SqlServerSourceExample { public static void main(String[] args) throws Exception { SourceFunctionString sourceFunction SqlServerSource.Stringbuilder() .hostname(localhost) .port(1433) .database(sqlserver) // monitor sqlserver database .tableList(dbo.products) // monitor products table .username(sa) .password(Password!) .deserializer(new JsonDebeziumDeserializationSchema()) // converts SourceRecord to JSON String .build(); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env .addSource(sourceFunction) .print().setParallelism(1); // use parallelism 1 for sink to keep message ordering env.execute(); } }2.4.0 之后则推荐使用支持并行读取与增量快照的SqlServerSourceBuilder实现位于 SqlServerSourceBuilder.java构建出的SqlServerIncrementalSource继承自 JDBC 增量 Source 框架import org.apache.flink.api.common.eventtime.WatermarkStrategy; import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.cdc.connectors.base.options.StartupOptions; import org.apache.flink.cdc.connectors.sqlserver.source.SqlServerSourceBuilder; import org.apache.flink.cdc.connectors.sqlserver.source.SqlServerSourceBuilder.SqlServerIncrementalSource; import org.apache.flink.cdc.debezium.JsonDebeziumDeserializationSchema; public class SqlServerIncrementalSourceExample { public static void main(String[] args) throws Exception { SqlServerIncrementalSourceString sqlServerSource new SqlServerSourceBuilder() .hostname(localhost) .port(1433) .databaseList(inventory) .tableList(dbo.products) .username(sa) .password(Password!) .deserializer(new JsonDebeziumDeserializationSchema()) .startupOptions(StartupOptions.initial()) .build(); StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); // enable checkpoint env.enableCheckpointing(3000); // set the source parallelism to 2 env.fromSource( sqlServerSource, WatermarkStrategy.noWatermarks(), SqlServerIncrementalSource) .setParallelism(2) .print() .setParallelism(1); env.execute(Print SqlServer Snapshot Change Stream); } }该示例把 Source 并行度设为 2增量快照框架下多个 subtask 并行读快照 chunk并开启 3 秒间隔的 CheckpointtableList参数接受schemaName.tableName形式的全限定表标识正则。无主键表的支持3.4.0从 3.4.0 版本开始SQLServer CDC 支持没有主键的表。使用无主键表时必须配置scan.incremental.snapshot.chunk.key-column并指定一个非空字段。有两个注意点如果表上存在索引尽量在scan.incremental.snapshot.chunk.key-column中使用索引中包含的列这能提升SELECT语句的执行速度无主键 SQLServer CDC 表的处理语义由scan.incremental.snapshot.chunk.key-column指定列的行为决定如果对该指定列不执行更新操作则保证 exactly-once 语义如果对该指定列执行了更新操作则仅保证 at-least-once 语义不过可以在下游自行指定主键并做幂等操作来保证数据正确性。警告使用非主键列作为 chunk key 可能造成数据不一致对有主键的 SQLServer 表使用非主键列作为scan.incremental.snapshot.chunk.key-column可能造成数据不一致官方文档给出了具体场景表结构主键为idchunk key 列设为pid非主键快照切分Split 01 pid 3Split 13 pid 5操作两个不同的 subtask 并发读取 Split 0 和 Split 1在读取过程中一条更新操作把id0记录的pid从2改为4且该更新发生在两个 split 的低水位与高水位之间结果Split 0 中包含了记录[id0, pid2]Split 1 中包含了记录[id0, pid4]由于这两条记录的处理顺序无法保证id0的pid最终值可能是2也可能是4从而产生数据不一致。因此建议优先使用主键列作为 chunk key若确需使用非主键列需自行接受并处理上述竞态。可用的 Source 指标连接器暴露如下 Flink metrics指标名称为 Flink 官方命名此处仅给出含义可用于观测快照与流读取进度GroupNameTypeDescriptionnamespace.schema.tableisSnapshottingGauge该表是否正在做快照namespace.schema.tableisStreamReadingGauge该表是否正在读取变更流namespace.schema.tablenumTablesSnapshottedGauge已完成快照的表数量namespace.schema.tablenumTablesRemainingGauge尚未完成快照的表数量namespace.schema.tablenumSnapshotSplitsProcessedGauge正在处理的 split 数量namespace.schema.tablenumSnapshotSplitsRemainingGauge尚未处理的 split 数量namespace.schema.tablenumSnapshotSplitsFinishedGauge已处理完成的 split 数量namespace.schema.tablesnapshotStartTimeGauge快照开始时间namespace.schema.tablesnapshotEndTimeGauge快照结束时间注意事项指标组名为namespace.schema.table其中namespace为实际数据库名schema为实际 schema 名table为实际表名对 SQL Server 而言组名形如test_database.test_schema.test_table。数据类型映射SQLServer CDC Connector 将 SQL Server 类型映射为 Flink SQL 类型完整映射关系如下SQLServer typeFlink SQL typechar(n)CHAR(n)varchar(n)、nvarchar(n)、nchar(n)VARCHAR(n)text、ntext、xmlSTRINGdecimal(p, s)、money、smallmoneyDECIMAL(p, s)numeric(p, s)DECIMAL(p, s)float、realDOUBLEbitBOOLEANintINTtinyintSMALLINTsmallintSMALLINTbigintBIGINTdateDATEtime(n)TIME(n)datetime2、datetime、smalldatetimeTIMESTAMP(n)datetimeoffsetTIMESTAMP_LTZ(3)类型转换的底层实现位于 SqlServerTypeUtils.java它根据 JDBCTypes常量把 DebeziumColumn映射为 FlinkDataType并依据列的optional属性决定是否追加NOT NULL约束对于无法识别的类型会抛出Dont support SqlSever type ... yet异常。server-time-zone配置项会参与TIMESTAMP相关类型的会话时区换算当业务库与 Flink 所在时区不一致时如部署在Asia/Shanghai应显式设置该参数避免时间偏移。小结SQLServer CDC Connector 是 Flink CDC 生态中接入微软 SQL Server 数据源的标准方案其能力闭环包括数据库端sp_cdc_enable_table开启捕获 → Flink SQLCREATE TABLEconnector sqlserver-cdc→ 全量快照 增量变更的连续读取 → 可选元数据列op_ts等与 Source 指标观测 → 下游任意 Sink 消费。生产使用时的关键决策点集中在三处并行快照开关scan.incremental.snapshot.enabled默认开启支持并行与 Checkpoint、chunk key 的选择默认主键第一列无主键表必须显式指定且关注 at-least-once 语义、启动位点scan.startup.mode默认initialtimestamp模式依赖增量快照框架。本文涉及的源码均可继续在仓库对应模块中深入研读例如 flink-connector-sqlserver-cdc 的source、table与dialect包。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考