多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Flink CDC实战:SQL Server到MySQL数据实时同步方案全解析

Flink CDC实战:SQL Server到MySQL数据实时同步方案全解析 简介面向大数据开发与数据集成工程师提供一套基于 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以CDCChange 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必须先对目标表开启CDCFlink才能读到数据如果跳过这一步直接跑作业日志里只会出现连接成功但没有任何数据输出的怪象。开启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 Ndbo, source_name NYourSourceTable, role_name NULL, capture_instance Ndbo_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 NQUERYSTATE, NSQLServerAGENT;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.2pom里需要显式把flink版本属性对齐properties flink.version1.14.2/flink.version scala.binary.version2.11/scala.binary.version /properties dependencies dependency groupIdcom.ververica/groupId artifactIdflink-connector-sqlserver-cdc_${scala.binary.version}/artifactId version2.3.0/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-java/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-table-planner-blink_${scala.binary.version}/artifactId version${flink.version}/version /dependency dependency groupIdorg.apache.flink/groupId artifactIdflink-clients_${scala.binary.version}/artifactId version${flink.version}/version /dependency /dependenciesMySQL的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?useSSLfalseserverTimezoneAsia/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里的serverTimezoneAsia/Shanghai必加MySQL驱动8.0版本对时区敏感不写的话报Server returns invalid timezone。useSSLfalse是测试环境省事生产内网其实也建议关掉避免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和StateBackendStreamExecutionEnvironment 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_enabled1且捕获作业的status不为空基本就恢复了。我踩过的一个变体是Agent虽然运行但服务器重启后CDC作业没有自动启动这时候要手动EXEC sys.sp_cdc_start_job Ncdc.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是单实例并行度太大会导致连接池被打满如果是集群或做了读写分离可以适当提升。经验值是一张普通的千兆网络MySQL10000 TPS左右的写入量sink并行度设为2-4就比较合理。调整方式很简单在INSERT INTO语句里加hintINSERT 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和主键这两个都是那种「作业启动成功、数据就是不动」的隐形坑排查起来极其消耗耐心。从那以后我每次搭同步链路都强制走完这套确认流程再上线希望这份笔记能帮你少走同样的弯路祝你的数据同步一次跑通。本文还有配套的精品资源点击获取
返回列表