news 2026/10/3 14:22:17

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Flink CDC实战:SQL Server实时同步到MySQL全链路指南

简介:一套基于 flink-connector-sqlserver-cdc 2.3.0 的 SQL Server 到 MySQL 实时同步工程示例,面向正在使用 Flink CDC 做数据接入、迁移或实时数仓同步的中高级工程师,也适合需要理解数据库变更捕获机制的开发人员快速上手。示例以 Maven 工程组织,包含 Java 源码、pom 依赖配置、数据库连接参数属性文件以及项目说明文档,重点展示源表与目标表的建表语句、数据流转插入逻辑、并行度设置、异常重试和一致性处理等关键写法。读者可按目录结构定位不同模块,将其中配置和代码直接移植到实际项目中,减少从零搭建的成本。压缩包共21个文件,以 Java 源码为主,辅以属性配置、XML 项目文件、说明文档和许可证声明,整体仅32KB,体量轻巧、便于阅读与修改。目前已有2313人学习,对于正在实施 SQL Server 变更数据捕获同步到 MySQL 的团队来说,是一份实用的落地参考。

1. 为什么实时同步这条路,我选 flink-connector-sqlserver-cdc 2.3.0

晚上十点半接到电话,说报表库数据对不上,定时任务半夜拉全量,源库一个慢查询就能让第二天的账全乱。把定时任务换成实时同步之后,我用 flink-connector-sqlserver-cdc 2.3.0 把数据从 SQL Server 实时同步到 MySQL 里,挂上之后基本不用再半夜爬起来看数据。这条链路不需要改业务表结构,只要在 SQL Server 上打开 CDC 功能,Flink SQL 写一段 DDL 就能从全量快照无缝切到增量日志。文章按落地顺序讲:依赖准备、最小链路、参数调整、真实踩坑,最后给验证和进阶用法。适合正在维护老 SQL Server、又要给报表继续喂 MySQL 数据的团队,尤其是每天被临时任务叫醒的那种。

2. 环境准备:开 SQL Server 的 CDC、配齐 2.3.0 依赖

2.1 先用 sysadmin 账号把 SQL Server 的 CDC 打开

开启 SQL Server CDC 是整个方案的第一个翻车点。很多人以为 Flink 那边装上 jar 就能读,其实连接器读的是源库的捕获实例,而捕获实例要靠 SQL Server 内建的 CDC 机制来写。先用一个 sysadmin 账号连到实例,确认当前数据库状态:

SELECT name, is_cdc_enabled FROM sys.databases WHERE name = 'TradeDB';

is_cdc_enabled返回 0 说明库还没开。然后执行数据库级开启和表级开启,这两步最好直接在 SSMS 的查询窗口里跑:

USE TradeDB; GO EXEC sys.sp_cdc_enable_db; GO EXEC sys.sp_cdc_enable_table @source_schema = 'dbo', @source_name = 'orders', @role_name = NULL, @supports_net_changes = 1; GO

sp_cdc_enable_db会在库里建出cdcschema,同时生成捕获和清理两个后台作业;sp_cdc_enable_table则会为dbo.orders生成cdc.dbo_orders_CT变更表,Flink 连接器只读这张变更表,不碰业务表。@role_name建议别直接传 NULL,我一般建一个cdc_reader角色,把权限收敛到这个角色里,给 Flink 的账号只要进角色就行,不用给它 db_owner。开启的表必须有主键,@supports_net_changes = 1才有效;没有主键的表,这一步会直接报错。

数据库级开启之后,SQL Server Agent 必须处于运行状态。CDC 的捕获是后台作业在刷,不是客户端拉的时候才去写。如果 Agent 没起来,或者正在跑的其他维护任务把作业挂起,你会遇到一个很迷惑的现象:Flink 连接器能连上,全量快照也能读到,但增量数据就是不动。排查时先看 Agent 的 Job 状态,再看sys.sp_cdc_help_jobs的输出里面有没有 capture 类型的任务。SQL Server 2008 到 2019 的 CDC 行为基本一致,2016 之后的版本对增量快照的并发处理更好一些,但前提都是 Agent 得活着。

