Flink CDC 实现达梦数据库到 MySQL 实时同步实战指南 📅 发布时间:2026/9/7 5:56:48 👁 浏览次数: 简介针对Flink CDC连接达梦数据库的实时同步需求一整套可运行的工程实现被整理为zip压缩包。方案面向有一定Flink基础、需要基于日志捕获达梦数据库增删改操作的开发者覆盖Java作业与SQL两种开发方式可支撑实时数仓、数据监控、事件驱动应用等场景。包内共315个文件以jar库、java/class源码及编译产物、xml/properties配置、sql脚本为主压缩包约341.71MB其中263个jar已内置Flink连接器与达梦驱动等依赖10个xml与3个properties提供作业和连接配置模板2个sql脚本可直接用于数据操作验证docx说明对工程结构与参数做了补充。与普通教程类资源不同这份工程包更强调开箱即用开发者拿到后可结合自身数据库地址、端口、用户名密码等参数修改配置快速启动同步作业并验证结果。目前已有759人学习下载对正在进行达梦数据库实时同步选型、二次开发或排障的工程师来说是节省重复搭建成本、加速落地的实用参考。 前阵子接了个信创项目要求把生产环境的达梦数据库核心表实时同步到下游 MySQL 分析库。我第一反应是翻 Canal、DataX 这些老熟人结果发现 Canal 官方根本没适配达梦DataX 只能做离线批量搬运做不到真正的日志级实时同步。直到我仔细翻了 Flink CDC 3.x 的官方连接器列表看到dameng-cdc躺在里面悬着的心才放下来。这应该是目前开源生态里为数不多能直接基于达梦日志做实时同步的正路。如果你也在为达梦数据库怎么接入实时数据链路发愁这篇文章应该能帮你省掉不少试错的成本。1. 为什么偏偏盯上“基于日志”的同步方案先说说我为什么对“基于日志”这件事这么执着。拿达梦当源库做同步方案其实不少但大多数是应付式方案真正扛得住生产环境的不多。1.1 传统手段为什么顶不住早年大家做数据同步无非是那几板斧应用双写业务代码里同时写达梦和下游。逻辑耦合重开发要改代码漏写、回滚、事务不一致都能把人逼疯。定时任务 增量字段同步任务每隔几分钟扫一次 update_time把新数据拉出来。这个方案的问题很致命源表没有时间字段就废了UPDATE 之后 update_time 没更新又废了物理 DELETE 根本抓不到而且倒腾一次增量要写不少脚本逻辑。数据库触发器在达梦上建触发器把变更写入中间表再靠调度任务搬运。触发器会拖累业务库性能生产环境 DBA 一般不会同意而且 DDL 变更、大批量 UPDATE 容易把触发器的性能问题放大成一个事故。JDBC 轮询模拟 CDCFlink CDC 也支持 JDBC connector 做周期拉取但轮询本质上还是“定时扫描”延迟取决于轮询间隔对表结构、数据量也很敏感大批量 UPDATE 每次都会把整表旧数据拉出来压力大还不实用。这些方案严格来说都不是“实时同步”更多是“准实时离线搬运”。真正的实时同步关注的是事务日志源库每产生一条 INSERT、UPDATE、DELETE目标库在秒级内原样收到。这个体验只有基于日志的方案能给到。1.2 “基于日志”和“JDBC 拉取”的本质区别达梦数据库和 Oracle 的体系很像运行时会持续产生 redo log在线日志。为了支持恢复和容灾还会把历史日志归档成 archive log归档日志。数据库的所有变更无论你是手动执行 SQL、存储过程批量更新还是走应用连接池最终都会落进日志里。基于日志的实时同步本质上就是去“读日志、解析日志、翻译成变更事件”。它有两个别的手段给不了的好处对业务库零侵入同步进程不碰表数据不加触发器不改业务代码只是以普通客户端的身份去读日志。能捕获一切 DML哪怕是UPDATE 全表、DELETE 全表这种不走应用逻辑的操作日志里都有记录同步任务也能完整翻译过去。用一句话总结JDBC 轮询是“问数据库要数据”基于日志是“等数据库把变更通知到你”。实时性、完整性、对源库的影响完全不是一个量级。2. Flink CDC 是如何消费达梦日志的想用好一个工具不能只看表面配置还得理解它底层到底做了什么。Flink CDC 的达梦连接器执行流程大概可以拆成两个阶段。2.1 快照和增量如何无缝衔接整个过程分两步走全量快照阶段任务刚启动时连接器会先对配置好的表做一次一致性快照把当前已有的全部数据读出来。这一步解决的是“历史数据”的问题。增量日志阶段快照启动的同时连接器会记录当前数据库日志的位点快照做完之后从那个位点继续往后读取归档日志把新增的变更一条条解析出来往下游发。这个设计和 Flink CDC 处理 MySQL、Oracle 的套路一脉相承快照和增量通过“日志位点”衔接起来。因此源库的归档日志必须保留足够长的时间至少要覆盖同步任务从启动快照到切到增量阶段的整个窗口否则任务会报“找不到日志位点”之类的错误这也是我在后面会重点强调的坑。2.2 达梦端要理解的几个前提用 Flink CDC 接达梦还有几个概念先搞清楚不然配置起来会一头雾水达梦的“用户”即“模式”达梦和 Oracle 类似一个用户对应一个模式schema表都挂在用户下面。这跟 MySQL 里“database”和“table”两层结构不一样配置同步任务时要明确指定 schema不然连接器找不到表。归档日志不等于在线日志Flink CDC 读的是归档日志。所以源库必须开启归档模式否则连接器拿不到足够的日志数据。网上很多帖子只讲了怎么配 Flink 端把达梦端要开归档这件事给漏了结果同步任务全量阶段好好的一切到增量就卡住。版本差异目前社区适配较好的是 DM8。DM7 虽然也能跑但动态视图、权限体系存在差异建议先在测试环境把连通性和日志读取验证一遍再上生产。如果你的需求不只是“把数据同步到 MySQL”而是想统一接管多源数据Flink CDC 3.x 的 pipeline 模式还能把达梦同步到 Kafka、Doris、StarRocks 等下游底层的变更事件格式是统一的下游加工起来很顺手。3. 同步前的达梦端改造归档日志与权限准备这一步是重中之重。很多 Flink CDC 接达梦的失败案例根源都在达梦端没准备好。别急着写 Flink SQL先把源库这边收拾利索。3.1 开启归档日志先检查达梦是否已经开启归档模式。用 DM 管理工具连上实例执行SELECT * FROM V$DATABASE; -- 关注 ARCHIVE_MODE 字段TRUE 表示已开启归档如果没开最简单的方式是使用 DM 管理工具的图形化界面右键实例 - 管理服务器 - 归档配置 - 勾选归档模式然后重启实例生效。命令行方式也支持DM8 上常见的操作序列类似这样-- 需要拥有 ALTER DATABASE 权限 ALTER DATABASE MOUNT; ALTER DATABASE ARCHIVELOG; ALTER DATABASE OPEN;注意开启归档一般要求重启数据库实例建议在业务低峰期操作并且提前确认归档目录的磁盘空间。归档日志增长速度取决于业务变更量磁盘规划要留足余量不然归档写满把库搞挂的事故可真不少见。3.2 创建同步账号与授权我的做法是单独建一个专用账号给 FlinkCDC 用不要拿 SYSDBA 到处跑安全审计过不了。DM8 下可以参考下面这套授权CREATE USER FLINKCDC IDENTIFIED BY Flink_123456; GRANT SELECT_CATALOG_ROLE TO FLINKCDC; GRANT SELECT ANY TABLE TO FLINKCDC; GRANT SELECT ON SYS.V$ARCHIVED_LOG TO FLINKCDC;SELECT_CATALOG_ROLE主要用来读系统视图SELECT ANY TABLE是为了让同步账号能读取业务库各张表。如果你只想同步个别几张表可以收窄权限只给具体表的 SELECT 权限但视图部分基本绕不开。权限给少了最常见的现象是全量快照能跑增量阶段一拉日志就报权限不足。排查这类问题可以先用达梦自带的工具或者 DBeaver 拿同一个账号登录手动执行一遍SELECT * FROM V$ARCHIVED_LOG能查到数据Flink 端权限这块基本问题不大。3.3 网络连通性与端口达梦默认端口是 5236。同步任务部署的机器要能直连源库防火墙、安全组都要放行。我之前还遇到过一次诡异问题本机用 DBeaver 连达梦很正常Flink 任务却连不上后来发现是服务器上的安全策略只放行了开发机的 IPFlink 所在的节点不在白名单里。这类网络问题比较隐蔽任务日志里报的又是连接超时容易让人误判成配置错误。4. 实战达梦到 MySQL 实时同步完整配置达梦端准备就绪后Flink 端的配置就相对轻松了。下面以一张典型的订单表为例演示完整的同步链路。4.1 场景说明与表结构假设达梦库里有张订单表PROD.FLINKCDC.T_ORDER核心字段如下字段名类型说明ORDER_IDBIGINT主键USER_IDBIGINT用户IDAMOUNTDECIMAL(12,2)订单金额STATUSINT状态CREATE_TIMETIMESTAMP(3)创建时间UPDATE_TIMETIMESTAMP(3)更新时间目标是把这张表近实时同步到 MySQL 的ods.t_order并且保持 DML 语义完整。4.2 Flink SQL 方式如果用的是 Flink CDC 2.x 或 3.x 的 SQL 模式先在 Flink SQL 客户端建一张映射达梦表的 Source 表CREATE TABLE dameng_t_order ( order_id BIGINT NOT NULL, user_id BIGINT, amount DECIMAL(12, 2), status INT, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector dameng-cdc, hostname 192.168.1.10, port 5236, username FLINKCDC, password Flink_123456, database-name PROD, schema-name FLINKCDC, table-name T_ORDER, scan.startup.mode initial, server-time-zone Asia/Shanghai );再建 MySQL 的 Sink 表然后通过INSERT INTO把数据写过去CREATE TABLE mysql_t_order ( order_id BIGINT NOT NULL, user_id BIGINT, amount DECIMAL(12, 2), status INT, create_time TIMESTAMP(3), update_time TIMESTAMP(3), PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( connector jdbc, url jdbc:mysql://192.168.1.20:3306/ods, username root, password Root_123456, table-name t_order ); INSERT INTO mysql_t_order SELECT order_id, user_id, amount, status, create_time, update_time FROM dameng_t_order;这里有几个容易踩的点Source 表必须声明主键。没有主键的话UPDATE 和 DELETE 事件在部分场景下安全语义会退化丢失变更类型信息。字段类型要一一对应。达梦的NUMBER、MySQL 的DECIMAL之间尽量做显式映射避免精度意外丢失。JDBC Sink 默认是 UPSERT 语义依赖主键做幂等写入。如果目标表没有主键数据重复问题会非常难受。4.3 更推荐的 Pipeline 方式如果你用的是 Flink CDC 3.x我更推荐用 YAML Pipeline 文件做配置。它不需要写一大堆 SQL而是以数据同步作业为单位配置直观版本升级后的兼容性也好一些source: type: dameng hostname: 192.168.1.10 port: 5236 username: FLINKCDC password: Flink_123456 tables: PROD.FLINKCDC.T_ORDER server-time-zone: Asia/Shanghai sink: type: jdbc url: jdbc:mysql://192.168.1.20:3306/ods username: root password: Root_123456 tables: ods.t_order pipeline: name: dameng-to-mysql-order parallelism: 2然后通过命令提交flink-cdc.sh \ --config pipeline.yaml这套 Pipeline 模式封装了下游的序列化和 connector 细节适合已经进入生产维护期的团队。而且它天然支持表级映射后续如果要把表同步到 Kafka 或 Doris只需要替换 sink 块源表配置完全不用动。4.4 同步效果验证配置跑起来之后怎么确认真的是“实时同步”我习惯做三轮验证全量验证先在达梦源表插几条初始数据再启动任务确认 MySQL 目标表出现同样数据。增量验证在达梦源表执行 INSERT、UPDATE、DELETE分别在几秒内观察目标表变化确认新增、修改、删除都能正确映射。断点验证停掉 Flink 任务再往达梦插入一批数据然后重启任务看数据能否续上。这一步同时能验证 checkpoint 是否生效。这里的第三点尤其重要。Flink CDC 本质是“至少一次”语义重启恢复时可能出现重复数据所以下游必须幂等。JDBC Sink 靠主键覆盖倒是问题不大但如果下游是 Kafka订阅方就得自己设计好去重逻辑。我建议在提交任务前打开 checkpoint否则重启丢数据就是家常便饭SET execution.checkpointing.interval 3s; SET execution.checkpointing.mode EXACTLY_ONCE;5. 实测中踩过的坑与调优建议任何一个新技术落地坑都是踩出来的。我把这段时间接达梦遇到的典型问题整理一下希望能给你提个醒。5.1 归档日志被清导致任务中断这是最要命的坑。某次测试环境同步任务跑了两天突然增量阶段报错看日志定位到“无法读取指定的日志文件”。查了一圈才发现达梦实例的归档日志被 DBA 手动清理过Flink CDC 记的位点已经指向了一个被删掉的归档文件。所以做基于日志的同步归档日志的保留策略必须和同步任务的生命周期对齐。如果业务变更量大建议给同步任务配置告警一旦出现日志滞后或者读取失败第一时间介入而不是等数据对账时才发现两边已经差了一大截。注意归档日志磁盘爆满和归档被清理是两个方向的问题前者会导致达梦写入阻塞甚至实例异常后者会导致同步任务断链运维上都要盯紧。5.2 大小写和 schema 匹配问题达梦默认对标识符大小写敏感未加双引号的表名会被统一转成大写。如果建表时用了双引号创建小写表配置里写大写名称就会“找不到表”。反过来说配置里写小写也可能匹配不上因为达梦内部实际存的是大写。我的建议是建表、建用户时统一用大写Flink 配置里也填大写。这样最省心。如果历史遗留已经建了小写表配置时就要精确填写匹配的大小写不能想当然。5.3 权限不足与视图兼容达梦连接器在读取日志时依赖系统的动态视图。不同版本达梦视图名可能不一样环境变量和参数设置在连接器文档里也可能有细微差异。权限不足的表现很多样有的是全量完后增量迟迟不输出有的是直接抛异常退出。遇到这类问题先用达梦账号手动执行相关视图查询确认有没有权限再检查视图名是否和当前版本一致。给自己建一个“最小验证脚本”每次部署新环境先在达梦端验证一遍比盲调 Flink 配置高效得多。5.4 大表快照与并行度全量快照阶段Flink CDC 会把大表按主键分段并行读取。但如果表的主键分布极不均匀比如主键是自增列前 1 亿行和后面 100 行占比悬殊默认的 chunk 划分可能导致某些子任务特别慢。遇到这种表我习惯先压一张数据量最大的表观察 Flink Web UI 上每个 subtask 的繁忙程度如果某个子任务明显是长尾热点就手动调整scan.incremental.snapshot.chunk.size来控制分片大小并适当提高并行度。快照阶段和增量阶段对资源的需求不一样不能指望一套参数通吃所有场景。5.5 关于 DDL 和下游幂等Flink CDC 对达梦的 DDL 同步支持还不够完善不像 MySQL 那样能比较从容地处理在线表结构变更。所以我通常建议上线期间尽量冻结业务表结构变更如果必须变更先停同步任务改完表结构后重启任务让任务重新做一次快照然后再恢复增量。这个流程笨但稳。至于下游幂等我之前因为偷懒目标表没设主键压测时一重启任务同一批数据插了两遍直接污染了报表数据。后来所有 JDBC Sink 的目标表都强制要求主键或者直接用主键唯一约束去重。Kafka 场景则是在消费端做 key 去重或状态去重宁可多设计一层幂等也不要把数据完整性的赌注押在“应该不会重发”上。回到开头那个项目基于 Flink CDC 的达梦实时同步最终平稳跑上了生产。现在每天凌晨看监控延迟稳定在秒级比之前定时脚本的方案不知道省心多少倍。如果你也打算把达梦接入实时链路建议先拿一张小表把环境彻底调通再逐步扩大同步范围。毕竟数据同步这事链路越早摸透后面才越不容易被突发流量和结构变更打得措手不及。本文还有配套的精品资源点击获取