简介:面向大数据开发与数据集成工程师,提供一套基于 flink-connector-sqlserver-cdc 2.3.0 的实时同步工程示例,解决将 SQL Server 变更数据持续同步到 MySQL 的典型场景。包内包含完整 Maven 项目结构,通过 SQL Server CDC 连接器捕获增删改事件,并借助 JDBC 接收器写入 MySQL,覆盖依赖配置、表结构定义、数据流转 SQL 及参数调优等关键内容。资源共 21 个文件,以 11 个 Java 源码为主,配合 properties 配置、XML 描述及 README 文档,压缩包仅 32KB,结构清晰,便于直接导入编译和改造。已有 2313 人浏览学习,适合需要快速落地 Flink CDC 同步任务、理解连接器配置与排错思路的开发者参考。借助这份源码与说明,可大幅缩短环境搭建和调试验证时间,并能在现有基础上扩展多表同步、调整并行度或接入其他下游存储。
1. 把SQL Server的数据实时搬到MySQL:先说结论和适用边界
做数据同步的工程师大概率都遇过这种需求——业务库在SQL Server,分析库在MySQL,每天凌晨跑批全量同步,第二天早上看数据还是昨天下午的。业务方催着要实时数据,DBA又不愿意让ETL工具直连业务表做轮询,怕影响线上性能。这个项目给的方案是用flink-connector-sqlserver-cdc 2.3.0,让Flink以CDC(Change Data Capture)的方式订阅SQL Server的事务日志,把增删改操作实时解析出来,再通过JDBC连接器写入MySQL。CDC不是轮询,是数据库主动把变更记录暴露给消费者,对业务库的压力小很多,实时性也能压到秒级甚至毫秒级。
这个资源适合两类人:一类是刚接触Flink CDC、想快速搭出一条SQL Server到MySQL同步链路的入门者,它能直接给你一份可跑的Maven工程,省去翻文档凑依赖的时间;另一类是已经在用其他同步工具、想评估Flink CDC方案能不能接手生产任务的工程师,你可以拿这份代码做基准测试,重点看它的增量快照机制、断点续传能力和主键处理策略。我拆完这份资源后最大的感受是:Flink CDC的上手门槛不在Flink本身,而在SQL Server端的权限和配置配合,这一块不弄对,后面全是白折腾。下面按我实际复现的顺序,从环境准备到参数调优一步步说。
2. 环境与前置配置:SQL Server开启CDC是第一步,没它一切免谈
2.1 为什么Flink CDC必须依赖SQL Server自带的CDC机制
flink-connector-sqlserver-cdc并不是自己发明了一套数据捕获方法,它是在SQL Server原生CDC功能之上做了一层封装。SQL Server从2008开始支持CDC,原理是:当你对某个表开启CDC后,SQL Server的捕获进程会扫描事务日志,把变更记录写入系统表cdc.<capture_instance>_CT,这个CT表里同时保留了操作类型(__$operation字段,1=删除、2=插入、3=更新前、4=更新后)和变更前后的数据快照。
Flink连接器干的事就是用sys.sp_cdc_get_captured_columns、sys.sp_cdc_get_ddl_history这些系统存储过程拿到表结构和变更历史,然后以轮询CT表的方式模拟流式读取。所以有个硬性前提:SQL Server必须先对目标表开启CDC,Flink才能读到数据,如果跳过这一步直接跑作业,日志里只会出现连接成功但没有任何数据输出的怪象。
开启CDC需要三个前置条件:SQL Server版本是2008及以上(企业版、开发版、标准版都行,但Express版不支持);SQL Server Agent服务必须处于运行状态,因为捕获和清理作业都是Agent调度的;用于同步的数据库账号需要有db_owner或至少VIEW SERVER STATE、SELECT ALL USER SECRETS等权限。很多时候Agent没启动,CDC作业压根不执行,CT表里从未写入过数据,Flink这边连接正常却收不到任何东西。
2.2 在目标库执行开启CDC的SQL脚本
我一般会先单独跑一遍开启脚本,确认CT表和CDC作业创建成功,再去动Flink那边的配置。下面这组脚本是复现这个项目时用的,你可以在SQL Server Management Studio里对目标库和目标表依次执行:
-- 1. 检查当前库是否已启用CDC SELECT name, is_cdc_enabled FROM sys.databases WHERE name = 'YourDatabase'; -- 2. 库级别开启CDC EXEC sys.sp_cdc_enable_db; -- 3. 对目标表开启CDC,并指定一个捕获实例名称,方便后续识别 EXEC sys.sp_cdc_enable_table @source_schema = N'dbo', @source_name = N'YourSourceTable', @role_name = NULL, @capture_instance = N'dbo_YourSourceTable'; -- 4. 确认捕获实例已生成 EXEC sys.sp_cdc_help_change_data_capture;这里重点说两个参数:@role_name设置为NULL意味着不限制访问CDC数据的角色,任何有权限的账号都能读CT表,如果公司安全审计严格,建议改成具体角色名;@capture_instance是自定义的捕获实例名,默认格式是<schema>_<table>,如果表名带下划线很容易读混,我习惯显式指定。
执行完后到系统表里验证一下,如果cdc.change_tables表里出现了刚才的捕获实例,说明开启成功。这一步最常见的失败原因就是SQL Server Agent没启动,sp_cdc_enable_table虽然返回成功,但后台的捕获作业根本没创建出来。所以务必备一下Agent状态:
-- 查看Agent是否在运行,status为4表示Running EXEC xp_servicecontrol N'QUERYSTATE', N'SQLServerAGENT';2.3 Maven依赖与Flink环境版本匹配
项目默认的Maven依赖坐标是flink-connector-sqlserver-cdc_2.11-2.3.0,注意这个命名里2.11是Scala版本,2.3.0才是连接器版本。你的Flink发行版如果是1.13.x或1.14.x,配这个连接器比较稳妥;如果Flink已经升到1.16以上,建议直接去找对应新版本的连接器,否则可能出现akka版本冲突。
还有个容易忽略的点:连接器自带的flink-cdc-base依赖会引入它指定的Flink版本相关类,如果你的Flink集群版本比连接器要求的低,ClassNotFound的报错会非常隐蔽。我复现这个项目时用的是Flink 1.14.2,pom里需要显式把flink版本属性对齐:
<properties> <flink.version>1.14.2</flink.version> <scala.binary.version>2.11</scala.binary.version> </properties> <dependencies> <dependency> <groupId>com.ververica</groupId> <artifactId>flink-connector-sqlserver-cdc_${scala.binary.version}</artifactId> <version>2.3.0</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-java</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-table-planner-blink_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency> <dependency> <groupId>org.apache.flink</groupId> <artifactId>flink-clients_${scala.binary.version}</artifactId> <version>${flink.version}</version> </dependency> </dependencies>MySQL的JDBC驱动不需要在pom里显式声明,flink-connector-jdbc会带一个兼容版本,但如果你的MySQL版本是8.0以上且用了caching_sha2_password认证插件,必须手动加一个最新版mysql-connector-java,否则连接会在握手阶段失败。
2.4 把Flink CDC连接器部署到集群
本地跑通了不代表集群上能用,连接器需要打进Flink的lib目录或通过-C参数随作业提交。生产环境我建议用-C方式,只影响当前作业,避免污染共享集群。常见的提交命令长这样:
# 编译项目,跳过测试 mvn clean package -DskipTests # 将连接器依赖和作业Jar一起提交给Flink flink run \ -m yarn-cluster \ -yn 2 \ -C ./target/flink-cdc-sql-server-mysql-1.0-SNAPSHOT.jar \ -c com.example.SqlServerCdcToMysql \ ./target/flink-cdc-sql-server-mysql-1.0-SNAPSHOT.jar如果你的Flink跑在Kubernetes上,思路一样,把连接器Jar挂到/opt/flink/lib或通过pod template的initContainer注入。注意一个坑:连接器的Jar包不能丢到Flink lib目录之后就以为完事了,它的依赖也得分发到所有TaskManager节点,否则作业在JobManager上能启动,一旦分片到TaskManager执行就报ClassNotFound。
3. 建表SQL与参数拆解:核心代码的每一行都值得抠
3.1 Source表定义:SQL Server CDC连接器的关键参数
Flink SQL中定义SQL Server数据源需要写一个完整的CREATE TABLE语句,每个WITH参数都直接影响读取行为。下面是这个项目里典型的source表定义,注释里标了容易踩坑的参数:
CREATE TABLE sql_server_table ( id INT, name STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'sqlserver-cdc', 'hostname' = '192.168.1.100', 'port' = '1433', 'username' = 'cdc_user', 'password' = 'YourPassword', 'database-name' = 'YourDatabase', 'schema-name' = 'dbo', 'table-name' = 'YourSourceTable', 'scan.startup.mode' = 'initial', 'chunk-meta.group.size' = '1000', 'chunk-key.even-distribution.factor.upper-bound' = '1000', 'connect.timeout' = '30s', 'poll.interval' = '1000ms' );逐一说下参数含义:scan.startup.mode有三个可选值,initial表示先做全量快照再增量读取,latest-offset直接从当前最新日志开始,timestamp指定从某个历史时间点开始。日常大多用initial,但要清楚它会先锁表做快照,全量阶段对源库有一定压力。
chunk-meta.group.size和chunk-key.even-distribution.factor.upper-bound是控制全量快照分片粒度的参数。当表数据量过亿时,Flink会按主键范围把快照任务拆成多个chunk并行读取,group.size决定每个批次处理多少个元数据分片,upper-bound用于判断主键分布是否均匀,顺便说一句,这个参数如果配得太大,遇到主键严重倾斜的表会触发拆分过细导致大量小任务堆积。poll.interval是轮询CT表的间隔,单位是毫秒,生产环境不建议低于500ms,太快会把SQL Server捕获作业的IO打满。
schema-name这个参数特别容易漏,在MySQL源里没有这个概念,换到SQL Server必须写,否则连接器不知道去哪个schema下找表。PRIMARY KEY (id) NOT ENFORCED是Flink语法中声明主键的方式,加NOT ENFORCED表示Flink不会校验数据是否违反主键约束,但它会用这个主键做分片和状态管理,所以源表必须有明确主键,否则initial模式的快照逻辑会退化,详细原因后面避坑章节专门说。
3.2 Sink表定义:MySQL JDBC连接器怎么接
MySQL目标表的定义相对直接,connector用的是jdbc,需要注意的参数是sink.buffer-flush.max-rows和sink.buffer-flush.interval,它们控制攒批写入的节奏:
CREATE TABLE mysql_table ( id INT, name STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://192.168.1.200:3306/YourTargetDb?useSSL=false&serverTimezone=Asia/Shanghai', 'table-name' = 'your_target_table', 'username' = 'sync_user', 'password' = 'YourPassword', 'sink.buffer-flush.max-rows' = '1000', 'sink.buffer-flush.interval' = '2s', 'sink.max-retries' = '3' );JDBC连接器默认是UPSERT语义,前提是目标表在MySQL中必须定义主键,且主键字段和source表一致。每攒够1000条或间隔2秒就刷一次缓存,这样能大幅降低MySQL端的写入频率。sink.max-retries控制写失败时的重试次数,但这个重试是同步的,如果MySQL宕机超过3次就直接让作业失败,实际生产中我会调大到5,配合Flink checkpoint来做恢复。
URL里的serverTimezone=Asia/Shanghai必加,MySQL驱动8.0版本对时区敏感,不写的话报Server returns invalid timezone。useSSL=false是测试环境省事,生产内网其实也建议关掉,避免SSL握手增加延迟。
3.3 数据流转与主键策略:INSERT INTO SELECT背后的行为
连接两张表的SQL很简单:
INSERT INTO mysql_table SELECT id, name, create_time, update_time FROM sql_server_table;这行语句看起来就是普通的数据搬运,但Flink CDC在底层做了不少事。initial模式下,先执行全量快照,此时SQL Server源表的UPDATE操作会同时产生两条CT记录(更新前和更新后),Flink会通过状态后端缓存旧值,等全量阶段完成后再合并增量变更。这个过程涉及Flink的checkpoint状态存储,如果状态比较大(几GB),建议用RocksDB状态后端而不是默认的HashMap,否则内存很容易撑爆。
主键在这里的作用是关联快照和增量的唯一键。Flink CDC构建快照分片时会先用SELECT MIN(id), MAX(id)拿到主键范围,然后按范围切分成多个查询任务。如果源表没有主键,连接器会回退到单分片全表扫描,数据量大时就是一场灾难。所以source表定义里PRIMARY KEY (id) NOT ENFORCED不是可写可不写的装饰,而是连接器行为的开关。
NOT ENFORCED还有一层含义是Flink不会去检查数据里是否有重复主键,如果源表本身存在重复主键(这种情况不该发生但确实发生过),Flink在合并多个分片的快照时会因为状态里的主键冲突报错。遇到这种脏数据,处理办法是清洗源表,或者在source表定义中去掉PRIMARY KEY声明,让Flink按rowkind逐条处理而不做状态合并,牺牲一部分精确性换作业稳定性。
3.4 整条作业放进一个Java类:从环境创建到执行
项目里的Java入口类只是把上面的SQL用TableEnvironment包起来,但有几个关键点需要注意。比如环境创建必须指定PLANNER_BLINK和StateBackend:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment(); env.setStateBackend(new EmbeddedRocksDBStateBackend()); env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink-checkpoints/cdc-job"); env.getCheckpointConfig().setCheckpointInterval(10000); env.getCheckpointConfig().setMinPauseBetweenCheckpoints(5000); env.getCheckpointConfig().setTolerableCheckpointFailureNumber(3); StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);代码逻辑分三块:先建表,再执行INSERT INTO,最后调用env.execute()。注意Flink CDC作业必须有checkpoint,因为增量读取的位置(LSN日志序列号)是存在状态里的,如果没有checkpoint,作业一重启就会丢位置——要么从头开始全量快照,要么直接从最新位置开始,中间的变化全部丢失。setCheckpointInterval(10000)是10秒一次,如果业务对重复数据容忍度低,可以压到3秒,代价是状态存储的IO压力变大。
有个细节值得提醒:ExecuteSql里如果同时创建了source表和sink表,表的字段顺序必须完全一致,INSERT INTO SELECT是按位置映射而不是按字段名映射。之前试过把source和sink的字段顺序排错,结果数值和字符串互串,数据写进MySQL后整列错位,排查了很久才发现是建表语句字段顺序不同。
4. 常见问题排查:五个容易翻车的地方与解决办法
4.1 作业启动成功但数据不更新:先查SQL Server Agent
- 现象:Flink作业从
RUNNING变到实际开始消费,但MySQL目标表始终没有新数据,source的CT表也查不到记录。 - 原因:SQL Server Agent服务没启动。CDC捕获作业依赖Agent调度,Agent停了,捕获线程也停了,CT表里根本没写入任何变更。Flink这边看着是正常的,因为它只是周期性轮询CT表,读不到就继续等。
- 解决:到SQL Server所在服务器上启动Agent服务:
# 用管理员权限的PowerShell或cmd执行 net start SQLSERVERAGENT启动后在SSMS里跑EXEC sys.sp_cdc_help_change_data_capture,如果返回结果里有is_cdc_enabled=1且捕获作业的status不为空,基本就恢复了。我踩过的一个变体是:Agent虽然运行,但服务器重启后CDC作业没有自动启动,这时候要手动EXEC sys.sp_cdc_start_job N'cdc.dbo_YourSourceTable_capture'。
4.2 全量快照阶段报decimal精度丢失
- 现象:同步完成后,MySQL中的数值和SQL Server源数据不一致,比如
DECIMAL(18,4)的金额,同步过去变成整数或小数点错位。 - 原因:Flink SQL type system在2.3.0版本对SQL Server的
DECIMAL(p,s)映射逻辑有默认精度推断,如果源表字段是DECIMAL(18,4),而Flink DDL里对应字段声明成了DECIMAL(10,2),Flink会按DDL声明截断。 - 解决:在source表的DDL中把字段精度写成和SQL Server完全一致:
CREATE TABLE sql_server_table ( amount DECIMAL(18, 4), ... ) WITH (...);同时确认MySQL目标表字段也是DECIMAL(18,4),这是最容易修但也最容易忽略的一环,我一般建表前先跑SELECT COLUMN_NAME, DATA_TYPE, NUMERIC_PRECISION, NUMERIC_SCALE FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_NAME = 'YourSourceTable'把元数据全拉出来,按这份清单写DDL。
4.3 源表没有主键导致快照退化成全表扫描
- 现象:
initial模式下,大数据量表的全量快照阶段非常慢,日志里出现单个task长时间运行,且没有分片并行。 - 原因:连接器需要主键来切分快照范围,表没主键时只能全表扫描,所有数据走一个分片。
- 解决:短期方案是用唯一索引替代主键,但要注意唯一索引的列值不能有NULL,因为分片查询会按
WHERE id >= ? AND id < ?构造条件,NULL值会被过滤掉导致数据丢失。长期方案是推动业务在源表上补主键,或者在Flink DDL里用一个不会为空的字段声明PRIMARY KEY。
如果实在没有合适字段,可以用SQL Server的ROWVERSION列辅助,但不建议刚上手就玩这么花的操作,先保证有主键再谈优化。
4.4 MySQL写入报Data too long或Out of range value
- 现象:作业运行一段时间后突然失败,日志里出现
Data truncation: Out of range value for column。 - 原因:SQL Server源表的字段长度和MySQL目标表不一致,通常是源表字段长度更长。比如SQL Server的
NVARCHAR(255)映射到MySQL的VARCHAR(100),超出就被截断。 - 解决:在所有同步链路建表前,先做字段类型和长度映射评审。我的习惯是写一个小脚本,读取SQL Server的
INFORMATION_SCHEMA.COLUMNS,生成MySQL的CREATE TABLE语句,字符串类型的长度统一放大到源表的1.2倍,NVARCHAR全部映射成VARCHAR但长度确保覆盖。项目里的做法是直接在Flink DDL里把字段声明为STRING,让底层自动映射,但生产环境我不会这么干,必须由DBA确认目标表结构。
4.5 检查点失败导致作业不断重启
- 现象:Running状态持续了一段时间后,作业突然自动取消,日志里报
Checkpoint expired before completing或Size of the state is larger than the maximum permissible size。 - 原因:状态后端存储配置不对,或者checkpoint超时时间设得太短。Flink CDC作业的全量快照状态下会保存大量中间状态,默认的checkpoint超时(10分钟)可能不够。
- 解决:调大checkpoint超时,并确认存储路径可用:
env.getCheckpointConfig().setCheckpointTimeout(60000L); env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);同时把state.backend换成RocksDB,并给TaskManager增加内存。这里有个偏方:如果只是测试环境,可以把setTolerableCheckpointFailureNumber调高到5,让作业在checkpoint连续失败几次后还能继续跑,给自己留出排查时间。
5. 进阶用法:并行度调优、增量快照验证与故障恢复
5.1 并行度设计:Source端和Sink端的合理配置
Flink CDC作业的并行度不能无脑加大。Source端的并行度决定了读取SQL Server CT表的并发线程数,但SQL Server的CDC捕获作业是单线程的,Flink并发轮询并不会提升数据产生速度,反而可能造成CT表上的锁竞争。常见做法是source并行度设为1,把并发集中到Sink端。
Sink并行度需要和MySQL写入能力匹配:如果MySQL是单实例,并行度太大会导致连接池被打满;如果是集群或做了读写分离,可以适当提升。经验值是一张普通的千兆网络MySQL,10000 TPS左右的写入量,sink并行度设为2-4就比较合理。调整方式很简单,在INSERT INTO语句里加hint:
INSERT INTO mysql_table /*+ OPTIONS('sink.parallelism' = '3') */ SELECT id, name, create_time, update_time FROM sql_server_table;5.2 验证同步正确性:从增量数据中检查Operation类型
数据同步完成后,怎么确定没丢数据?我的验证方式分三步:先比全量行数,再比增量时间戳,最后抽查差异数据。全量行数对账简单,SELECT COUNT(*)两边跑一下就行。增量验证稍微讲究一点,Flink CDC会在数据中附带一个隐藏字段__$operation,你可以在source表定义里加一行把它透传出来:
CREATE TABLE sql_server_table ( id INT, name STRING, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH (...); -- 注意:这行不是语法正确示例,需要配合流式处理才能访问,但思路是: -- 在SQL Server CT表里通过 __$operation 判断当前事件是插入、更新还是删除用Flink SQL做流式计算时,__$operation不会直接出现在表字段里,但你可以通过Table API的RowKind访问,或者干脆在源表变更时给同步链路加一个审计字段,记录操作类型。具体做法是把source表包一层视图,通过ROW_KIND函数标记每条记录是+I(插入)、-U(更新前)、+U(更新后)还是-D(删除)。
我在生产上用的一个土办法是:在MySQL目标表加一个sync_op字段,用CASE语句把操作类型映射到字符串并写入,这样DBA查数据时一眼就能看出来哪条是删除哪条是更新。
5.3 从checkpoint恢复:让同步不丢不重
Flink CDC的断点续传能力是所有实时同步方案里我最看好的一点。假设MySQL目标表某天半夜挂了,Flink作业跟着失败,修复MySQL后,不需要手工启动一个同步工具去补数据,直接恢复Flink的job即可。前提是SQL Server Agent还在跑,CDC的CT表数据还在保留期内。恢复命令:
# 从最近一次成功的checkpoint恢复 flink run -s hdfs:///flink-checkpoints/cdc-job/<checkpoint-uuid> \ -c com.example.SqlServerCdcToMysql \ ./target/flink-cdc-sql-server-mysql-1.0-SNAPSHOT.jar恢复后Flink会从上次记录的位置继续读CT表,这段期间的数据变更不会丢,但因为checkpoint是间隔触发的,恢复位置和真实故障时间之间可能有几秒的数据重复,需要下游幂等处理。MySQL的JDBC连接器是UPSERT写入,主键相同就覆盖,天然适配这种场景,这也是我建议目标表必须设主键的原因之一。
5.4 一个值得养成的习惯:上线前先做全链路演练
资源文件里那份Maven工程我在多个环境跑过几次,积累了一个固定流程:先在本地起一个SQL Server容器、一个MySQL容器,用测试表模拟增删改,跑通后再切到真实业务表。每次切真实表前,我会强制走一遍这几步:确认SQL Server Agent运行、确认目标表有主键、确认源表有主键、确认Flink作业的checkpoint目录可写、确认MySQL账号权限。这五个确认看起来繁琐,但每一个都在实际项目中绊过我。尤其是Agent和主键这两个,都是那种「作业启动成功、数据就是不动」的隐形坑,排查起来极其消耗耐心。从那以后我每次搭同步链路都强制走完这套确认流程再上线,希望这份笔记能帮你少走同样的弯路,祝你的数据同步一次跑通。
本文还有配套的精品资源,点击获取