Flink CDC 2.3.0实战:SQL Server实时同步到MySQL全指南

📅 发布时间:2026/10/11 0:34:58
Flink CDC 2.3.0实战:SQL Server实时同步到MySQL全指南
简介面向大数据实时同步场景这套工程基于 flink-connector-sqlserver-cdc 2.3.0提供从 SQL Server 到 MySQL 的完整 CDC 同步示例适合具备 Flink 基础、正在搭建实时数仓或跨库同步链路的开发者。包体共 21 个文件压缩后约 32KB主要包含 11 个 Java 源文件、5 个 properties 配置、2 个 XML 配置、2 个 Markdown 说明与 1 份许可证Java 代码覆盖连接器参数构建、源表与目标表映射及数据写入逻辑properties 与 XML 便于快速替换连接器、JDBC 连接参数README 中英文文档提供导入与启动指引。项目以可运行 Maven 结构呈现核心涵盖依赖引入、CREATE TABLE WITH 配置、INSERT INTO 数据流转等关键环节可直接用于扩展多表同步或调整并行度也可作为 CDC 连接器与 SQL Server、MySQL 版本兼容性排查的参考。当前已有 2314 人学习适合需要实时数据入仓、跨库同步或升级 Flink CDC 的开发者下载对照。1. 为什么把 SQL Server 实时同步到 MySQL我用了 flink-connector-sqlserver-cdc 2.3.0业务库落在 SQL Server报表库和下游数仓落在 MySQL 的场景在很多做数据服务的团队里都遇到过。凌晨跑批同步越来越赶不上业务要数的节奏最直接的诉求就是让 SQL Server 的变更在秒级内落到 MySQL 里。我最后选的是 flink-connector-sqlserver-cdc 2.3.0 这套连接器它把全量快照和增量日志拉成一条流用 Flink SQL 几十行就能把同步任务跑起来不需要自己维护采集程序也不需要碰 SQL Server 的高可用配置。这套方案适合要搭实时数仓、替换手工导数、或者给报表系统做准实时同步的从业者前提是你愿意接受“同步链路里多了一套 Flink 集群”这个事实。2. 版本配对与 CDC 开启SQL Server 侧先立住再谈同步2.1 连接器版本、Flink 版本与驱动包怎么配对Flink CDC 2.3.0 里的 SQL Server 连接器官方发布包里有带 sql 字样的 fat jar名字是 flink-sql-connector-sqlserver-cdc-2.3.0.jar。常见做法是把这个 jar 直接放到 FLINK_HOME/lib 下再放一份 MySQL 的 JDBC 驱动因为 Flink 官方的 lib 目录里默认没有 mysql-connector。这里踩过一个大坑只放了连接器 jar忘了放 MySQL 驱动任务在提交时报 ClassNotFound指向 com.mysql.cj.jdbc.Driver查了半天才发现是 lib 目录缺东西。版本配对建议按“Flink 小版本 连接器小版本”一起锁。Flink CDC 2.3.0 对应 Flink 1.13 到 1.15 这几个小版本都能跑但我一般会锁在 Flink 1.14.6 上组件少、社区案例多、排错资料好找。连接器内部依赖 Debezium 的 SQL Server 插件fat jar 已经把 Debezium 相关 class 打进去了不需要再单独引。MySQL 驱动建议用 8.0.27 以上版本配合 MySQL 5.7 和 8.0 都行驱动版本太老会碰上 caching_sha2_password 认证不兼容任务起不来。jar 包放好之后先别急着写 SQL。你还需要确认 SQL Server 侧有没有开 CDC没开的话任务一启动就会报“找不到捕获实例”之类的错误根本进不到读取日志的阶段。2.2 在 SQL Server 上开启 CDC 的完整 SQL 模板SQL Server 的 CDC 分两级数据库级和表级两个都得开。数据库级开启后SQL Server 会创建一个名为 cdc 的架构并启动捕获作业和清理作业表级开启后才会真正记录变更。下面这段是标准模板我在 2012、2014、2019 三个版本上都跑过USE trade_db; GO -- 数据库级开启 CDC需要 sysadmin 或 db_owner 权限 EXEC sys.sp_cdc_enable_db; GO -- 表级开启 CDCorders 表不限制访问角色 EXEC sys.sp_cdc_enable_table source_schema Ndbo, source_name Norders, role_name NULL, capture_instance Ndbo_orders; GO -- 确认数据库和表都已经开启 SELECT name, is_cdc_enabled FROM sys.databases WHERE name trade_db; SELECT name, is_tracked_by_cdc FROM sys.tables WHERE name orders; GO这段脚本的逻辑分三步先开库再开表最后查状态。role_name 传 NULL 表示不限制 CDC 访问角色生产环境建议传一个具体角色名比如 gating_role否则所有有权限的账号都能查变更数据。capture_instance 是捕获实例名默认是 架构_表名后续排查和重建捕获实例时都要用到它建议按 dbo_orders 这种格式显式指定比默认生成的乱序名好认。开启 CDC 之后SQL Server Agent 必须处于启动状态。捕获作业和清理作业都挂在 Agent 下面Agent 没启动事务日志里的变更不会被扫进 cdc 表Flink 端会一直卡在等待状态从外表看像是连接断了。很多新手在这块翻车逻辑上以为开完 CDC 就完事了实际上 Agent 服务才是真正干活的进程。2.3 开启 CDC 前先排查表和库的三个条件不是每张表都适合直接开 CDC。我通常先过三关三关都过了才上连接器。第一关表必须有主键。Flink CDC 2.3.0 在快照阶段依赖主键做分片和断点记录无主键表虽然能开 CDC但全量快照时容易出问题恢复位点也不可靠。SQL Server 这边如果表只有唯一索引没有主键我建议先补一个主键再开 CDC。业务上不允许加主键的就走 latest-offset 模式只取增量放弃全量历史数据。第二关数据库的 recovery model 必须是 full 或 bulk_logged。simple 模式下事务日志会被周期性截断CDC 捕获作业还没来得及扫到的变更就没了Flink 端会报 LSN 断档。这块很多 DBA 会忽略因为他们平时习惯把业务库设成 simple 模式省日志空间。第三关估算表的数据量和变更频率。超过两亿行的表第一次全量快照会把 SQL Server 的读 IO 打满常见做法是错峰做初始快照或者先把这张表从同步任务里摘出去用一次性初始化脚本导入 MySQL再让 Flink 从 latest-offset 追增量。SQL Server 的 writelog 和日志增长在这类大表上会非常明显后面第 5 章会专门讲这个现象。3. Flink SQL 建表与提交把 SQL Server 变更变成一条实时流3.1 source 表定义sqlserver-cdc 连接器的核心参数连接器、驱动和 SQL Server 侧的 CDC 都准备好之后打开 Flink SQL Client 或者写一个 .sql 脚本文件。source 表定义是整个同步任务的入口字段类型要和源表对齐主键声明用 NOT ENFORCED 让 Flink 只做逻辑识别不强制在 Flink 内部建约束CREATE TABLE orders_source ( id INT NOT NULL, order_no VARCHAR(32), user_id BIGINT, total_amount DECIMAL(10,2), status TINYINT, created_at TIMESTAMP(3), updated_at TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector sqlserver-cdc, hostname 192.168.1.21, port 1433, username cdc_user, password your_password, database-name trade_db, schema-name dbo, table-name orders, server-time-zone Asia/Shanghai, scan.startup.mode initial, scan.incremental.snapshot.chunk.size 4096 );这段建表 SQL 的参数分三组连接组、表定位组、行为组。连接组的 hostname、port、username、password 不用多说账号建议单独建一个只读的 CDC 账号权限给 db_datareader 加上对 cdc schema 的 select不要直接用 sa。表定位组的 database-name、schema-name、table-name 要精确匹配table-name 支持正则比如 orders.* 可以把 orders 开头的表都拉进来但实际生产里我很少这么用一张表一条任务出了问题好排查。行为组里最值得花时间的是 scan.startup.mode两个值要分清initial 表示先做全量快照再无缝接增量适合第一次跑latest-offset 表示只从当前日志位点开始读适合已经有初始化数据的场景。scan.incremental.snapshot.chunk.size 默认是 8096 行这是一个玄学参数调小到 4096 或 2048 能降低对源库的瞬时压力但会拉长快照时间调大则快照快但容易造成源库 IO 抖动。大表建任务时我会先按 4096 跑观察源库的 waits 再决定要不要改别一上来就 20000。3.2 sink 表定义写 MySQL 的 JDBC 参数怎么设source 建好后建 sink。同步到 MySQL 不需要专门的 mysql-cdc 连接器做目标端Flink 官方提供的 JDBC connector 就够用。sink 表字段要和 MySQL 表结构对齐类型别照抄 source第 4 章会展开类型问题。这里先给一个最小可跑通的建表语句CREATE TABLE orders_sink ( id INT NOT NULL, order_no VARCHAR(32), user_id BIGINT, total_amount DECIMAL(10,2), status TINYINT, created_at TIMESTAMP(3) NULL, updated_at TIMESTAMP(3) NULL, PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://192.168.1.30:3306/ods_db?useSSLfalserewriteBatchedStatementstrueallowPublicKeyRetrievaltrue, table-name orders, username sync_user, password your_password, sink.buffer-flush.max-rows 1000, sink.buffer-flush.interval 1s, sink.max-retries 3 );这里的 url 里有三个参数值得记住。useSSLfalse 是避免 MySQL 8.0 默认要求 SSL 握手造成不必要的连接延迟rewriteBatchedStatementstrue 是让 MySQL 驱动把多条 INSERT 合并成一条多值 INSERT吞吐能差出好几倍allowPublicKeyRetrievaltrue 是给 caching_sha2_password 认证用的缺了它连接会报公共密钥检索错误。sink.buffer-flush.max-rows 和 sink.buffer-flush.interval 控制攒批写库的节奏。1000 行或 1 秒哪个先到就触发写入这是生产里用得比较顺手的组合。调成 100 行会很慢调成 10000 行会在任务重启时丢更多未 flush 的数据。sink.max-retries 默认 3MySQL 侧死锁频繁时重试能救一部分临时故障但救不了结构性冲突比如源表和目标表主键不一致导致的重复写。3.3 提交任务checkpoint 与 INSERT INTO 的配合source 和 sink 都建好后一行 INSERT INTO 就把整条链路串起来了。任务提交前checkpoint 必须开。Flink CDC 连接器的增量阶段是靠 checkpoint 记录消费位点的不开 checkpoint任务重启后会从 latest-offset 重新开始增量数据直接丢一大段。-- 脚本开头建议先设置 checkpoint SET execution.checkpointing.interval 10s; SET execution.checkpointing.mode EXACTLY_ONCE; SET parallelism.default 2; INSERT INTO orders_sink SELECT id, order_no, user_id, total_amount, status, created_at, updated_at FROM orders_source;checkpoint interval 我一般从 10s 起步。太短比如 3sSQL Server 端会频繁请求 LSN源库压力大太长比如 60s任务重启后追丢的数据窗口也大。EXACTLY_ONCE 模式在 checkpoint 对齐时会暂停 source 读取对延迟敏感的场景可以先用 AT_LEAST_ONCE 跑通再切救火阶段这是个常见的妥协手段。parallelism.default 建议从 2 开始source 和 sink 都是单并行度也能跑但 MySQL 写入侧受限于单连接吞吐很难上去。提交命令用 SQL Client 跑脚本文件是最省事的cd $FLINK_HOME ./bin/sql-client.sh -f /data/flink_jobs/sync_orders.sql任务提交后第一件事不是看 MySQL 有没有数据而是去 Flink UI 上看 Checkpoint 是不是 10s 一个稳定完成。SQL Server CDC 是流式任务source 端有没有数据进来从 UI 的 Received Records 曲线一眼就能看出来曲线是平的说明连接器没读到变更先查 SQL Server Agent 和 CDC 捕获作业别急着调 Flink 参数。4. 类型映射与字段雷区SQL Server 到 MySQL 的六个翻车点4.1 SQL Server 到 MySQL 的类型映射速查表Flink SQL 在 source 和 sink 中间充当了翻译层但翻译不是无损的。下面这张表是我做同步前必核对一遍的映射清单照着它建 MySQL 表基本不会出结构性错误SQL Server 类型Flink SQL 类型MySQL 目标类型说明intINTINT直接对应bigintBIGINTBIGINT注意 MySQL 8.0 的 int 最大到 2^31-1smallintSMALLINTSMALLINT直接对应tinyintTINYINTTINYINT注意 Flink 的 TINYINT 是 -128 到 127bitBOOLEANTINYINT(1)MySQL 没有原生 booleandecimal(p,s)DECIMAL(p,s)DECIMAL(p,s)精度必须一致否则四舍五入差异numericDECIMALDECIMALnumeric 是 decimal 的同义词moneyDECIMAL(19,4)DECIMAL(19,4)不要映射成 DOUBLEfloatDOUBLEDOUBLEFlink 没有 float映射为 doublerealFLOATFLOAT4 字节浮点char / ncharCHAR(n)CHAR(n)按字符数不是字节数varchar / nvarcharVARCHAR(n)VARCHAR(n)注意 nvarchar 的 n 是字符数MySQL 端同样按字符数建text / ntextSTRINGLONGTEXT大文本统一走 longtextdateDATEDATE直接对应timeTIMETIME直接对应datetimeTIMESTAMP(3)DATETIME(3)精度按毫秒datetime2TIMESTAMP(3)DATETIME(3)高精度会截断见 4.2smalldatetimeTIMESTAMP(0)DATETIME分钟精度datetimeoffsetSTRINGVARCHAR(40)见 4.2uniqueidentifierSTRINGCHAR(36)见 4.2binary / varbinaryBYTESVARBINARY(n)二进制按字节imageBYTESLONGBLOB大二进制这张表的核心原则是“能精确就精确别图省事全用 STRING”。全用 STRING 虽然能跑通但下游做聚合和排序时全得 cast性能损耗大而且 Flink 里 STRING 排序和 MySQL 里 VARCHAR 排序规则是有差异的数据对不上容易产生“看起来一致、实际不一致”的假象。4.2 datetime2、datetimeoffset、uniqueidentifier 的映射细节三个最容易翻车的类型单独拿出来说。datetime2 的精度可以到 7 位小数也就是 100 纳秒而 Flink 的 TIMESTAMP(3) 只能到毫秒。源表如果用了 datetime2(7)同步到 MySQL 后精度直接丢四位业务上对时间敏感的场景会炸。我一般会先在源库查一下 datetime2 列的真实精度如果是 datetime2(7) 且业务无法妥协目标表用 DATETIME(6) 并配 Flink 的 TIMESTAMP(6) 来做虽然还丢一位但影响小很多。这是标准解法里最贴近业务实际的一种。datetimeoffset 更麻烦它带时区偏移比如 12:00:00 08:00。Flink 的 TIMESTAMP 类型不存时区信息强行映射会把偏移丢掉数据到了 MySQL 就变成无时区标记的墙钟时间。常见做法是直接映射成 STRING在源库查出来的格式类似 2024-01-01 12:00:00.1234567 08:00原样存进 MySQL 的 VARCHAR(40)下游要算时间再去解析。这样做放弃了时间字段的索引能力但保证了信息完整比乱转时区导致数据错误强得多。uniqueidentifier 是 SQL Server 的 GUID 类型映射到 MySQL 没有原生对应。有人喜欢转成 CHAR(36) 保格式有人转成 BINARY(16) 省空间。我默认选 CHAR(36)因为下游排查问题时要拿这串 ID 去对单可读性比那十几字节的空间值钱。Flink 侧映射为 STRING写进 MySQL 的 CHAR(36)。4.3 主键、排序规则与大小写敏感的隐性坑主键坑往往不在建表阶段暴露而在跑了一段时间后以主键冲突的形式出现。SQL Server 的主键如果带自增标识MySQL 侧的主键一定要建 AUTO_INCREMENT且初始值要大于源表当前最大值否则第一批数据写进来就会撞键。源表主键是 bigint 时MySQL 侧不能用 int会溢出。之前接手过一个任务源端主键最大到了 32 亿MySQL 端用的 int同步到一半任务直接报 duplicated entry查了半天才反应过来是类型写小了。排序规则这个坑很阴。SQL Server 默认排序规则常见的是 SQL_Latin1_General_CP1_CI_AS大小写不敏感MySQL 端如果建表用了 utf8mb4_bin等价条件下的字符串匹配结果就会不同。比如源表有个 status 字段存 Success 和 successSQL Server 认为两者相同唯一索引不会挡住MySQL 用 bin 排序规则会当成两条不同记录sink 端 upsert 逻辑就可能乱。我的一般做法是源库排序规则带 _CI 的MySQL 侧统一用 utf8mb4_general_ci 对齐大小写不敏感语义源库带 _CS 的MySQL 侧用 utf8mb4_bin。这样能避免大量“为什么同一行数据被重复插入”的诡异问题。MySQL 排序规则这块水很深尤其是查询结果里 ORDER BY 的表现两个库不一致时下游报表排序会跟预期差得很远。5. 避坑指南2.3.0 同步到 MySQL 的高频故障与排查顺序5.1 快照阶段一直读不完SQL Server 日志疯涨现象任务提交后进入了 initial 快照模式半天过去 Flink UI 上 source 的 Received Records 还在涨但 MySQL 里新增条数远低于源表总量同时 SQL Server 的日志文件增长很快writelog 等待明显。原因SQL Server 大表做全量快照时Flink 需要扫描大批量数据并记录位点。表的数据量过大、chunk.size 设置过大会让单次查询扫描太久事务日志持续累积捕获记录。另一个常见原因是源表没有主键快照读无法分片全部压在单线程上。解决先把 scan.incremental.snapshot.chunk.size 调到 2048 或 4096观察源库的 IO 和日志增长是否放缓。如果大表实在跑不动改用 latest-offset 模式只做增量历史数据用一次性初始化脚本先导到 MySQL。SQL Server 侧还要检查 CDC 清理作业把 retention 调小一点避免历史捕获数据堆积把磁盘打满。5.2 任务一启动就报 LSN 定位失败现象任务提交后没有任何数据进 MySQL日志里报“找不到 LSN 对应的捕获实例”或者“指定的 LSN 已无效”每次重启都是同一个位置挂掉。原因Flink 在 checkpoint 里记录了消费位点但 SQL Server 侧的 CDC 捕获实例被清理掉过或者数据库做了还原、日志被截断导致记录在案的 LSN 已经不存在。最常见的是 DBA 在维护窗口把 CDC 作业停了几天任务这边还在老位点上重启。解决确认 SQL Server Agent 一直在跑然后到 Flink 侧找到这个任务清除它的 state让任务从 latest-offset 重新开始。我的操作顺序是停任务在 SQL Server 执行 sys.sp_cdc_disable_table 再 sys.sp_cdc_enable_table 重建捕获实例Flink 端删除旧 savepoint再用 latest-offset 模式重启。注意这中间窗口期的变更不会被重复消费需要业务补数。5.3 MySQL 侧死锁和主键冲突反复出现现象任务能跑但经常报 Deadlock found 或者 Duplicate entry重试 3 次后任务失败延迟越来越大。查看 MySQL 的日志发现锁等待集中在 orders 表的同一批数据上。原因Flink JDBC sink 默认攒批写一批里对同一行的更新如果和另一批写入碰撞就会触发 MySQL 的锁竞争。MySQL 的锁分类里间隙锁在这类批量 upsert 场景中最容易冒出来尤其是 MySQL 默认隔离级别是 REPEATABLE READ 的情况下。另一个常见诱因是源表主键和目标表唯一键不一致导致 Flink 判断不出“这一行应该更新而不是插入”。解决先核对源表和目标表的主键、唯一键定义必须一致。然后在 JDBC url 里加 rewriteBatchedStatementstrue 并把 sink.buffer-flush.max-rows 从 1000 调低到 500减少单批写入的持锁时间。如果 MySQL 侧业务高并发并且对同一订单号反复更新考虑把事务隔离级别调到 READ_COMMITTED可以让间隙锁问题大幅缓解但这是 MySQL 侧的全局配置需要和 DBA 确认。5.4 checkpoint 频繁失败延迟从秒级涨到分钟级现象Flink UI 上 Checkpoint 一直 failed后台日志报 checkpoint expired任务没有挂但数据延迟从秒级变成分钟级source 一直堵着。原因EXACTLY_ONCE 模式下checkpoint 对齐需要所有 subtask 的输入都处理完如果某个时刻 sink 写 MySQL 变慢barrier 对齐就会超时checkpoint 失败后又会触发下一次形成恶性循环。解决我一般先加宽超时execution.checkpointing.timeout 从默认的 10 分钟改成 20 分钟同时调大 execution.checkpointing.min-pause-between-checkpoints保证两次 checkpoint 之间有空隙。再没用就把 sink.buffer-flush.max-rows 调小降低单次 flush 的耗时。都要还解决不了临时把模式切到 AT_LEAST_ONCE 救火数据不会丢但可能有少量重复等高峰期过去再切回来。这里的经验是别一上来就堆并行度先确认瓶颈在 MySQL 写入侧还是 checkpoint 对齐侧。5.5 上游 ALTER TABLE 加列任务直接挂掉现象上游业务在 orders 表加了一列 delivery_type加完没过几分钟Flink 任务报错退出重启后依然失败。原因Flink CDC 2.3.0 对 schema evolution 的支持非常有限SQL Server 端表结构变了但连接器持有的 schema 还是旧的读到变更数据时字段数量对不上整个任务就崩了。这是所有 CDC 同步方案里最头疼的问题因为 SQL Server 的捕获实例在表结构变更后也可能失效。解决先把任务停下按顺序处理。SQL Server 这边把该表的捕获实例禁用再重建让 CDC 重新捕获新结构的数据MySQL 这边同步执行 ALTER TABLE 加列Flink 这边修改 source 和 sink 表定义加新字段然后从 latest-offset 重启。这中间窗口期变更的数据需要靠业务侧或一次性脚本补齐。生产环境我一般会在变更窗口前通知业务暂停写一小会儿把风险降到最低。6. 验证与进阶从“能跑”到“敢上生产”的三件小事6.1 用 count 抽样双校验验证同步一致性任务跑通后第一件事不是看监控面板而是做一个朴素的校验分别查 SQL Server 源表和 MySQL 目标表的行数以及 MAX(updated_at)。两个数都对上再做一条抽样 SQL把订单金额和状态在两个库各自聚合一次对比。这个办法土但有效能挡住 90% 的字段映射错误。我用这个方案验证过三套任务两次都揪出了类型隐式转换问题一次是 money 被转成了 DOUBLE一次是 NCHAR 列同步后尾部空格差异。6.2 先入 Kafka 再入 MySQL值不值得加这一层同步链路里加一层 Kafka很多人觉得多余但如果你有多个下游要消费同一份 SQL Server 变更它就值了。链路从 SQL Server CDC → MySQL 变成 SQL Server CDC → Kafka upsert → MySQL等于把 SQL Server 的 CDC 捕获结果变成一份可多次消费的日志数据。代价是延迟从秒级变成秒级加毫秒级几乎可以忽略换来的是下游解耦和重放能力。我的建议是现在只有一个目标端直接写 MySQL预期会有第二个、第三个目标端就提前上 Kafka后面能少改很多任务。这套组合已经是实时数仓里比较常见的落法了。6.3 收尾建议留一个 24 小时观察窗口我自己带这类同步任务从不上线当天就宣告完工。任务跑起来后至少留 24 小时观察窗口看两条曲线Flink 的 checkpoint 完成时间曲线以及 MySQL 侧的写入延迟曲线。24 小时内如果能经受住业务高峰和一次手动重启这套链路基本算立住了。还有一个个人习惯每次改完同步 SQL 都要清一次 Flink 的 savepoint把它当作一次“不能后悔药的练习”强迫自己想清楚改动的影响边界。希望帮到你。本文还有配套的精品资源点击获取