2.2 配齐 Flink 运行环境和连接器 jar

我的常见做法是下载 Flink 发行版自己搭 session,不依赖云厂商的托管环境。版本上,我用 Flink 1.14.x 配 flink-connector-sqlserver-cdc 2.3.0 跑过,也见过同事用 1.13 和同样版本连接器跑通。2.3.0 的坐标在 Maven Central 上,用 Maven 拉最省事:

mvn dependency:get \ -Dartifact=com.ververica:flink-connector-sqlserver-cdc:2.3.0

把拉到的 jar 和 flink-connector-jdbc 一起放进 Flink 的 lib 目录。如果机器上没装 Maven,也可以直接去 Maven Central 搜flink-connector-sqlserver-cdc,把 2.3.0 对应的 jar 下载下来丢进 lib。flink-connector-jdbc 我用 2.x 这条线,比如 2.2.0,版本太老会和 2.3.0 的序列化结构对不上。

SQL Server 的驱动一般不用单独放,连接器的依赖关系里会带 mssql-jdbc;但 MySQL 驱动不会自带,还得手动放一个 mysql-connector-java 8.x 到 lib。三个文件摆齐后重启 Flink session,在 SQL Client 里先试一条命令验证环境:

./bin/sql-client.sh embedded

看到状态栏没有 ClassNotFound 之类的红色告警,再执行SHOW CURRENT CATALOG;确认连接器加载没问题。这步别急着跑业务 SQL,先确认 lib 目录下没有两个冲突版本的 MySQL 驱动,否则后面 sink 报错,你分不清是驱动冲突还是 SQL 写错。

2.3 MySQL 侧的目标库、账号和最后确认

MySQL 侧相对简单,提前建好目标库、写权限账号,再把目标表的字符集确认一下。常用建库语句是这样的:

CREATE DATABASE IF NOT EXISTS analytics DEFAULT CHARACTER SET utf8mb4 COLLATE utf8mb4_unicode_ci; CREATE USER 'sync_writer'@'%' IDENTIFIED BY '******'; GRANT SELECT, INSERT, UPDATE, DELETE ON analytics.* TO 'sync_writer'@'%';

目标表不用手工建也行,让 Flink 去写;但生产上我建议先手工建表,原因在第 4 章的类型映射里会细说。MySQL 8.0 对 utf8mb4 的支持比 5.7 更干净,排序规则用utf8mb4_unicode_ci还是utf8mb4_bin会影响下游 LIKE 查询的敏感度,和同步本身关系不大,但报表 SQL 会踩。MySQL 8.0 的驱动建议直接用 mysql-connector-java 8.x,5.7 反而容易在 SSL 握手上报错,后面的连接串里统一显式加useSSL=false,内网环境少一层加密开销,也能避开不少版本兼容问题。

3. 最小同步链路:两张 DDL 加一条 INSERT 跑通全量加增量

3.1 sqlserver-cdc 源表 DDL 与参数拆解

环境确认无误后,最小同步链路只需要两张表 DDL 和一条 INSERT。源表 DDL 是连接器工作的核心,先看订单表:

CREATE TABLE orders_src ( id INT, order_no STRING, user_id BIGINT, amount DECIMAL(12, 2), status STRING, create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'sqlserver-cdc', 'hostname' = '10.0.0.8', 'port' = '1433', 'username' = 'cdc_user', 'password' = '******', 'database-name' = 'TradeDB', 'schema-name' = 'dbo', 'table-name' = 'orders' );

这段配置是 2.3.0 里 source 最简写法,我只列了必填项。逻辑说明:PRIMARY KEY 那行不是给 SQL Server 看的,是 Flink 端用来识别变更主键、决定 upsert 语义的,NOT ENFORCED表示连接器不做约束校验,但会按主键去重。connector 固定写sqlserver-cdc,大小写敏感,写成sqlserver_cdc会直接报找不到连接器。

