简介FlinkCDC 与达梦数据库结合的日志实时同步方案资源适合需要将达梦数据变更实时接入 Flink 的开发者、数据工程师与数仓建设者。该方案基于日志解析捕获插入、更新、删除操作既可用于数据仓库实时同步也适合构建实时报表、监控告警等事件驱动应用对构建低延迟、高可靠数据管道有直接参考价值。压缩包共 5 个文件大小约 35.48MB构成以 2 个 jar 包为主体包含连接器与 JDBC 驱动另附 1 个 SQL 初始化脚本、1 个参考程序压缩包及 1 份 docx 用户手册覆盖环境准备、启动配置与 Java/SQL 两种接入方式。已有 2083 人学习下载。使用者可依据手册快速搭建 Flink CDC 连接达梦的同步作业参照程序理解如何捕获变更数据并结合实际业务做处理与分析减少自行摸索底层日志解析的耗时。1. FlinkCDC 同步达梦数据库为什么基于日志解析是唯一靠谱的路做过实时同步的工程师应该都见过这种尴尬业务库换成达梦数据库之后原本在 MySQL 上跑得好好的 FlinkCDC 实时同步链路到达梦这儿突然哑火了。直接 JDBC 轮询的话慢查询日志里全是罪证延迟还只能做到分钟级大事务一多还会把源库拖垮真正要做到秒级实时同步唯一靠谱的路就是解析达梦数据库的归档日志。基于日志解析既不侵扰源库性能又能拿到完整的事务和行变更这正是 FlinkCDC 实时同步方案能对接达梦的原因。这篇就把我从零搭起来的达梦实时同步链路完整拆一遍——归档配置、Flink 代码、参数设置、上线后踩过的坑希望能帮到正在接达梦的你。2. 达梦数据库开启归档与补充日志同步前的三个前提条件2.1 为什么必须开启归档模式达梦数据库的重做日志是循环复用的默认情况下只保留最近的变更。基于日志解析的同步组件启动后要从一个稳定的位点开始持续读取如果重做日志被覆盖前面的变更就全丢了只能从当前时刻开始接入历史数据拿不到。所以第一步就是把达梦切到归档模式。归档模式会在重做日志切换时把日志文件复制到独立目录同步组件读取这个目录里的文件做解析。这里有个容易被忽略的点达梦开启归档需要数据库先处于 mount 状态不是所有版本都能在 open 状态直接改。我第一次操作时直接在 open 状态下执行 ALTER DATABASE ADD ARCHIVELOG结果报“数据库状态不允许该操作”后来才意识到要按 mount、改配置、open 的顺序走。另一个容易被忽略的点是归档目录要提前建好并且确保数据库进程有写权限否则配置完成后归档文件写不进去同步任务启动后一直等不到新日志。2.2 用 disql 完成归档和补充日志配置达梦数据库常用命令里最核心的就是这一组 disql 操作。先登录实例确认当前状态再切换模式。# 使用 disql 登录默认账密 SYSDBA/SYSDBA端口 5236 disql SYSDBA/SYSDBAlocalhost:5236 # 确认当前是否已经是归档模式1 表示归档0 表示非归档 SELECT ARCH_MODE FROM V$DATABASE; # 切到 mount 状态这一步会短暂断开业务连接 ALTER DATABASE MOUNT; # 添加归档日志配置dest 指向归档目录file_size 是单个文件上限 ALTER DATABASE ADD ARCHIVELOG DEST/dm8/arch,TYPElocal,FILE_SIZE64,SPACE_LIMIT0; # 重新打开数据库 ALTER DATABASE OPEN;ARCH_MODE 查询结果如果是 0就按上面的流程走一遍。TYPElocal 表示写本地磁盘dest 目录要确保已存在且有写权限。SPACE_LIMIT0 表示不限制归档总空间生产上我一般会设一个具体上限而不是 0单位是 MB比如 SPACE_LIMIT51200 就是 50G防止归档无限增长把磁盘写满。FILE_SIZE64 指单个归档文件 64MB文件写满后自动切换新文件。归档配置完成之后还要开启补充日志。默认情况下达梦日志里只记录变更后的数据拿不到变更前的镜像下游做 upsert 或者数据对账时就缺了关键信息。# 开启全列补充日志记录变更前镜像 ALTER DATABASE ADD SUPPLEMENTAL LOG DATA (ALL) COLUMNS; # 确认补充日志状态返回值应是 YES SELECT SUPPLEMENTAL_LOG_DATA_ALL FROM V$DATABASE;开启全列补充日志之后每次 update 都会把完整行镜像写进日志日志量会明显变大。如果业务表字段特别多、更新又频繁建议只对需要同步的关键表开启补充日志而不是全库全列否则对磁盘空间和 IO 的压力会翻好几倍。我后来的习惯是核心业务表开全列普通配置表只依赖默认日志这样既保证同步数据完整又控制住日志规模。2.3 用 Navicat 连接达梦数据库验证配置结果配置完成后最直接的验证方式就是用 Navicat 连接达梦数据库跑一遍查询确认状态。Navicat 连接达梦数据库时要注意选对驱动和端口默认端口是 5236不是 MySQL 的 3306也不是 Oracle 的 1521。连接参数如下表连接项值主机达梦实例所在 IP端口5236用户名SYSDBA密码安装时设置的密码模式需要同步的业务模式例如 DMUSER连接成功后重新执行 SELECT ARCH_MODE FROM V$DATABASE确认返回值是 1。再往业务表插一条数据观察归档目录下是否生成新的 .log 文件。这一步是为了确认归档链路是真的通了不是配置完就完事。如果 Navicat 报“无效连接”先检查达梦的监听服务是否启动再看端口是否被防火墙挡住。2.4 归档对源库性能和权限的影响开启归档之后每个重做日志切换都会伴随一次文件写入对磁盘 IO 有一定占用但整体影响可控。真正影响性能的是全列补充日志开启之后每次 update 都会写完整行镜像如果一个表有几十个字段、更新又频繁日志量会翻好几倍。我的做法是只对需要同步的表开启补充日志并且定期检查归档目录的清理情况。同步任务消费完的归档文件如果没有自动清理机制需要在确认位点已经消费过去之后再手工删最安全的做法是保留最近 24 小时的归档之前的归档由脚本按时间清理。3. 借道 Debezium 引擎构建 FlinkCDC 任务核心代码与参数拆解3.1 为什么不是直接写一个 Flink CDC 连接器达梦数据源可以接入 Flink 实时链路的方案并不多。自己从零写一个达梦归档日志解析器成本太高归档日志的格式不是公开的版本升级后格式变不变也不确定风险全压在自己身上。常见做法是借道 Debezium 引擎Debezium 负责对接达梦日志解析把变更事件封装成标准格式Flink 侧用 SourceFunction 接住再交给下游 Kafka 或者数仓。这条链路的好处在于解析逻辑完全由 Debezium 承担Flink 只负责 source 的并发、状态和 offset 管理两边各管各的出了问题也好定位。Flink CDC 本身也可以理解成 Debezium 在 Flink 生态里的封装所以用 Flink 算子接 Debezium拿到的还是标准的 Flink 实时链路后续接 Flink SQL 或者 DataStream API 都不受影响。项目里我一般会先把 Debezium 独立跑通确认日志解析没有问题再把它包进 Flink 工程这样排查问题的时候不用同时怀疑两边。3.2 Debezium 达梦连接器核心参数连接器的关键配置项如下表其中 offset 存储方式直接决定了任务重启后能不能续上位点我单独拎出来讲。配置项推荐值说明connector.classio.debezium.connector.dm.DmConnector达梦连接器入口需把对应 jar 放进 Flink liboffset.storageorg.apache.kafka.connect.storage.FileOffsetBackingStore本地文件存 offset便于排查offset.storage.file.filename/data/flink/dm/dm-offsets.datoffset 落盘路径重启后读取offset.flush.interval.ms60000每隔 60 秒把位点持久化一次snapshot.modeinitial启动时先做全量快照再做增量table.include.listDMUSER.ODS_ORDER,DMUSER.ODS_USER只同步指定的表多张表用逗号分隔max.batch.size2048单批最大变更条数防止 OOMdecimal.handling.modedouble小数类型转成 double避免精度字符串问题offset 的持久化间隔很关键。如果任务正常运行时宕机内存里记录的位点还没有落盘重启后连接器会从最后一次持久化的位点重新解析这一段日志会被重复消费一遍下游要做幂等处理。offset.flush.interval.ms 调得越小丢失的位点越少但磁盘写入也越频繁我一般设置在 30 到 60 秒之间。3.3 Flink Source 函数完整代码下面这个 DmLogSource 是同步任务的核心。它把 Debezium 引擎包在 Flink 的 RichSourceFunction 里启动时读 offset运行中出事件取消时关闭引擎。import io.debezium.engine.ChangeEvent; import io.debezium.engine.DebeziumEngine; import io.debezium.engine.format.Json; import org.apache.flink.configuration.Configuration; import org.apache.flink.streaming.api.functions.source.RichSourceFunction; import java.util.Properties; public class DmLogSource extends RichSourceFunctionString { private volatile boolean running true; private DebeziumEngineChangeEventString, String engine; private final Properties props; public DmLogSource(Properties userProps) { this.props userProps; } Override public void open(Configuration parameters) throws Exception { // 固定连接器入口和 offset 存储位置避免每次重启都重读全量 props.setProperty(connector.class, io.debezium.connector.dm.DmConnector); props.setProperty(offset.storage, org.apache.kafka.connect.storage.FileOffsetBackingStore); props.setProperty(offset.storage.file.filename, /data/flink/dm/dm-offsets.dat); props.setProperty(snapshot.mode, initial); if (props.getProperty(table.include.list) null) { props.setProperty(table.include.list, DMUSER.ODS_ORDER); } engine DebeziumEngine.create(Json.class) .using(props) .notifying(record - { // record.value() 是 JSON 格式的变更事件可直接发到下游 Kafka }) .build(); } Override public void run(SourceContextString ctx) throws Exception { engine.run(); } Override public void cancel() { running false; if (engine ! null) { try { engine.close(); } catch (Exception ignored) { // 关闭引擎失败不影响 Flink 任务退出 } } } }逻辑说明open 方法里做的两件事一是补全连接器参数二是构建 Debezium 引擎。run 方法启动引擎后进入阻塞循环引擎每解析出一条变更事件notifying 回调就会触发一次。cancel 方法用于 Flink 任务取消时释放引擎资源避免连接句柄泄漏。参数说明snapshot.modeinitial 表示任务第一次启动时先对 table.include.list 指定的表做全量快照快照完成后才进入增量解析。如果表里数据量很大全量阶段会持续一段时间要提前评估对源库的查询压力。table.include.list 建议只列真正需要的表列太多全量阶段会拉得很长而且快照期间表结构变更会中断快照所以上线前业务表结构要冻结。3.4 在 Flink 主程序里装配并跑通StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.setParallelism(1); env.enableCheckpointing(60000); Properties dmProps new Properties(); dmProps.setProperty(database.hostname, 10.20.30.40); dmProps.setProperty(database.port, 5236); dmProps.setProperty(database.user, SYSDBA); dmProps.setProperty(database.password, your_password); dmProps.setProperty(database.dbname, DMSERVER); dmProps.setProperty(database.schema, DMUSER); DataStreamSourceString source env.addSource(new DmLogSource(dmProps)); source.map(json - json \n).print(); env.execute(dm-log-sync);逻辑说明主程序里只做三件事一是开启 checkpoint 保证位点状态定期持久化二是把达梦连接信息放进 Properties三是把 source 的 JSON 事件打印出来。第一次跑这个工程建议先把 print 打开看数据内容确认能持续输出变更事件后再接下游。参数说明database.dbname 对应达梦实例名database.schema 是要同步的模式名这两个值不区分大小写但必须和达梦实际配置一致。setParallelism(1) 是因为 Debezium 引擎本身是单线程解析并行度设高没有意义反而会重复消费日志。4. 常见问题与避坑归档清理、快照衔接和连接器配置4.1 归档目录把磁盘写满现象告警平台凌晨弹出磁盘使用率 98%达梦实例无法新增连接应用侧开始报连接超时连 Navicat 连接达梦数据库都进不去。原因达梦服务器上 /dm8/arch 目录写满了。我的归档配置里 SPACE_LIMIT0表示不限制归档总空间而同步任务因为网络抖动停了一整晚归档文件只增不减最终把磁盘撑爆。解决先手工把确认已经消费过的归档文件挪走释放磁盘空间再修改归档配置加一个空间上限。修改命令是 ALTER DATABASE MOUNT 之后执行 ALTER DATABASE MODIFY ARCHIVELOG DEST/dm8/arch,TYPElocal,FILE_SIZE64,SPACE_LIMIT51200再 ALTER DATABASE OPEN。从那以后我所有达梦实例的归档 SPACE_LIMIT 都设为 50G 或 100G绝不再用 0。4.2 全量快照结束后丢了一段增量现象同步任务第一次启动走了全量快照全量完成后对账发现快照期间产生的部分更新没进 Kafka。原因表在快照开始之后、快照结束之前又发生了写入而连接器记录快照起点时归档日志的位点还没跟上快照做完之后跳过中间那段变更。解决把上下线流程改成规范操作停业务写入、清空 offset 文件、重启任务完整走一遍 initial 快照。如果已经上线了不想停业务可以在快照结束后手动补拉一段归档日志但补拉逻辑复杂我一般不推荐。最靠谱的还是控制变更窗口快照期间冻结写入。4.3 大事务把 Flink 内存打满现象某天业务跑了一次几十万行的批量 updatesource 算子内存立刻飙到几个 G随后 GC 频繁任务整体延迟拉高。原因Debezium 会将一个事务内的所有变更一次性产出一个大事务等于一条超大消息流Flink source 拿到后要先在内存里缓冲再逐条下发缓冲越大越容易 OOM。解决在 source 后面加一个限流缓冲算子把变更事件按条数切小之后再发往下游。另一个思路是在业务侧给大事务单独打标同步任务识别到超过阈值的变更时只记录位点不逐条解析等业务低谷期再补。4.4 DDL 变更不进入同步链路现象业务侧给达梦表新增了一个字段下游 Kafka 里的 JSON 结构还是旧字段后续解析作业按新字段取值时全部为空。原因达梦补充日志默认记录的是 DML 变更DDL 操作默认不会通过日志解析产出变更事件Flink 链路自然感知不到表结构变化。解决同步链路上先屏蔽 DDL表结构变更走线下流程通知下游。如果业务无法接受就额外跑一个表结构比对工具定时检测达梦和下游 schema 的差异发现不一致时触发下游刷新。现阶段这是最稳的方案日志解析 DDL 在各数据库里都是最麻烦的部分。4.5 offset 文件丢失导致重复全量同步现象Flink 任务重启后source 没有从上次位点续上而是重新跑了一轮全量快照下游 Kafka 里出现大量重复数据。原因我用 FileOffsetBackingStore 把 offset 存在本地文件部署机器的磁盘被清理过或者任务跑在容器里容器重建后本地文件没了连接器找不到位点只能重新快照。解决生产环境把 offset 存到 Kafka 里的 KafkaOffsetBackingStoretopic 销毁重建之间位点不会丢。如果坚持用本地文件就把 offset 文件挂到持久化磁盘并且定期备份。5. 同步任务验证与参数调优从对账到延迟监控5.1 用全量 count 做上线前对账同步任务上线前先做一轮行数对账。从达梦侧统计目标表数量再从 Kafka 消费侧统计同一张表的数量两边数值一致才能说明基础链路是通的。# 从达梦侧统计目标表数量 SELECT COUNT(*) FROM DMUSER.ODS_ORDER; # 从 Kafka 侧消费同 topic 统计消息行数 kafka-console-consumer.sh --bootstrap-server localhost:9092 \ --topic dm_ods_order --from-beginning \ --property print.valuetrue | wc -l行数一致只能说明数量对得上内容是否完整还要抽主键做差集比对。常见做法是从达梦导出主键集合从 Kafka 里解析出主键集合两边做差差集为空才放行上线。这个比对脚本不用写得多复杂关键是要在业务低峰期跑否则写入中的新数据会干扰比对结果。5.2 延迟监控怎么量达梦日志同步和 MySQL binlog 同步不一样归档文件的切换是异步的所以延迟会有一个基础间隔不是实时到秒级。我习惯让业务应用在写入达梦表的同一时刻往一个小 Kafka topic 里发一条心跳记录包含当前数据库时间。同步任务正常消费到对应表的变更之后比对心跳时间和变更到达时间差值就是链路延迟。如果没有心跳机制可以用慢查询日志辅助判断业务写入高峰时段再手动对比归档文件的切换时间。慢查询日志在生产上要谨慎开启只开短时间窗口用完就关不然对达梦性能有影响。5.3 核心参数调优对照表参数调整场景建议值max.batch.size大事务频繁调大到 4096减少 flush 次数poll.interval.ms日志切换频繁调小到 500降低空转等待offset.flush.interval.ms怕重启丢位点调到 5000每隔 5 秒落盘一次heartbeat.interval.ms长事务场景开启避免空闲连接被误判断开参数说明max.batch.size 不是越大越好调太大会增加单批处理时间下游消费跟不上就产生背压。poll.interval.ms 调小会让连接器更频繁地检查归档目录日志量小的场景下会造成无谓的 IO 开销。这些参数要在启动前改改完重启任务切记确认 offset 文件路径没有同时被改动否则恢复时容易翻车。6. 把同步链路做成可运维的数据管道offset 交给 Kafka 与故障恢复技巧6.1 用 KafkaOffsetBackingStore 替代本地文件本地文件存 offset 在单机场景够用但 Flink 任务重启时如果本地文件被误删整个任务会重新走一遍全量快照。生产环境我习惯把 offset 存到 Kafka 的 topic连接器自己提交位点不依赖单机磁盘。props.setProperty(offset.storage, org.apache.kafka.connect.storage.KafkaOffsetBackingStore); props.setProperty(offset.storage.kafka.topic, dm-sync-offsets); props.setProperty(offset.storage.kafka.bootstrap.servers, localhost:9092);逻辑说明offset 存到 Kafka 之后只要 topic 不被误删任务不管怎么重启都能从已提交的位点接续解析。topic 建议用单分区保留时间设置为保留全部因为这个 topic 的数据量本身很小但丢一个位点段就可能导致重复全量快照。6.2 用 Flink Checkpoint 兜底恢复Flink 侧开启 checkpoint 之后source 的 offset 会进入算子状态下游落库失败时可以从最近一次 checkpoint 恢复。env.enableCheckpointing(60000);说明checkpoint 周期设 60 秒恢复时最多重复 60 秒的数据下游要做幂等。如果下游是 Kafka建议把写入 Kafka 的 flush 间隔也调小避免 checkpoint 恢复时重新发送大量消息。6.3 故障恢复的固定巡检流程从那次归档磁盘被打爆之后我每次部署达梦同步任务都强制走一遍固定流程先查 V$DATABASE 的 ARCH_MODE再确认补充日志是 ALL接着检查归档目录空间和 SPACE_LIMIT 是否有限制然后对一遍 table.include.list 和实际业务表数量最后启动任务盯着延迟监控看完 10 分钟再离开。这套流程看着土但救过我好几次一次是归档没开一次是 offset 路径写错都是在巡检阶段发现的。希望帮到你。本文还有配套的精品资源点击获取