简介: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/SYSDBA@localhost:5236 # 确认当前是否已经是归档模式,1 表示归档,0 表示非归档 SELECT ARCH_MODE FROM V$DATABASE; # 切到 mount 状态,这一步会短暂断开业务连接 ALTER DATABASE MOUNT; # 添加归档日志配置,dest 指向归档目录,file_size 是单个文件上限 ALTER DATABASE ADD ARCHIVELOG 'DEST=/dm8/arch,TYPE=local,FILE_SIZE=64,SPACE_LIMIT=0'; # 重新打开数据库 ALTER DATABASE OPEN;ARCH_MODE 查询结果如果是 0,就按上面的流程走一遍。TYPE=local 表示写本地磁盘,dest 目录要确保已存在且有写权限。SPACE_LIMIT=0 表示不限制归档总空间,生产上我一般会设一个具体上限而不是 0,单位是 MB,比如 SPACE_LIMIT=51200 就是 50G,防止归档无限增长把磁盘写满。FILE_SIZE=64 指单个归档文件 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.class | io.debezium.connector.dm.DmConnector | 达梦连接器入口,需把对应 jar 放进 Flink lib |
| offset.storage | org.apache.kafka.connect.storage.FileOffsetBackingStore | 本地文件存 offset,便于排查 |
| offset.storage.file.filename | /data/flink/dm/dm-offsets.dat | offset 落盘路径,重启后读取 |
| offset.flush.interval.ms | 60000 | 每隔 60 秒把位点持久化一次 |
| snapshot.mode | initial | 启动时先做全量快照,再做增量 |
| table.include.list | DMUSER.ODS_ORDER,DMUSER.ODS_USER | 只同步指定的表,多张表用逗号分隔 |
| max.batch.size | 2048 | 单批最大变更条数,防止 OOM |
| decimal.handling.mode | double | 小数类型转成 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 RichSourceFunction<String> { private volatile boolean running = true; private DebeziumEngine<ChangeEvent<String, 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(SourceContext<String> 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.mode=initial 表示任务第一次启动时先对 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"); DataStreamSource<String> 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_LIMIT=0,表示不限制归档总空间,而同步任务因为网络抖动停了一整晚,归档文件只增不减,最终把磁盘撑爆。
解决:先手工把确认已经消费过的归档文件挪走,释放磁盘空间,再修改归档配置加一个空间上限。修改命令是 ALTER DATABASE MOUNT 之后执行 ALTER DATABASE MODIFY ARCHIVELOG 'DEST=/dm8/arch,TYPE=local,FILE_SIZE=64,SPACE_LIMIT=51200',再 ALTER DATABASE OPEN。从那以后我所有达梦实例的归档 SPACE_LIMIT 都设为 50G 或 100G,绝不再用 0。
4.2 全量快照结束后丢了一段增量
现象:同步任务第一次启动走了全量快照,全量完成后对账,发现快照期间产生的部分更新没进 Kafka。
原因:表在快照开始之后、快照结束之前又发生了写入,而连接器记录快照起点时,归档日志的位点还没跟上,快照做完之后跳过中间那段变更。
解决:把上下线流程改成规范操作:停业务写入、清空 offset 文件、重启任务完整走一遍 initial 快照。如果已经上线了不想停业务,可以在快照结束后手动补拉一段归档日志,但补拉逻辑复杂,我一般不推荐。最靠谱的还是控制变更窗口,快照期间冻结写入。
4.3 大事务把 Flink 内存打满
现象:某天业务跑了一次几十万行的批量 update,source 算子内存立刻飙到几个 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 里的 KafkaOffsetBackingStore,topic 销毁重建之间位点不会丢。如果坚持用本地文件,就把 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.value=true | 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 路径写错,都是在巡检阶段发现的。希望帮到你。
本文还有配套的精品资源,点击获取