参数说明:hostname 要能连通 SQL Server 实例,如果用 AlwaysOn 的只读路由,要确保路由本身稳定;port 默认 1433,改过实例端口的要对应改。database-name 和 table-name 都支持正则,常见做法是'database-name' = 'TradeDB|OrderDB'这样一个作业同步多个同构库,前提是多张表结构完全一致。schema-name 默认 dbo,多 schema 时可以写成'dbo|sales'。这些正则写法在 2.3.0 上验证过,省去一个表一个作业的重复劳动。

3.2 jdbc sink 表 DDL 与 flush 参数设置

sink 表用 flink-connector-jdbc,字段顺序要和 source 对齐,主键必须和源表一致:

CREATE TABLE orders_sink ( id INT, order_no STRING, user_id BIGINT, amount DECIMAL(12, 2), status STRING, create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://10.0.0.20:3306/analytics?useSSL=false&rewriteBatchedStatements=true', 'table-name' = 'sync_orders', 'username' = 'sync_writer', 'password' = '******', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '2s', 'sink.max-retries' = '3' );

最后一条 INSERT 就是整个同步任务:

INSERT INTO orders_sink SELECT id, order_no, user_id, amount, status, create_time FROM orders_src;

参数说明:url 里useSSL=false的前提是 MySQL 端没有强制 SSL,否则写着 false 也会报 SSL 连接错误,这是 MySQL 8.0 的常见新坑;rewriteBatchedStatements=true对多行批量写入提升非常明显,增量高峰每分钟几千行时,没有这个参数 batch 提交效率会掉一个量级。sink.buffer-flush.max-rows控制攒多少行刷一次,sink.buffer-flush.interval控制最多等多长时间,两个条件满足其一就触发 flush。我一般把 max-rows 放在 1000 到 2000,interval 放在 1 到 3 秒,太小会频繁建连接,太大下游报表会看到延迟。sink 表必须有主键,flink-connector-jdbc 才能在源表发生 UPDATE 时生成按主键更新的语句,否则它只会 append,源表一改数据,MySQL 里就出现重复行。

3.3 开 checkpoint 再提交,别裸跑

这条任务能跑通和能稳定跑是两回事。SQL Client 模式下,把 checkpoint 当成同步链路的一部分,启动前先设置:

SET execution.checkpointing.interval = 10s; SET execution.checkpointing.externalized-checkpoints.retention = RETAIN_ON_CANCELLATION;

第一行是每 10 秒做一次快照,任务失败时 Flink 能从最近一次 checkpoint 恢复,已经同步到 MySQL 的数据不会重复消费;第二行让作业取消后保留 checkpoint 文件,这相当于后悔药,集群重启后能接上次的位置继续。提交方式我习惯用 SQL Client 批处理跑:

./bin/sql-client.sh embedded -f orders_sync.sql

如果想脱离 session 长时间跑,先把 job graph 生成出来,再用flink run提交,否则 session 一关作业就没了。生产环境里我一般先-f跑一遍看日志,确认 source 和 sink 都连通后,再用 detach 模式提交。第一次提交后注意观察 Flink UI 里的 Source 算子记录数,只有看到记录数开始增长,才说明全量快照已经启动。

4. 三个必调 source 参数与 SQL Server 到 MySQL 的类型映射边界

4.1 scan.startup.mode、chunk size、并行度调哪里

第一个必调参数是scan.startup.mode。2.3.0 支持initial和latest-offset两种,默认initial。initial 的含义是先做一次全量快照,再由连接器自动记住快照完成时的日志位置,之后从那里开始读增量。latest-offset 适合已经跑过一段时间、只想要新数据的场景。我给同事的建议是:新表第一次接实时同步,永远用 initial,宁肯让全量多跑一会儿,也不要为了省时间切 latest-offset 然后发现中间丢了一小段。切 latest-offset 时必须保证 SQL Server 的 CDC 捕获作业一直在跑,否则业务写进去的变更没有落在变更表里,丢了就是丢了。

第二个是scan.incremental.snapshot.chunk.size,默认 8096。它的含义是连接器把全量快照按主键区间切成多个 chunk,一个 chunk 一条 SELECT,避免一条大 SQL 压垮源库。大表第一次同步时,我会把 8096 调成 1024 到 2048,让 SQL Server 的负载更平缓;小表维持默认即可。反过来,源库是空闲主机、压力也不大的时候,可以调大到 20000 加速全量,但要注意 SQL Server 的锁和 tempdb 压力。这个参数在 source DDL 的 WITH 里直接写:

WITH ( 'connector' = 'sqlserver-cdc', ... 'scan.incremental.snapshot.chunk.size' = '2048' )

第三个是并行度。flink-connector-sqlserver-cdc 2.3.0 对 SQL Server 的增量读并没有 MySQL CDC 那种多 binlog 分片的能力,把 source 并行度写成 4,日志读取仍然是单线程在干,所以别指望并行度救性能。真正该并行的是下游 sink 和后续计算,source 并行度保持 1 到 2 就行。想要多张表并行,就在同一个作业里注册多个 source 表,让 Table 层的算子各自跑,而不是去调单表并行度。

4.2 类型映射:decimal 精度、datetime2 和 uniqueidentifier 按表设

连接器读出什么类型,决定了 sink 表建什么样。前面两张表 DDL 里已经默认 Flink 类型能装下 SQL Server 的值,但生产表往往复杂得多。按踩过的坑整理成一张对照表:

SQL Server 类型Flink 端实际映射MySQL 端建议列类型注意点
intINTINT无
bigintBIGINTBIGINT无
varchar / nvarcharSTRINGVARCHAR(长度)建议按源列长度手动建表
decimal(p,s)DECIMAL(p,s)DECIMAL(p,s)p 超过 38 时 MySQL 建不了列,需要提前 CAST
datetime2TIMESTAMP(3)datetime(3)毫秒以下精度会丢,讲究的话源表收 STRING
datetimeoffsetSTRINGvarchar(50)不建议直接映射 datetime,时区信息会丢
uniqueidentifierSTRINGchar(36)小写横线格式,下游要处理
bitBOOLEANtinyint(1)MySQL 端按习惯选

表里最坑的是 decimal 和 datetime2。SQL Server 允许 decimal(38,10) 这类超高精度,MySQL 8.0 的 decimal 最大也能到 65 位精度,但实际建列时很多 DBA 只给到 decimal(20,4),字段精度不够写不进去。常见做法是在 source 表 DDL 里直接 CAST,从源头把精度压到目标库能接受的档位:

SELECT id, order_no, user_id, CAST(amount AS DECIMAL(18,2)) AS amount, status, create_time FROM orders_src;

datetime2 的问题在于 Flink 端默认映射成 TIMESTAMP(3),源库存的是2024-01-01 12:23:45.1234567,读出来就是2024-01-01 12:23:45.123,丢了三位小数。对审计类字段,我建议把源表 DDL 里对应字段声明成 STRING,sink 也建 varchar(30),精度原样保留,代价是下游 SQL 要做一次日期转换。uniqueidentifier 同理,直接当字符串走最省事,省得两端类型不一致导致写不进去。

4.3 jdbc sink 的三个参数决定写入延迟

sink 侧除了前面提过的 flush 两件套,还有一个容易被忽略的关键点:连接串上的配置。MySQL 端默认连接数 151,过高的并行度会把连接池打爆,出现Too many connections。我把 JDBC sink 的并行度控制在 2 到 4,配合连接池上限 200 左右即可。另一个可调选项是sink.max-retries=3,避免偶发网络抖动直接把任务搞失败。

调完参数后再看数据延迟:在 SQL Client 里查指标不方便,最快的方式是到 MySQL 里查某一行的更新时间,再和源库同一条记录的 update 时间对比。增量场景下,从 SQL Server 提交事务到 MySQL 可见,通常能到秒级;如果出现分钟级延迟,优先怀疑 source 并行度是不是被调得过大,或者 JDBC sink 的 batch 堆积。rewriteBatchedStatements只对 INSERT 生效明显,如果同步的任务里大量是 UPDATE,还得看主键是否稳定,主键不稳定,sink 的更新语句会退化成一条条执行,延迟直线上升。

5. 避坑排查清单:五条从 SQL Server 同步到 MySQL 的真实踩坑记录

5.1 增量一直不动:先怀疑 SQL Server Agent 而不是 Flink

现象:作业正常跑,全量数据都过去了,之后源表一改,MySQL 侧纹丝不动。原因:SQL Server 的 CDC 捕获作业没在跑,或者 Agent 服务没启动。连接器从捕获实例读不到新日志,自然没数据。解决:先看 SQL Server Agent 的服务状态,再执行下面查询确认捕获会话:

EXEC sys.sp_cdc_help_jobs;

如果没有 capture 类型的 job,说明建库时 Agent 没有正常创建作业,重启 Agent 后手动调用sys.sp_cdc_start_job重新拉起。很多团队装完 SQL Server 顺手把 Agent 服务禁了,这是 CDC 方案里最常见的翻车点。

5.2 SSL 加密导致连接失败

现象:Flink 作业启动后周期性报The driver could not establish a secure connection to SQL Server by using Secure Sockets Layer (SSL) encryption,任务一直重启失败。原因:SQL Server 2019 之后很多实例默认把加密选项调成 Force Encryption,mssql-jdbc 驱动在握手阶段要验证服务器证书,自签名证书或内网证书不被信任就会断。解决分两条,按环境选一条。如果内网传输可信,可以在 SQL Server 配置管理器里把 Force Encryption 改成 Optional,属于源头开关;如果 DBA 不允许动实例配置,就把 SQL Server 的证书导入 Flink 运行环境的 truststore。注意别盲目在连接参数里关验证,不同版本的连接器对trustServerCertificate这类选项支持不一致,最稳的做法是在实例侧和团队确认清楚。

5.3 全量快照阶段把源库拖垮

现象:第一次同步大表,SQL Server CPU 冲到 100%,业务侧开始出现锁等待。原因:chunk size 太大,连接器用大范围查询反复扫表,加上正好撞上业务高峰。解决:把scan.incremental.snapshot.chunk.size从默认 8096 调到 1024,把 source 并行度设为 1,全量阶段相当于给连接器限速;再不行就把 initial 同步安排在凌晨执行,先跑全量,之后切回正常窗口做增量。全量慢不是病,源库被杀才是事故。

5.4 无主键表和 datetime2 低精度一起翻车

现象:A 表能同步,B 表建 source 时直接报错;C 表源库时间是 7 位小数,MySQL 里只有 3 位。原因:B 表没有主键,SQL Server CDC 不允许无主键表开启捕获;C 表的 datetime2 被 Flink 默认降成 TIMESTAMP(3)。解决:无主键表先在 SQL Server 侧补主键,哪怕是业务主键;datetime2 字段在 source DDL 里声明为 STRING,在 sink 里建 varchar,把精度责任转移给 MySQL 端再用 CAST 转换。同步任务是长期跑的,字段不能将就,将就的结果是下游对不上账。

5.5 作业停了太久,重启后增量从日志中间断掉

现象:任务因为发布停机两周,重启后 MySQL 里少了一段某几小时的数据。原因:SQL Server CDC 的清理作业会把超过保留期的变更记录清掉,日志留不住就补不回来。解决:停机前确认 CDC 清理作业的保留时间,临时把sys.sp_cdc_change_job的 retention 调大,或者在恢复时手动补一次全量。对长期依赖实时同步的团队,我的习惯是每周至少看一次任务健康度,而不是等月底对账才发现少数据。这条是血泪经验:实时链路不比定时任务省心,只是把问题提前暴露了。

6. 进阶:用窗口核对校验同步、多表宽表合并与作业健康监控

6.1 用时间窗口核对校验同步没丢

验证同步没丢,最扎实的方式是抽最近 10 分钟有更新的记录做对比。源库执行一条统计:

SELECT COUNT(*) AS cnt FROM TradeDB.dbo.orders WITH (NOLOCK) WHERE update_time >= DATEADD(MINUTE, -10, GETDATE());

再到 MySQL 目标表跑同样的时间条件统计:

SELECT COUNT(*) AS cnt FROM analytics.sync_orders WHERE update_time >= NOW() - INTERVAL 10 MINUTE;

两边数量一致说明最近十五分钟链路是通的;不一致就把主键列表拉出来做差集。这类核对要写进运维脚本,每天早上跑一次,比等月底对账体面得多。

6.2 多表 join 宽表与整库同步的正则写法

如果 SQL Server 里有多张要同步的表,比如订单主表和订单明细,可以注册两个 source、两个 sink,用一段 SQL 做 join 后再写入 MySQL 宽表,省掉一次中间库存储。整库同步则用 database-name 和 table-name 的正则参数拆作业,按业务域分开跑,避免单作业背太多表、一处慢全链路阻塞。

6.3 用 Checkpoint 做心跳监控

日常监控不用搞复杂,看 Checkpoint 完成时间就能判断任务健康度。长时间不 checkpoint 基本就是卡死,优先查 source 是否在等待锁、sink 是否被 MySQL 慢查询拖住。再配合上面两个核对脚本,实时同步就能做到半夜不接电话。希望帮到你。

本文还有配套的精品资源,点击获取

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/3 14:20:51

MySQL DDL 执行指南:锁表原理、三种算法与生产环境避坑

MySQL 上执行 DDL,看起来就是几条 ALTER 语句的事,但真正在业务环境里动过手的朋友都清楚,这是数据库运维里最容易“翻车”的一类操作。尤其是表到了千万级、亿级,哪怕只是加一个索引,一条 SQL 都可能把主库拖到报警&a…

作者头像 李华
网站建设 2026/10/3 14:20:49

LightTools菲涅尔透镜手动建模五大技巧,避开杂散光与效率坑

做光学设计的同行应该都有这种感觉:LightTools里做非球面、自由曲面,靠自带的建模工具就能搞定,但一到菲涅尔透镜就头疼。尤其是需要手动创建的时候,齿高、环距、脱模角、圆角这些参数一旦处理不好,lighttools里看起来…

作者头像 李华
网站建设 2026/10/3 14:20:20

正则匹配实战:从规则原理到性能避坑指南

正则匹配这个东西,很多人第一反应是“不就是查个字符串吗”,但真到了线上日志排查、数据清洗、接口参数校验的时候,才发现自己写出来的表达式要么匹配不到、要么误杀一片。我过去几年里在项目里被正则坑过无数次,也靠它救过急&…

作者头像 李华
网站建设 2026/10/3 14:20:19

Altium Designer批量替换元器件全攻略:原理图、封装与库同步

1. 为什么你一定会用到批量替换元器件 1.1 实际项目中的高频替换场景 做硬件设计的朋友,几乎没人逃得过“替换元器件”这件事。最常见的情况就是:板子画到一半,采购跟你说某颗料缺货、交期排到明年,或者代工厂反馈你选的封装工艺…

作者头像 李华
网站建设 2026/10/3 14:20:13

基于Python的差分隐私协同过滤推荐系统设计与实现

简介:一份基于Python实现差分隐私与协同过滤相结合的推荐系统毕业设计资源,适用于计算机相关专业学生、推荐系统研究者及隐私保护技术学习者。内容覆盖推荐系统隐私保护研究背景与国内外现状、协同过滤算法主要步骤、差分隐私概念与常用实现机制&#xf…

作者头像 李华
网站建设 2026/10/3 14:18:42

SOEM与Qt实现EtherCAT软主站:从交叉编译到嵌入式部署实战

QT和SOEM这套组合,说实话在工业自动化圈子里不算新鲜,但真正把这套东西从源码编译一路玩到嵌入式板子上、还要配上Qt界面做成人机交互的,网上能一篇讲透的教程并不多见。我最早接触SOEM是因为一个产线改造项目——要用EtherCAT总线带动8个伺服…

作者头像 李华