公司业务库有一百多张表,报表和数据团队天天催实时数据。以前我都是写 DataX 离线任务,每天凌晨跑批,T+1 的延迟在业务那边已经被反复吐槽。后来一个朋友给我推荐了 Dinky 的 CDCSOURCE 方案,把 MySQL 整库同步到另一个 MySQL 实例,一条 SQL 就能把源库的表实时搬到目标库,不用为每张表单独写 source、sink,省了大量重复劳动。
我用这套方案已经跑了大半年,从最初的调研、验证到生产落地,中间踩过不少坑:版本不兼容、binlog 参数没开、server-id 冲突、目标端表建不出来……这篇文章把完整的配置过程、参数含义和排错经验整理出来,给准备做 MySQL 实时同步的同学一个可以直接照抄的参考。
不管你是刚接触 Flink CDC 的新手,还是已经在用 Dinky 想盘整库同步的老手,这套方案的思路和避坑点都有参考价值。
1. 方案整体设计思路与适用场景
1.1 CDCSOURCE 是怎么做到"一条 SQL 同步整库"的
先说明最核心的一点:Dinky 是一个基于 Flink SQL 的开发运维平台,而 CDCSOURCE 是它内置的一种扩展语法。一条 EXECUTE CDCSOURCE 语句,在后台会自动完成"读取源端数据库元数据、为每张表生成 source 表与 sink 表、把整库的表统一注册并启动实时同步"这一整套流程。对使用者来说,你需要写的东西只有一条 SQL。
底层核心是 Flink CDC。Flink CDC 基于 Debezium 实现,它会伪装成 MySQL 的从节点,订阅 binlog,拿到行级别的增删改事件。Debezium 捕获到的是标准化的 change event,再经过 Flink CDC 转成 Flink 能消费的数据流。CDCSOURCE 做的事情,就是把这一能力封装成一句 SQL,同时解决"每张表都要手写 DDL 和 insert into select"的繁琐问题。
用一个生活化的类比来解释:手动写 Flink CDC 相当于给每家门店单独装修、单独进货、单独配店员;CDCSOURCE 相当于开连锁超市,总部统一盘点商品清单(元数据)、统一装修(建表)、统一配货(同步)。它没有引入新的同步引擎,本质上是把 Flink SQL 的开发过程"模板化"了,这才是它省事的根因。
1.2 为什么选 CDCSOURCE 而非手写 Flink CDC 任务
很多人看到这里会问:我自己写 Flink CDC 也能做到,凭什么要用 Dinky 的 CDCSOURCE?我把两种方式的差异拉了一个对比表,大家感受会更直观。
| 对比维度 | 手动 Flink CDC | Dinky CDCSOURCE |
|---|---|---|
| 任务编写 | 每张表都要写 CREATE TABLE 和 INSERT INTO,100 张表就是 300 条 SQL | 一条 EXECUTE CDCSOURCE 搞定 |
| 表结构变更 | 源表加了字段,要手动改任务里的 DDL 再重启 | 启动时自动读取元数据建表,人工介入少 |
| 运维管理 | 需要额外写提交脚本、管理 checkpoint 和 savepoint | Dinky 后台统一管理任务、日志、保存点 |
| 扩展性 | 灵活,二次加工随便写,适合复杂场景 | 语法封装完整,复杂加工需要退出模板写原生 Flink SQL |
| 上手门槛 | 需要对 Flink 开发比较熟 | 只要懂 SQL 就能配置 |
这套方案的适用场景很明确:表多、单表逻辑不复杂、目标端表结构能对齐源端。比如从 MySQL 同步到另一个 MySQL 报表库、从 MySQL 同步到 Doris / StarRocks / ClickHouse 数仓。反过来,如果涉及多表 join、字段清洗、bean 处理、窗口聚合这些加工逻辑,我不建议硬套 CDCSOURCE,那种情况老老实实回到 Flink SQL 手动开发更可控。
还有一点要注意,CDCSOURCE 虽说是"整库同步",但它对运行过程中新增表、源端 ALTER TABLE 这类结构变更的自动感知能力,跟 Dinky 版本、Flink CDC 版本有很大关系。这个后面我会单独展开,因为这是我踩过最深的坑。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 环境准备:版本选型、binlog 配置与 Dinky 部署
2.1 版本到底怎么配才不踩坑
先说结论,我生产环境实际使用的是这样一套组合,跑了半年没有遇到兼容性问题:
- Dinky 1.0.3
- Flink 1.17.2
- flink-sql-connector-mysql-cdc-2.4.1.jar
- flink-connector-jdbc-3.1.1-1.17.jar
- mysql-connector-java-8.0.30.jar
为什么这么选?有几个原因。
Flink CDC 2.4.1 对 Flink 1.17 的适配非常成熟,而且支持按主键分片并行读取快照,也就是 chunk 分片机制。这个机制对大表初始同步很关键,并行度设上去之后,一张几千万行的表能通过多个 chunk 同时拉取,速度提升明显。
Flink CDC 3.x 虽然引入了官方原生 pipeline 语法,整库同步更"正统",但它的依赖组织和 Dinky 的适配相对较新,Dinky 不同版本对 3.x 的支持情况不一样,新手一上来直接用 3.x 很容易被依赖问题劝退。如果你先跑通 2.4.x 方案,后续再迁移到 3.x 会平滑很多。
另外要注意,Dinky 的 lib 目录和 Flink 的 lib 目录里不要放重复或冲突的 jar。比如两个不同版本的 flink-connector-jdbc 同时存在,提交任务时很可能出现方法签名冲突,报错极其难查。
2.2 源端 MySQL 必须改的 binlog 参数和账号授权
CDCSOURCE 要读源端变更,源库 MySQL 必须开启 binlog 且格式为 ROW。我见过很多同学配置到一半报错,根源就是 binlog 没开或者格式不对。
修改 MySQL 配置文件 my.cnf(Windows 下是 my.ini),在 [mysqld] 段加:
ini复制[mysqld]
server-id = 1
log-bin = mysql-bin
binlog_format = ROW
binlog_row_image = FULL
expire_logs_days = 7
逐个解释一下这些参数为什么重要:
- server-id:MySQL 集群内唯一。Flink CDC 伪装成从库连接主库时,如果多个任务用了相同 server-id,主库会认为是同一个从库重连,产生冲突报错。
- log-bin:开启 binlog 日志,文件名前缀是 mysql-bin。
- binlog_format = ROW:必须设为 ROW 行级格式。Debezium 靠这个才能解析出每一行的变更前值和变更后值;STATEMENT 或 MIXED 格式下拿不到完整的行级变化。
- binlog_row_image = FULL:记录所有字段的镜像,保证 update 事件里 before 和 after 都是完整数据。默认值其实已经是 FULL,但有的云数据库或历史配置会不一样,建议显式写上。
- expire_logs_days:根据你的同步链路重要程度设置 binlog 保留时间。太短会导致任务故障恢复时找不到历史 binlog 位置。
改完配置后重启 MySQL,然后执行下面几条 SQL 验证:
sql复制SHOW VARIABLES LIKE 'log_bin';
SHOW VARIABLES LIKE 'binlog_format';
SHOW BINARY LOGS;
看到 log_bin 为 ON、binlog_format 为 ROW,就说明 binlog 侧就绪。
源库同步账号也需要特殊权限,不能只给 SELECT。我习惯创建独立的 CDC 账号,避免复用业务账号:
sql复制CREATE USER 'cdc_user'@'%' IDENTIFIED BY 'StrongPass_2024';
GRANT SELECT, RELOAD, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO 'cdc_user'@'%';
FLUSH PRIVILEGES;
这里 SELECT、SHOW DATABASES 是读取初始快照和元数据用的,REPLICATION SLAVE 和 REPLICATION CLIENT 是 Debezium 订阅 binlog 的必备权限。少了 REPLICATION 相关权限,任务启动时会报 binlog 连接被拒绝或者权限不足一类的错误,后面排查章节会细说。
2.3 Dinky 安装并挂载 Flink 实例
Dinky 本身是 Java 应用,解压部署后连接元数据库进行初始化,再注册 Flink 集群实例就能用了。我快速梳理一下操作顺序。
第一步,下载 Dinky 安装包,解压到 /opt/dinky。解压后目录里有 bin、lib、config 等目录,不要急着启动,先把配置改了。
第二步,修改 config 目录下的 application.yml,把 Dinky 自身的元数据库配置成你自己的 MySQL。元数据库用于保存任务定义、集群配置、保存点信息等。Dinky 启动时如果检测到库里有表就不会重复初始化,首次启动会自动建表。
第三步,把前面提到的几个 jar 包放进 Dinky 的 lib 目录:flink-sql-connector-mysql-cdc、flink-connector-jdbc、mysql-connector-java。不同版本的 Dinky 放置 jar 的位置可能略有差异,如果是带插件中心的版本,也可以通过后台界面上传。
第四步,启动 Dinky。
bash复制cd /opt/dinky/bin
sh start.sh
第五步,登录后台,在"注册中心"里注册 Flink 实例。调试阶段我推荐先注册一个 Standalone Session 集群:先去 Flink 目录启动一个 session,bin/start-cluster.sh,然后在 Dinky 后台填 Flink Rest 地址,比如 http://localhost:8081。生产环境建议用 per-job 或 application 模式,每个任务独占资源,互相不干扰,也方便从 savepoint 精确恢复。
第六步,注册完成后在 Dinky 里能测试连接,看到集群心跳正常,环境就通了。
这里有一个实操细节:Dinky 元数据库和同步目标库最好不要共用同一个 MySQL 实例,避免 Dinky 自己的元数据变更也被 CDC 捕获,造成不必要的循环。
3. CDCSOURCE 核心语法与参数逐个拆解
3.1 一段能直接跑的整库同步语句
先给一段完整的语句,我以 MySQL 到 MySQL 为例:
sql复制EXECUTE CDCSOURCE mysql2mysql_demo
WITH (
'connector' = 'mysql-cdc',
'hostname' = '127.0.0.1',
'port' = '3306',
'username' = 'cdc_user',
'password' = '123456',
'source-db' = 'mall',
'sink.connector' = 'jdbc',
'sink.db' = 'mall_report',
'sink.servers' = 'jdbc:mysql://192.168.1.50:3306/mall_report?useUnicode=true&characterEncoding=utf8&useSSL=false',
'sink.username' = 'sync_user',
'sink.password' = '123456',
'sink.driver' = 'com.mysql.cj.jdbc.Driver',
'table-name' = 'ods_.*',
'parallelism' = '2',
'checkpoint' = '10000',
'scan.startup.mode' = 'initial'
);
注意,这个示例里 source-db 和 sink.db 写的是不同名字,这是 CDCSOURCE 完成库名映射的方式。如果源库和目标库同名,直接把两个参数写成一样即可。
3.2 关键参数解析:source 端、sink 端、调度与恢复
我把这段 SQL 里的参数按职责拆成三类来理解:source 源端、sink 目标端、任务控制。
先看 source 端参数。
| 参数 | 含义 | 说明 |
|---|---|---|
| connector | 源端 CDC connector 类型 | mysql-cdc / postgres-cdc / oracle-cdc 等 |
| hostname / port | 源库地址与端口 | 写实际连通地址,别写 localhost,容器部署时尤其注意 |
| username / password | 源库同步账号 | 使用前面创建的 cdc_user,权限不能是普通只读账号 |
| source-db | 源库名 | CDCSOURCE 会读取这个库的所有表结构 |
| table-name | 表名匹配正则 | 支持 Java 正则,默认 .*,灵活控制同步范围 |
| database-name | Flink CDC 原生 database-name 参数 | 部分版本用于和 source-db 配合,一般保持一致 |
| scan.startup.mode | 启动位置 | initial 先全量后增量;latest-offset 只从当前 binlog 最后位置开始;earliest-offset 从最早 binlog 开始 |
再来看 sink 目标端参数。
| 参数 | 含义 | 说明 |
|---|---|---|
| sink.connector | 目标端连接器 | 这里用 jdbc,同步到 Doris 可换成 doris,StarRocks 换 starrocks |
| sink.db | 目标库名 | 和 source-db 不同名时就是库名映射 |
| sink.servers | 目标端 JDBC URL | 注意带上 useSSL=false 和 utf8 参数,减少连接问题 |
| sink.username / sink.password | 目标端账号 | 需要有目标库建表、写入、删改的权限 |
| sink.driver | JDBC 驱动类 | MySQL 8 用 com.mysql.cj.jdbc.Driver |
最后是任务级控制参数。
| 参数 | 含义 | 说明 |
|---|---|---|
| parallelism | 并行度 | 影响 source 分片并发和下游写入并发,需要结合源库负载调整 |
| checkpoint | checkpoint 间隔(毫秒) | 决定故障恢复的粒度,10000 表示每 10 秒做一次快照 |
scan.startup.mode 是很多人忽略但很关键的参数。第一次跑全量加增量同步,必须用 initial,它会先把历史存量数据读出来,再无缝切换为 binlog 增量。如果你已经通过其他方式迁完存量数据,只需要追增量,那用 latest-offset 能从当前 binlog 位置开始,省掉全量扫描。这里有个容易犯的错误:存量还没同步完就设 latest-offset,导致目标端丢数据,全量数据永远不会补回来。
parallelism 这个参数,在 Flink CDC 2.4.x 里会触发分片读取机制。源表有主键时,CDC 会按主键把表拆成多个 chunk 并行读;没有主键的表回退到单 chunk 串行读,速度会明显慢。所以配置并行度前,最好先确认源表主键情况。
3.3 表匹配、库名映射与 DDL 同步的处理逻辑
CDCSOURCE 的表匹配逻辑基于正则。我例子里的 ods_.* 表示只同步以 ods_ 开头的表,比如 ods_order、ods_order_item。如果你确定要同步整库所有表,把 table-name 设为 .* 即可。这个参数用正则的另一个好处是,如果业务方要求只同步订单相关表,写 .*(order|trade).* 就行,非常灵活。
库名映射是 CDCSOURCE 里比较简单的逻辑:source-db 是源库名,sink.db 是目标库名。两条语句里分别配置后,所有从 mall 库选中的表都会自动落到 mall_report 库里,表名保持不变。这个映射对于"同一份数据源同时同步到多个目标场景"特别有用,你可以灵活自由调整目标库。
目标表结构是怎么来的?CDCSOURCE 在任务启动时会读取源端元数据,为每个匹配的表生成镜像建表语句,自动在目标库创建同名表。这意味着目标端连接账号必须有 CREATE 权限。生产环境建议在目标端预先规划好表空间、字符集和索引策略,因为 CDCSOURCE 自动建出来的表,默认情况下只保留字段类型,不会自动复制源端的索引和主键。这一点非常影响同步后的查询性能,目标端是报表场景时尤其要注意,需要在任务跑起来之后补充索引。
DDL 同步是另一个重点。Flink CDC 2.x 对源端 DDL 事件的默认处理方式是:把 schema change 作为事件流转出来,但 Dinky CDCSOURCE 对 DDL 是否执行、以什么形式执行,不同版本行为不一样。我使用的 1.0.3 + CDC 2.4.1 组合,源端 ALTER TABLE 新增字段,目标端表结构不会自动变更。所以我的建议是:生产环境不要把"运行中改表结构"当成 CDCSOURCE 的能力,表结构变更应该走专门的维护流程,比如夜间窗口改完表结构后,从 savepoint 重启任务让 CDCSOURCE 重新读取元数据。
4. 实操全流程:从订单库整库同步到报表 MySQL
4.1 这次要解决的问题和具体配置
我这边的实际场景是:源库 127.0.0.1:3306 上有个 mall 库,里面是订单中心相关的 32 张表,包括订单主表、订单明细、支付流水、退款记录等。目标端是另一台机器 192.168.1.50:3306 上的 mall_report 库,报表系统从这里取数。
需求很简单:源库任意表发生增删改,报表库要在秒级看到。报表团队原来一直靠凌晨跑批脚本同步,数据延迟一天,业务方等不及。
我在 Dinky 里新建的 FlinkSQL 任务,实际用的语句是这样:
sql复制EXECUTE CDCSOURCE mall_to_report
WITH (
'connector' = 'mysql-cdc',
'hostname' = '127.0.0.1',
'port' = '3306',
'username' = 'cdc_user',
'password' = 'xxxx',
'source-db' = 'mall',
'sink.connector' = 'jdbc',
'sink.db' = 'mall_report',
'sink.servers' = 'jdbc:mysql://192.168.1.50:3306/mall_report?useUnicode=true&characterEncoding=utf8&useSSL=false&rewriteBatchedStatements=true',
'sink.username' = 'sync_user',
'sink.password' = 'xxxx',
'sink.driver' = 'com.mysql.cj.jdbc.Driver',
'table-name' = '.*',
'parallelism' = '4',
'checkpoint' = '10000',
'scan.startup.mode' = 'initial',
'debezium.snapshot.fetch.size' = '8192'
);
这里相比前面的示例多了两个参数,一个是 rewriteBatchedStatements=true,一个是 debezium.snapshot.fetch.size=8192。
rewriteBatchedStatements 是 JDBC 连接参数,开启后能让批量写入以 multi-values 形式发给 MySQL,写入性能能提升一个量级。做数据同步的人应该都知道,这个参数不加,默认的批处理其实是假批量,一条一条执行,速度上不去。
debezium.snapshot.fetch.size 控制全量快照阶段每次从源端读取的行数。默认值偏保守,调大后全量阶段的行读取批次更大,快照能快不少。
4.2 在 Dinky 上新建任务并提交运行
在 Dinky 后台的实际操作流程不复杂。
登录后进到"数据开发"模块,新建一个 FlinkSQL 任务,名字随意,比如 mall_2_report_sync。编辑器里粘贴上面那一段 EXECUTE CDCSOURCE 语句。
然后选择集群实例。调试阶段我会先选 session 集群,点"执行"按钮。这里要注意 Dinky 里"执行"和"提交"的差异——执行是在当前页面调试运行,能够快速看到错误日志;提交是把任务正式 detach 到集群。第一次跑通之前建议先点执行,便于观察日志。
任务进入 RUNNING 状态后,切到 Flink Web UI,能看到 Dinky 自动拆出的所有子任务。你会看到每个源表都有对应的 source 算子、transform 算子和 jdbc sink 算子,数量很多但结构规整,这正是 CDCSOURCE 把模板化的 SQL 提交到集群后的正常表现。
从我的实测来看,32 张表全量快照阶段大约 8 分钟跑完,之后进入持续的 binlog 增量消费阶段。这里我要提醒一句,全量阶段如果表特别多,source 端负载会明显上升,MySQL CPU 和网络 IO 都会有一定波动,生产大库建议在业务低峰期做首次任务,或者先用 table-name 正则小范围试跑。
4.3 数据验证:增删改都能实时反映到目标端吗
任务跑起来后,验证工作非常关键。我会按以下清单逐项确认。
先确认全量数据完全一致:
sql复制SELECT COUNT(*) FROM mall.orders;
SELECT COUNT(*) FROM mall_report.orders;
两张表记录数一致只能说明基本同步成功,还要抽查特定字段内容是否一致,特别是 decimal、datetime、varchar 这类精度敏感字段。
然后测增量。在源库插入一条订单:
sql复制INSERT INTO mall.orders(order_id, user_id, amount, status, create_time)
VALUES (999999, 10001, 199.90, 'PAID', NOW());
等一两秒,去目标库查:
sql复制SELECT * FROM mall_report.orders WHERE order_id = 999999;
我这里实测延迟在 1 秒以内,任务跑得越久延迟越稳定,因为 checkpoint 之后,sink 提交是准实时的。
再测 update:
sql复制UPDATE mall.orders SET status = 'REFUND' WHERE order_id = 999999;
目标库对应记录的 status 字段很快变成 REFUND。
最后测 delete,源库删除后,目标库记录也被删掉。这一步验证很重要,说明 CDCSOURCE 跑的是真正的 CDC 增量同步,target 端所有数据变更都是完整的事件流,不是简单 append 一次快照就结束。
增量逻辑验证通过后,我会再做一张历史大表的针对性测试,比如源库里几张千万行的表,确保 chunk 分片并行读没有把数据读丢。常见问题就是快照阶段多冲合一,分片边界处理不当导致丢数据,用大表验证最稳妥。
4.4 checkpoint 与断点续传:重启不丢数据
同步任务跑起来之后,真正考验它的是故障恢复能力。比如我某次因为目标端 MySQL 重启导致写入连接断开,Flink 任务会自动重启,但因为配置了 checkpoint,它会从最近一次 checkpoint 状态恢复 binlog offset,不会把断流期间的变更数据丢掉。
实际操作中断点续传的步骤是这样的。
在 Dinky 的任务运维页面找到"SavePoint"按钮,手动触发 savepoint。Dinky 会把当前任务的状态保存下来,包括所有算子状态和 binlog 消费 offset。任务暂停后,如果你需要改同步表范围、调并行度、或者源端数据库做了维护,就可以在任务配置里选择"从 SavePoint 恢复",重新提交。
这里有一个容易踩坑的地方:修改并行度之后从旧 savepoint 恢复,如果并行度的变化导致算子状态无法重新映射,任务会启动失败。所以我的习惯是,只要改并行度,宁愿重新走一次 initial 全量同步,也不要硬从旧 savepoint 恢复。反过来,如果只是新增表或者改表名正则,并行度不变,那从 savepoint 恢复基本没问题。
有一个运维经验分享给各位:Dinky 的 savepoint 和 checkpoint 不要混为一谈。checkpoint 是 Flink 周期性的状态快照,用于宕机自动恢复;savepoint 是管理员主动触发的业务快照,用于任务升级、迁移、范围调整。生产环境每次版本迭代或者表范围变更,都要手动打一个 savepoint,再操作变更,这比你完全依赖 checkpoint 可靠得多。
5. 常见问题与排查技巧实录
5.1 连不上源库:权限、网络、host 三个坑
任务提交后如果立刻报错,先别怀疑 CDC 配置,大概率是连接源库这一步出了问题。
典型的报错是:
text复制java.sql.SQLException: Access denied for user 'cdc_user'@'%' (using password: YES)
这个报错的直接原因是权限不足。注意我们前面创建的账号虽然给了 REPLICATION 权限,但如果你用的是普通业务账号,只开了 SELECT,那启动初始快照阶段能通过,但到了 binlog 订阅阶段就会报权限错误。解决办法是确认账号权限:
sql复制SHOW GRANTS FOR 'cdc_user'@'%';
确保包含 REPLICATION SLAVE 和 REPLICATION CLIENT。
另外一个常见问题是网络不通。Dinky 服务所在机器和源库 MySQL 之间的网络不通,或者云安全组、防火墙规则拦了 3306 端口。这种情况下报错通常是 Communications link failure,解决思路很直接:在 Dinky 所在机器上用 mysql -h 源库IP -P 3306 -u cdc_user -p 手动连一下,能连通再谈后续。
还有一个容易忽略的 host 问题:Dinky 配置里 hostname 写 localhost,实际源库在其他机器,导致走回了本地 socket。这个问题在容器部署场景尤其常见,检查时把源库 IP 显式写上,别偷懒。
5.2 读取 binlog 报错:server-id 冲突和 binlog 不存在
binlog 相关报错是整库同步方案里最让人头大的问题。
第一个高频报错内容是:
text复制ERROR: A slave with the same server_uuid/server_id as this slave has connected to the master
这意味着同时有多个 Flink CDC 任务在连接同一个 MySQL 主库,但都用了相同的 server-id。Flink CDC 的 server-id 默认如果不配,会随机生成,但可能和其他任务撞上。解决办法是在每个 CDCSOURCE 语句里显式指定唯一值:
sql复制'debezium.server.id' = '5401'
不同任务用不同编号:5401、5402、5403。这样主库才能把它们当作不同的从库来对待。
第二个报错是:
text复制Could not find first log file name in binary log index file
这个往往是因为 scan.startup.mode 配了 earliest-offset,但 MySQL 里历史 binlog 已经被 expire_logs_days 清掉,导致 CDC 找不到最早的日志文件。解决办法是改用 initial 或 latest-offset 启动,同时合理设置 binlog 保留时间。如果业务允许,直接把启动模式切回 initial 最省事,全量重来一遍保证数据完整。
5.3 任务启动成功但没建表没数据
这种问题比报错更隐蔽,因为任务状态是 RUNNING,看起来一切正常,但目标库一张表都没建,数据更是没有。
我从实际排查中总结出三个主要原因。
第一个原因是 sink.db 对应的目标库在目标实例上不存在,而 CDCSOURCE 自动建库需要 CREATE DATABASE 权限。如果同步账号没有这个权限,建库失败,但任务整体不会报致命错误。解决思路是提前在目标实例上把空库建好,比如 CREATE DATABASE mall_report DEFAULT CHARSET utf8mb4,然后再启动同步任务。
第二个原因是 table-name 正则写错。比如源库的表名是 mall_order,你写 ^ods_.* 去匹配,一张表都匹配不上,CDCSOURCE 变成了空跑。排查方法是在 Dinky 执行前先用 SQL 在源库里验证正则:
sql复制SELECT TABLE_NAME FROM information_schema.TABLES WHERE TABLE_SCHEMA = 'mall' AND TABLE_NAME REGEXP '^ods_.*';
看返回结果是不是你想要的表集,确认没问题再提交任务。
第三个原因是 Dinky lib 目录缺了目标端连接器。任务能启动,但 sink 算子初始化时报 ClassNotFound,日志里会看到 flink-connector-jdbc 相关的类找不到。解决方法是把 flink-connector-jdbc 和对应数据库驱动 jar 放进 Dinky lib 目录后重启 Dinky,再重新提交任务。
5.4 同步速度慢、大表卡住怎么办
整库同步跑完存量后,增量阶段一般很稳,但初始快照阶段如果遇到大表,性能就是瓶颈。
碰到大表卡住,我第一个动作是看源表的索引结构。Flink CDC 2.4.x 的分片算法依赖主键或唯一键,如果表连主键都没有,并行度再高也白搭,只能单 chunk 顺序读。这种情况下最直接的优化是给源表补一个主键,或者用一个单一唯一索引的字段做分片依据。
第二个手段是调并行度。parallelism 从 1 调到 4,分片数量会相应增加,多线程并行拉取,全量快照速度接近线性提升。但并行度不是越大越好,一定要盯着源库的 CPU 负载,别把业务库打挂了。我一般从 2 开始试探,观察源库负载和任务吞吐,再逐步往上涨。
第三个手段是前面提到的 debezium.snapshot.fetch.size。把它从默认调大到 8192,单次从源库拉取的行数变多,网络往返次数减少,对全量速度提升很明显。增量阶段这个参数用不上,但不会造成副作用。
还有一个容易忽略的优化点:目标端写入的 rewriteBatchedStatements 参数。就算 source 端读得再快,目标端写入不行也是白搭。这个参数开启后批量写入性能差异在压测里能差出 5 到 10 倍,务必加上。
5.5 目标端表结构变更与 DDL 事件处理
前面反复提到 DDL 问题,这里集中展开一下。
CDCSOURCE 在任务启动时会读源端元数据生成目标表结构,但任务运行期间,如果源表执行 ALTER TABLE 新增字段、修改字段类型、加索引等操作,目标端表结构不会自动跟着变。这个行为在我的 Dinky 1.0.3 版本上是这样的,新版 Dinky 对 schema evolution 的支持会更好,但生产环境不能赌,必须按"不支持自动 DDL 同步"来设计流程。
我的处理流程是这样的:
第一,源端表结构变更前,在目标端同步执行一次相同的 DDL,保持两边结构一致,之后任务继续跑。这是最省事的办法,但容易因为两边先后顺序问题导致目标端写入失败。
第二,或者利用维护窗口,在源库变更完后,手动给 Dinky 任务打 savepoint,停任务,再重新从 savepoint 启动。任务重启时会重新读取源端元数据吗?实测表明,CDCSOURCE 在启动时会重新获取源表结构,所以重启后会自动把新字段补上。这也侧面说明,从 savepoint 重启比直接恢复 checkpoint 更可靠,因为 savepoint 保存的是业务状态,表结构元数据会在启动时刷新。
第三,如果是加字段这种小变更,我在实际中更喜欢这两种方法组合:先在目标库把新字段加上,接着重启 CDCSOURCE 任务,让它重新对齐元数据。这样能避免写入过程中目标端表结构比写入的 schema 旧而报错。
强调一个细节:目标端自动建表不会复制源端主键和索引。同步完成后,报表查询如果经常按某些字段过滤,一定要自己补索引,否则数据同步没问题但查询能慢到让你怀疑人生。
6. 实操心得与经验补充
6.1 生产上线前建议先做这几件事
这套方案虽然上手简单,但要上生产,我建议不要跳过下面这些准备工作。
先小范围试跑。不要一上来就同步全部 100 张表,先选 2 到 3 张核心表,把前面说的全量、增量、update、delete、重启恢复全部验证一遍再扩范围。我当年就是因为想一步到位,结果 100 张表里混杂了几张没有主键的历史表,快照阶段卡了将近一小时,又不敢停任务,只能干等。
再核对字段类型的一致性。MySQL 到 MySQL 的映射一般没问题,但 unsigned bigint 到有符号 bigint 可能溢出,datetime 的时区处理也可能有偏差,json 类型在不同版本表现也不一样。先抽样核对目标端数据,别只看行数。
最后设置监控告警。Dinky 任务失败、Flink 作业重启、长时间无 checkpoint、同步延迟超过阈值,都应该有提醒。至少保证同步任务停了之后有邮件或钉钉通知,不然用户比你先发现数据断了,那感觉很难受。
6.2 成本与替代方案的一些思考
很多人拿到 CDCSOURCE 后容易走极端:要么觉得它万能,要么觉得它鸡肋。我的观察是,它就是定位在"快速起步、省人工"的整库同步工具,适合自己的场景才叫好方案。
如果目标只是做一个低成本只读从库,MySQL 原生主从复制更成熟、延迟更低,没必要引入 Flink 全家桶。如果需要同步到 Doris、StarRocks 这种分析型库,CDCSOURCE 一样支持,把 sink 参数换掉即可,Doris/StarRocks 的 sink 对批量导入的优化非常强。
如果同步链路里需要复杂加工,CDCSOURCE 的模板就不够用了,这种场景我的建议是:用 CDCSOURCE 做原始表的实时落地,再用 Flink SQL 或 ClickHouse 物化视图在目标端做加工,各司其职。这比强行在一个 CDCSOURCE 里做清洗要灵活得多。
6.3 最后说一个提升幸福感的小技巧
由于 CDCSOURCE 自动建表不复制主键和索引,我用一个小技巧降低了大量手工操作:提前把建表模板沉淀到自动化脚本里。具体做法是,在源端通过 SHOW CREATE TABLE 拿到表结构,用脚本自动把 ENGINE 改为 InnoDB、加入主键、修正字符集为 utf8mb4,然后先在目标端批量建好表,再启动 CDCSOURCE,把 sink.db 指向这个已准备好的库。这样既能享受 CDCSOURCE 自动同步的红利,又能保证目标端表结构和索引完全在自己的控制中,同步任务本身也更稳定。
我第一次在 Dinky 里提交 CDCSOURCE 的时候,最大的感受是省下来的工作量确实明显,但也别把它当成万能钥匙。它适合的是表多、单表逻辑不复杂、目标端结构能对齐源端的同步场景,一旦涉及二次加工,还是老老实实回去写 Flink SQL。这套工具把我从"每张表手写一条 SQL"的重复劳动里解放了出来,把精力放到了应对外层的数据治理和使用上,这才是它真正值钱的地方。
