ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

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; GOsp_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 \ -Dartifactcom.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.x5.7 反而容易在 SSL 握手上报错后面的连接串里统一显式加useSSLfalse内网环境少一层加密开销也能避开不少版本兼容问题。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?useSSLfalserewriteBatchedStatementstrue, 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 里useSSLfalse的前提是 MySQL 端没有强制 SSL否则写着 false 也会报 SSL 连接错误这是 MySQL 8.0 的常见新坑rewriteBatchedStatementstrue对多行批量写入提升非常明显增量高峰每分钟几千行时没有这个参数 batch 提交效率会掉一个量级。sink.buffer-flush.max-rows控制攒多少行刷一次sink.buffer-flush.interval控制最多等多长时间两个条件满足其一就触发 flush。我一般把 max-rows 放在 1000 到 2000interval 放在 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 建不了列需要提前 CASTdatetime2TIMESTAMP(3)datetime(3)毫秒以下精度会丢讲究的话源表收 STRINGdatetimeoffsetSTRINGvarchar(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 里对应字段声明成 STRINGsink 也建 varchar(30)精度原样保留代价是下游 SQL 要做一次日期转换。uniqueidentifier 同理直接当字符串走最省事省得两端类型不一致导致写不进去。4.3 jdbc sink 的三个参数决定写入延迟sink 侧除了前面提过的 flush 两件套还有一个容易被忽略的关键点连接串上的配置。MySQL 端默认连接数 151过高的并行度会把连接池打爆出现Too many connections。我把 JDBC sink 的并行度控制在 2 到 4配合连接池上限 200 左右即可。另一个可调选项是sink.max-retries3避免偶发网络抖动直接把任务搞失败。调完参数后再看数据延迟在 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 Encryptionmssql-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 慢查询拖住。再配合上面两个核对脚本实时同步就能做到半夜不接电话。希望帮到你。本文还有配套的精品资源点击获取
返回列表