
简介面向 Flink 实时计算开发者的完整工程模板聚焦 flink-connector-sqlserver-cdc 2.3.0 在 SQL Server 到 MySQL 增量同步场景中的落地实践。资源包内为 Maven 标准结构共 21 个文件包含 11 个 Java 源文件实现数据源定义、转换逻辑与输出、5 个 properties 配置数据库连接、CDC 参数及作业配置、2 个 XMLpom 依赖与构建配置、2 个 Markdown 说明文档及许可证文件压缩包仅 32KB轻量便于快速导入。已有 2313 人浏览学习。通过阅读源码与配套说明可掌握 SQL Server CDC 连接器依赖引入、源表与 MySQL 目标表的 CREATE TABLE 写法、INSERT INTO 数据流转方式以及并行度、容错等调优思路适合需要快速搭建实时同步任务或学习 Flink CDC 原理的开发者参考。1. 为什么把 SQL Server 实时同步到 MySQL 绕不开 flink-connector-sqlserver-cdc 2.3.0业务库落在 SQL Server 2019 上报表和运营后台却要求 MySQL 能查到准实时数据这是很多数据中台项目开场就要面对的局面。flink-connector-sqlserver-cdc 2.3.0 解决的就是这条链路它借助 SQL Server 自身的事务日志 CDC 机制把 insert、update、delete 解析成变更事件再由 Flink 作业把数据写入 MySQL。整个过程既不用给业务系统加双写逻辑也不用每天凌晨跑批拉数延迟可以控制在秒级甚至更快。适合正在做异构数据库同步、订单中心拆分、或者维护着 MySQL 离线数仓但不想继续等 T1 的人。你只要会基础 SQL 和一点点 Flink SQL就能把它跑起来。2. 同步链路设计与选型SQL Server CDC 机制和 Flink 作业拓扑2.1 SQL Server CDC 是怎样产生变更数据的SQL Server 从 2008 版本开始内置 Change Data Capture 能力。它的工作方式不是靠前端应用上报而是靠 SQL Server Agent 启动一个捕获作业持续扫描事务日志中被标记的区域把变更后的数据连同__$operation、__$start_lsn这类元信息写入目标库的cdcschema 下的捕获实例表中。比如你在dbo.orders表上启用 CDC就会生成一个cdc.dbo_orders_CT表里面既有变更前快照也有变更后快照。Flink CDC 连接器做的是把这一层隐藏起来它以一种兼容 Debezium 的方式读取这些捕获表或日志输出成标准ChangeRecord。所以在 Flink SQL 里你看到的 source 表和普通 MySQL 表没有明显区别但背后已经自动完成了 LSN 追踪、事务边界识别和变更事件解析。这里有个很容易搞错的关键点如果数据库级 CDC 没开启Flink 作业仍然能连接 SQL Server却永远等不到增量数据。所以无论你用的是 flink-connector-sqlserver-cdc 2.3.0 还是其他版本第一步永远是先确认 SQL Server 自己的 CDC 开关已经打开并且 SQL Agent 服务正在运行。2.2 为什么轮询和 DebeziumKafka 在此场景都不是首选常见的传统做法是写一个 JDBC 轮询任务每几秒执行一次select * from orders where create_time ?。这个方案最大的问题是天然漏掉 deleteupdate 也只拿得到最新值拿不到变更前的值。而且高频轮询会让 SQL Server 的 CPU 和 IO 压力集中在某几张热表上数据量过千万后会越来越难维护。也有的团队选择 Debezium KafkaDebezium 抽取 SQL Server 日志Kafka 做消息缓冲再起一个 Consumer 写入 MySQL。这套链路很完整也能做到实时但你的基础设施里需要多出 Kafka、Schema Registry 和至少一个消费者服务。如果项目目标仅仅是“把几张业务表实时同步到 MySQL”用 Flink CDC 一条链直接拉通维护成本要低一个量级source 和 sink 都在同一个 Flink job 里状态、Checkpoint、恢复逻辑也都由 Flink 统一管理。2.3 三种同步方案对比选型不是越复杂越好方案捕获删除需要额外组件链路维护成本适合场景JDBC 定时轮询否无低数据量小、可容忍几分钟延迟Debezium Kafka是Kafka、消费者服务高上下游都要复用同一份变更流flink-connector-sqlserver-cdc 2.3.0是仅 Flink中异构库间实时同步、单/多表同步选 flink-connector-sqlserver-cdc 2.3.0 还有一个实际理由它用的是一套增量快照框架。老版本是全表锁住做完快照再切增量2.3.0 会把表按主键或指定键拆成多个 chunk分块读取并断点续做源库压力明显更小。这也是它在生产环境里比 1.x 更值得投入的原因。3. 搭建最小同步作业从开启 SQL Server CDC 到 Flink SQL 双表对接3.1 前置环境与 jar 包准备先确认你的环境满足几个基础条件SQL Server 2016 及以上版本建议选 2017/20192022 也可以但要确保 SQL Server Agent 服务已经安装并启动MySQL 目标库建议用 5.7 或 8.0Flink 集群可以是 standalone 或者基于 YARN/Kubernetes 的部署方式。你需要把两个 jar 放到 Flink 安装目录的lib/下一个是flink-sql-connector-sqlserver-cdc-2.3.0.jar另一个是flink-connector-jdbc的对应版本。jar 放进去之后重启 Flink 集群或 SQL Client然后用SHOW JARS;确认能看见它们。这一点别跳过很多作业提交后报 ClassNotFoundException问题都出在 jar 没有真正加载。SHOW JARS;如果你刚装完 SQL Server 2019建议顺手安装 SSMS 并登录一次目标库后面查看 CDC 捕获作业、监控 LSN 都会方便得多。目标 MySQL 侧也要提前建好需要的库和表不要指望 Flink 自动建表JDBC sink 默认没有建表能力。3.2 在 SQL Server 上开启数据库级和表级 CDC下面这段 SQL 需要以 sysadmin 身份在源库执行USE MyDB; GO -- 开启数据库级CDC EXEC sys.sp_cdc_enable_db; GO -- 对 dbo.orders 开启表级CDC EXEC sys.sp_cdc_enable_table source_schema Ndbo, source_name Norders, role_name NULL, filegroup_name NPRIMARY, supports_net_changes 1; GO先说第一行USE MyDB它把当前上下文切到你真正要同步的库。sys.sp_cdc_enable_db只做一件事在库级别打开 CDC生成的元表会放在cdcschema 下。sys.sp_cdc_enable_table里的role_name可以传NULL表示所有能读该库的用户都可以访问捕获表supports_net_changes 1表示额外支持“净变更”查询方式Flink 的某些增量读取路径会依赖它。执行完成后可以验证SELECT name, is_cdc_enabled FROM sys.databases WHERE name MyDB; SELECT * FROM cdc.change_tables WHERE source_schema_name dbo AND source_table_name orders;第一张表确认数据库级 CDC 是否生效第二张表能看到系统生成的捕获实例通常是dbo_orders_CT。如果这里查不到记录说明表级开启没有成功后面 Flink 作业即使跑起来也不会有真正的变更数据输出。3.3 用 Flink SQL 定义 Source 和 MySQL Sink环境准备好后写一个sync_orders.sql文件内容如下-- 源表SQL Server 的 dbo.orders CREATE TABLE sqlserver_orders ( order_id INT, user_id INT, amount DECIMAL(10,2), status STRING, create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector sqlserver-cdc, hostname 10.0.0.10, port 1433, username cdc_user, password Cdc_Pass_2024, database-name MyDB, schema-name dbo, table-name orders, scan.startup.mode initial ); -- 目标表MySQL 里的 orders CREATE TABLE mysql_orders ( order_id INT, user_id INT, amount DECIMAL(10,2), status STRING, create_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://10.0.0.20:3306/mysql_replica?useSSLfalseserverTimezoneAsia/Shanghai, table-name orders, username sync_user, password Sync_Pass_2024 ); INSERT INTO mysql_orders SELECT * FROM sqlserver_orders;这段 DDL 需要特别注意几个参数scan.startup.mode initial表示作业启动时先做一次全量快照再切到增量日志历史数据和实时变更都会一起同步。如果你想跳过全量、只从当前 LSN 开始接增量可以改成scan.startup.mode latest-offset但第一次搭建时建议用initial确保源库已有的数据也落到 MySQL。PRIMARY KEY (order_id) NOT ENFORCED并不是真的去目标库建主键而是告诉 Flink 这个表有主键方便它做 state 管理和 upsert 优化。MySQL 目标表本身必须提前建好主键否则 JDBC sink 无法保证精确一次语义。提交作业的常见命令是bin/sql-client.sh embedded -f ~/jobs/sync_orders.sql如果你想在独立模式下提交也可以把 SQL 转成flink run -py或通过INSERT INTO的连招在 SQL Client 里手动执行。跑起来之后用SHOW JOBS;或 Flink Web UI 看作业状态一个正常的作业应该显示 RUNNINGSource 和 Sink 并行度各为 1。3.4 checkpoint 与并行度设置同步作业最怕的是网络抖动或者 MySQL 短暂不可用Flink 的 checkpoint 就是后悔药。建议在作业里加上基础参数SET execution.checkpointing.interval 10s; SET execution.checkpointing.mode EXACTLY_ONCE; SET execution.checkpointing.min-pause 5s; SET execution.checkpointing.tolerable-failed-checkpoints 3;Checkpoint interval 不宜设得太短源库在高峰期可能每秒上万条变更每 2 秒做一次 checkpoint 反而拖累日志读取10 秒是个比较稳的起步值。min-pause 5s表示上一次 checkpoint 结束到下一次开始至少隔 5 秒避免频繁触发。并行度方面flink-connector-sqlserver-cdc 2.3.0 默认是单并行度只有对做了 chunk 拆分的表才能提升 source 端并行度。多并行度读取不是设置一个parallelism.default4就完事了还需要满足两个条件表要有明确主键或者指定scan.incremental.snapshot.chunk.key-column否则作业会回退到单线程全量扫描快照速度没有任何提升。初学阶段先把并行度保持为 1专注跑通链路再考虑压快照。4. 连接参数与时区陷阱让 2.3.0 稳定跑起来的细节4.1 SSL 加密导致的连接失败与 encrypt 参数很多人第一次跑这个作业报错并不是 CDC 相关而是连接本身就失败了。日志里出现类似下面的内容[08001] [Microsoft][ODBC Driver 17 for SQL Server] SSL Provider: 证书链是由不受信任的颁发机构颁发的原因是 SQL Server 实例开启了强制加密默认客户端又会对服务端证书做公网验证而本地部署的 SQL Server 通常用的是自签名证书链验证自然失败。flink-connector-sqlserver-cdc 2.3.0 支持通过debezium.*前缀透传底层参数给 Debezium 的 SQL Server connector这时可以在源表 WITH 里追加debezium.database.encrypt false关闭 encrypt 后底层以非加密方式连接 SQL Server。如果你的网络环境安全、且跑在内网这个做法常见也够用。如果安全规范不允许关闭加密也可以尝试debezium.database.trustServerCertificate true但要注意这个参数能否透传成功受连接器版本影响。我一般先看日志只要报错停留在协商阶段优先用encryptfalse快速打通打通后再根据安全团队要求决定是否需要正式 CA 证书。4.2 TIMESTAMP 时区当 SQL Server 遇上是 MySQLSQL Server 的datetime、datetime2本身不带时区信息Flink 会把它映射成TIMESTAMP(3)。如果 Flink 作业部署机器的时区是 UTCMySQL 连接串又设置了serverTimezoneUTC同步过去后你会在 MySQL 里看到比实际时间早 8 小时的数据。最简单的规避方法是在 Flink SQL 环境里显式声明SET table.local-time-zone Asia/Shanghai;同时把 MySQL JDBC URL 里的serverTimezone也统一成Asia/Shanghai让两端解释时间口径一致。更严谨的做法是把源表create_time定义为TIMESTAMP_LTZ(3)但 SQL Server 的 CDC 捕获表本身不记录带时区的时间多数场景没这个必要。这个坑很隐蔽因为作业不会失败只有对账时才会发现时间对不上。建议第一次全量同步完成后立刻在源库和目标库各查一条最新数据的create_time做比对不要等到第二天报表出数据才炸。4.3 用户权限、SQL Agent 与事务日志的常规体检Flink CDC 作业的源库账号实际需要读取捕获表的权限。给一个固定同步账号直接授予db_owner是省事做法生产环境如果权限收紧至少也要保证该账号能读cdcschema 下的所有对象并拥有对源表的 SELECT 权限。数据库级 CDC 的开启最好由 DBA 执行普通账号一般不应有sysadmin。SQL Server Agent 服务是另一个重点是对象。CDC 捕获任务依赖 Agent如果 Agent 停了捕获作业会在队列里堆积事务日志会因为没有及时读取而不断膨胀甚至把磁盘撑满。作业跑起来后建议用下面这段检查最近的日志扫描状态SELECT session_id, start_time, end_time, duration, log_record_count, tables_processed FROM sys.dm_cdc_log_scan_sessions ORDER BY start_time DESC;如果经常看到duration特别长但log_record_count很少说明扫描异常或捕获作业被卡住。这种问题在 Flink 侧很难通过日志发现必须回到 SQL Server 侧排查 Agent 作业、捕获线程和服务账号状态。5. 常见问题优先排查5 条踩坑记录与解决思路5.1 只有全量没有增量现象作业显示 RUNNINGMySQL 里能查到启动时同步过来的历史数据但源库继续 insert、updateMySQL 完全没有新数据进来。原因SQL Server 数据库级或表级 CDC 没有真正开启Flink 的增量部分拿不到任何变更日志或者 SQL Server Agent 没有启动捕获作业从未运行。解决先用 3.2 节里的验证 SQL 确认cdc.change_tables有该表记录再检查 Agent 服务。如果 Agent 是启动之后才启用的 CDC捕获作业会延迟一段时间才能补扫历史日志等一两分钟再看sys.dm_cdc_log_scan_sessions里有没有新增 session 即可。5.2 连接报 08001 SSL 证书链错误现象作业启动时抛出[08001] SSL Provider: 证书链是由不受信任的颁发机构颁发的或者类似客户端无法建立连接。原因SQL Server 强制加密客户端校验了自签名证书证书链不受本地 JDK truststore 信任。解决参考 4.1 节在 source WITH 中加入debezium.database.encrypt false临时关闭加密快速打通内网链路。如果关闭后仍然报错再检查hostname是否配成了 SQL Server 的计算机名而不是 IP主机名解析到错误的证书信息也会导致握手失败。5.3 MySQL 侧主键冲突和重复数据现象同步作业本身没有报错但 MySQL 中同一个主键出现两条记录或者数据库日志里出现 duplicate key。原因最常见的是目标 MySQL 表没有主键或联合唯一键JDBC sink 只能追加写无法识别同一行是 update另一种情况是 Flink 作业从 checkpoint 恢复后MySQL 里已经有部分未提交的数据而 sink 的 flush 缓冲又造成重复衔接。解决先给 MySQL 目标表补上主键再检查 Flink DDL 里的PRIMARY KEY定义和 source 表主键是否完全一致。JDBC sink 生产环境建议设置sink.buffer-flush.max-rows 0让每批写入按内部连接池节奏刷新不要依赖小批次 flush。还要确认用的是 2.3.0 配套的 JDBC connector而不是老版本部分老版本对 upsert 支持不完整。5.4 大表快照长时间停留在 0%现象有几万或上百万行的大表作业启动后 Flink Web UI 显示总算子长时间不变快照进度始终很低。原因scan.startup.mode initial在 2.3.0 里会走增量快照框架但表如果只有nvarchar列而没有明确的 chunk 拆分键连接器可能退化为单线程读取又或者源表在快照期间产生了大量 DML导致日志读取和快照读取相互争抢。解决为无主键大表显式指定拆分键例如scan.incremental.snapshot.chunk.key-column order_id同时调低快照阶段并行度避免 4 个并行读取对源库同时发起大查询。快照阶段不要盲目把 checkpoint 间隔调得很短快照本身是一次长生命周期作业checkpoint 间隔 30 秒到 60 秒更合理。5.5 源表加列后作业变慢甚至报错现象业务在源表上新增了一个列同步作业开始报 schema 不兼容错误或者 JSON 序列化异常。原因Flink SQL 在作业启动时已经基于 DDL 锁定了字段列表源表 schema 变更不会自动映射到已经运行的作业。SQL Server CDC 捕获表在源表加列后也可能重新映射捕获实例旧的列序对不上。解决如果新增列不是必选项目可以先不处理如果需要同步新列必须停止作业、更新 Flink DDL、重新保存 checkpoint 后启动作业。不能依赖动态 schema 更新这是当前 SQLServer CDC 连接器在生产上的边界。同理源表 drop 列时更建议先停作业以免日志解析阶段出现无法反序列化的历史记录。6. 再把扩展性和延迟监控补上这步做完才敢上生产单表同步跑通只是第一步上生产前我至少会做三件事。第一把延迟监控接入到现有告警。Flink 本身没有直接的“数据延迟”指标但你可以通过 SQL Server 侧观察捕获作业状态或者在下游 MySQL 里放一张带update_time的心跳表每 10 秒写一次当前时间再用 Flink 同步这张心跳表最后通过查询 MySQL 最大时间差算出真实链路延迟。比起看 source 物理延迟这种方式能直接反映可见数据延迟。第二做一次源表和目标表的对账。对账不复杂先分别数主键和行数SELECT COUNT(*) AS order_cnt, SUM(order_id) AS checksum FROM MyDB.dbo.orders;目标库执行同样两条语句然后对比order_cnt和checksum。行数一致不一定内容一致但至少能快速发现漏数据或重复同步。正式对账建议每天低峰期跑一次如果持续差异再考虑抽样比对单行内容。第三动态扩表不要一上来就用正则同步整库。flink-connector-sqlserver-cdc 2.3.0 的table-name支持带正则或|分隔符的多表配置但这些表 schema 必须完全一致实际项目中很少条件成立。我现在的习惯是一个业务域一张 source 表对应一个 JDBC sink至多把同结构的分库分表用table-name dbo.orders_[0-9]统一到一个作业里其余情况宁可多配几个独立作业也不要为了省资源把不同 schema 的表塞进同一个作业。跑这类同步作业踩得多了我更相信 SQL Server 侧的状态比 Flink 日志更有话语权事务日志增长、Agent 捕获 session、CDC 捕获表堆积这些才决定同步能不能长期稳定。先承认连接器有边界再基于检查点和基础监控把链路兜住这套方案就足够拿去交付了。希望帮到你。本文还有配套的精品资源点击获取