
简介Flink CDC对接达梦数据库的成套实现面向数据工程师与实时计算开发者解决基于日志解析的达梦数据库变更数据捕获与实时同步问题。相比查询式同步日志级捕获对源库影响更小可持续捕捉插入、更新、删除等增量操作支撑实时数仓、数据监控、告警等事件驱动场景。压缩包共5个文件、约35.48MB含Flink CDC连接器jar包、达梦JDBC驱动、参考示例程序、SQL初始化脚本以及dm8用户手册。jar与驱动可直接集成进Flink工程示例程序演示作业编写与调试的典型流程SQL脚本便于在sql-client中快速建立同步任务手册则围绕配置参数、连接方式、兼容性及常见排错展开覆盖从环境搭建到联调上线的关键环节。目前已有2083人学习下载适合需要构建达梦实时同步管道、异构数据入仓或事件驱动应用的团队也适合希望深入理解Flink CDC连接器原理的进阶开发者。1. 达梦实时同步为什么绕不开日志把 FlinkCDC 放在正确的位置业务库跑在达梦上下游数仓、大屏、消息系统都想拿到秒级变更定时任务撑不住轮询查询也很快被频繁 SQL 压垮。FlinkCDC 搭配达梦自身日志能力可以绕过业务表查询直接从日志里还原每一条增删改实现真正的实时同步。这个方案对业务库无侵入不建触发器、不加轮询索引尤其适合表量大、写入频繁、需要把数据流进 Kafka 或数仓的团队。它能让你在达梦数据库使用教程里常见的定时同步玩法直接换成一条可持续扩的流式链路。2. 达梦日志解析原理为什么只有日志能保证完整增量2.1 达梦的 redo 与归档日志解析增量的事实来源达梦数据库的日志机制和 Oracle 很接近日常写入会产生 redo log文件循环复用写满会覆盖旧内容。想要做可追溯的增量同步必须依赖归档日志。归档模式开启后达梦会把写满的 redo 按时间顺序拷贝成归档文件这些文件才是 FlinkCDC 真正要读取的事实来源。生产上我一般用 disql 确认归档状态# 登录达梦实例当前用户需有 dba 权限 disql SYSDBA/xxxxxlocalhost:5236 # 查看归档配置 SELECT NAME, STATUS$ FROM V$ARCHIVE;输出里如果 STATUS$ 为 0说明归档没开。开启归档不能只靠一条 SQL达梦多数版本要求先配置 dmarch.ini再修改 INI 参数最后重启实例。核心配置项是这样ARCH_WAIT_APPLY 0 ARCH_NUM 6 ARCH_SPACE_LIMIT 10240ARCH_NUM 是保留归档文件数量ARCH_SPACE_LIMIT 是归档目录空间上限。两个参数共同决定你最多能回溯多少日志这直接关系到 FlinkCDC 断点续传的可用窗口。生产环境至少保留 24 小时以上的归档否则同步作业稍一抖动就会丢日志。日志有了怎么解析成可读 SQL达梦提供了一套与 Oracle 高度相似的 LogMiner 能力通过 DBMS_LOGMNR 包注册日志文件启动解析后从 V$LOGMNR_CONTENTS 读取 SCN、表名、操作类型、SQL_REDO、SQL_UNDO。做过 Oracle 日志同步的工程师迁移到达梦第一反应都是找这套接口因为它成熟、可控而且不需要侵入业务库。2.2 FlinkCDC 在日志链路里负责什么很多人以为 FlinkCDC 是现成的连接器拿来就连达梦。实际生产链路里FlinkCDC 承担的是框架职责不是黑匣子。它提供 SourceFunction 接口我们在里面实现快照读取和日志解析它负责 checkpoint帮我们把“上次解析到哪个日志、哪个 SCN”这个状态存下来它管理并行度和反压让下游消费不过来的时日志解析能自动放慢。快照阶段通常用 JDBC 直接查询表数据增量阶段用 LogMiner 读 V$LOGMNR_CONTENTS两者在同一个水位线上衔接。这个“先查全量再从日志追增量”的模式是 Flink CDC 的标准玩法。快照期间业务库照常写入这些写入会落到归档日志里快照结束后再从刚快照的时间点开始解析就能保证数据不丢也不重复。有的团队直接用社区提供的达梦连接器但它覆盖的场景有限尤其是 DDL 变化、自定义字段类型、特殊字符集自己写 Source 反而更好控制。我的建议是小规格、少表、日志量不大可以用现成连接器先验证生产链路最好理解 LogMiner 接口后自己实现至少要有能力改解析逻辑。2.3 同步方案对比轮询、触发器和日志解析怎么选拿“达梦数据库使用教程”里最常见的几种同步方案来做对比方案延迟业务侵入实现成本场景定时批处理分钟到小时无低离线数仓轮询增量字段秒到分钟需加时间戳或自增索引中小表、非核心链路触发器 中间表实时高每张表都要触发器和维护高基本不做维护成本失控日志解析 FlinkCDC秒级无高生产主流核心链路首选轮询在数据量小的时候看起来还行但有两处天生缺陷一是 DELETE 无法通过时间戳感知要对比全表才知道删了哪条二是多表更新顺序容易乱下游无法还原事务边界。触发器方案能实时但每张表都要建触发器对业务写入性能影响明显而且达梦上触发器调试和 DDL 维护都很折腾。日志解析虽然前期实现成本高但无侵入、能拿到事务和 SCN 顺序是真正适合长期跑的方案。3. 跑通最小工程从开启归档到消费 Kafka 里的变更事件3.1 达梦侧准备账号权限与归档文件可达性开始写 Flink 代码之前先把达梦侧环境理顺。用 Navicat 连接达梦数据库执行以下 SQL 创建专用账号不建议直接用 SYSDBA 跑同步任务CREATE USER FLINK_USER IDENTIFIED BY flink_123; GRANT CREATE SESSION TO FLINK_USER; GRANT SELECT ON SYS.V$LOGMNR_CONTENTS TO FLINK_USER; GRANT EXECUTE ON SYS.DBMS_LOGMNR TO FLINK_USER;注意达梦的视图名在不同版本略有差异有的环境是 V$LOGMNR_CONTENTS有的是 DBA_LOGMNR_CONTENTS。执行前先查一下当前系统视图名SELECT TABLE_NAME FROM ALL_TABLES WHERE TABLE_NAME LIKE %LOGMNR%;权限只给查询和计算包执行权限不要给业务表 DML 权限。同步账号只需要读日志和归档文件的能力。还有一个容易被忽略的点运行 Flink 的服务器必须能访问达梦数据库服务器上的归档日志目录。如果 Flink 在远程机器要么把归档目录做成 NFS 共享要么通过达梦的日志服务接口暴露否则 LogMiner 只能解析本机文件拉不到历史归档。这里我吃过亏Flink 作业日志一直报“无法定位日志文件”查了半天发现是网络路径没权限。3.2 编写基于 LogMiner 的 Flink SourceFlink 侧我用 DataStream API 实现一个自定义 SourceFunction。核心逻辑是注册日志文件启动 LogMiner循环查询增量 SCN 以上的变更记录然后输出给下游。public class DamengLogMinerSource extends RichSourceFunctionRowData { private String jdbcUrl; private String username; private String password; private long lastScn; Override public void run(SourceContextRowData ctx) throws Exception { try (Connection conn DriverManager.getConnection(jdbcUrl, username, password)) { // 注册第一个日志文件常见做法是最近的一个归档 String addLogSql {call DBMS_LOGMNR.ADD_LOGFILE(?,?)}; try (CallableStatement cs conn.prepareCall(addLogSql)) { cs.setInt(1, 0); cs.setString(2, /dm/arch/ARC000000123.log); cs.execute(); } // 启动 LogMiner String startSql {call DBMS_LOGMNR.START_LOGMNR(?,?,?)}; try (CallableStatement cs conn.prepareCall(startSql)) { cs.setInt(1, 1); // 选项不限 LOB 等按需打开 cs.setInt(2, 1); cs.setLong(3, lastScn); cs.execute(); } // 每次轮询增量 SCN 之后的记录 String sql SELECT SCN, TABLE_NAME, OPERATION, SQL_REDO, SQL_UNDO FROM V$LOGMNR_CONTENTS WHERE SCN ? ORDER BY SCN; PreparedStatement ps conn.prepareStatement(sql); ps.setFetchSize(500); while (running) { ps.setLong(1, lastScn); ResultSet rs ps.executeQuery(); while (rs.next()) { long scn rs.getLong(SCN); String tableName rs.getString(TABLE_NAME); String operation rs.getString(OPERATION); String sqlRedo rs.getString(SQL_REDO); RowData row buildRow(tableName, operation, sqlRedo); ctx.collect(row); lastScn scn; } rs.close(); Thread.sleep(1000); } } catch (Exception e) { throw new RuntimeException(e); } } }代码里的几个参数要特别说明。ps.setFetchSize(500)是 JDBC 每次从达梦取回的行数设小了导致频繁网络往返设大了占用内存500 到 1000 是比较稳妥的区间。Thread.sleep(1000)控制轮询空转频率1 秒一次满足大多数秒级同步诉求日志量大的业务可以改成 500ms但会加重数据库服务器负担。lastScn必须在 checkpoint 里持久化这样作业重启后从上次解析到的 SCN 继续而不是从最早的归档重新扫一遍。RowData 的构建这里不展开但要注意SQL_REDO 拿到的是文本 SQL最好在 Flink 端解析成结构化字段而不是直接整条丢给下游。否则下游消费时还得再解析一次而且类型信息会丢失。常见做法是在 buildRow 里按表名匹配 schema然后逐字段拆分。3.3 把变更事件写到 Kafka形成实时链路Source 写好后配上 Kafka Sink 就能形成一个实时同步管道。DataStreamSourceRowData stream env.addSource( new DamengLogMinerSource(jdbcUrl, username, password, lastScn)); stream.addSink(new FlinkKafkaProducer( dm-dbtable-cdc, new SimpleStringSchema(), // 生产环境换成 JsonSchema 或 AvroSchema kafkaProps )); env.enableCheckpointing(10_000);Kafka 生产者参数里acksall必须开这是保证下游不丢数据的前提。retries调大Flink 端如果因为网络抖动提交失败重试机制能兜底。如果对延迟敏感linger.ms不要超过 50msbatch.size也不用刻意调大。数据流起来后在 Navicat 连接达梦数据库随便改一张被监听表里的数据然后消费 Kafka 对应 topic。能看到 UPDATE 语句对应的 SQL_REDO说明这条链路已经通了。我第一次验证时改了一行数据Kafka 里几分钟没动静最后发现是归档日志还没生成——达梦生成归档的节奏和 redo 切换有关修改后需要等日志切换或手动触发归档。3.4 如何确认同步没有丢数据最直接的方式对一张表做全量 count再统计同步到 Kafka 里的操作数。但更细的验证是选一个有明确业务含义的字段比如“更新时间”在源端插入一条记录写入时间到 Kafka 里看这条数据到达的时间两个时间差就是端到端延迟。另一个朴素但有效的验证把源表按照主键排序导出再把 Kafka 里的变更按主键 replay 成最终视图两者比对。这个逻辑可以写成离线校验作业定时跑。日志同步最怕“看起来在跑实际上丢了一部分”的状态尤其大了之后没人发现。这个校验作业值得在第一天就搭好。4. 同步任务稳定性4 组必调参数与 Nacos 配置管理4.1 快照阶段参数fetchSize 和并行度决定全量速度全量快照如果用 JDBC 直接从达梦读表最容易踩的坑是默认 fetchSize。达梦驱动默认一次取全部结果集表和内存大时会 OOM如果手动设 fetchSize 又可能因版本差异失效。建议在 JDBC URL 上显式开启游标String jdbcUrl jdbc:dm://host:5236?fetchSize1000resultSetTypeFORWARD_ONLY;全量阶段并行度不是越高越好。一张大表用 1 个并行度读取一批小表用 4 到 8 个并行度我一般是按表大小手动分配。达梦的并发读能力不像分布式库那么强并行度太高会把实例 IO 打满影响业务。4.2 LogMiner 读取参数日志窗口与事务缓冲LogMiner 每次解析的日志范围要控制在一个合理窗口。一次注册几十个归档文件会把内存撑爆而且解析速度骤降。常见做法是维护一个调度逻辑下一次解析前只注册最近 1 到 2 个归档发现没有新文件就等 1 秒有新文件再 ADD_LOGFILE。// 每个轮询周期检查新归档而不是一次注册所有 String checkSql SELECT NAME FROM V$ARCHIVE WHERE SEQUENCE# ? ORDER BY SEQUENCE#; // 发现新日志后DBMS_LOGMNR.ADD_LOGFILE 追加再重新 START_LOGMNRV$LOGMNR_CONTENTS 的读取结果默认会包含很多中间事务状态如果只关心最终提交的变更建议在 START_LOGMNR 时打开提交事务过滤选项否则同一行数据可能出现多次中间态Flink 下游做聚合时会很乱。这个属于易踩点后面避坑章详说。4.3 Checkpoint 参数状态保存是断点续传的后悔药Flink 的 checkpoint 是日志同步的后悔药没有它作业一挂就要从头扫归档。CheckpointConfig config env.getCheckpointConfig(); config.setCheckpointInterval(10_000); config.setMinPauseBetweenCheckpoints(5_000); config.setCheckpointTimeout(60_000); config.setMaxConcurrentCheckpoints(1);setCheckpointInterval设置成 10 秒既保证恢复粒度又不至于频繁触发导致日志源重复提交。setMaxConcurrentCheckpoints(1)强制串行否则两个 checkpoint 交叉LogMiner 的 lastScn 状态可能互相覆盖。状态后端我建议用 RocksDBenv.setStateBackend(new RocksDBStateBackend(hdfs:///flink/ckpt/dm-cdc));RocksDB 能支撑大状态而且增量 checkpoint 对达梦这种长时间运行的同步任务更友好。4.4 用 Nacos 管理多任务的运行时配置同步任务多了以后JDBC URL、表名、起始 SCN、归档目录这些参数散落在各个作业里改一处就打包重启效率太低。可以把配置全部丢到 Nacos 配置中心Flink 作业启动时拉取运行期间监听刷新。在 Nacos 里配置项可以这样组织# dm-cdc-task.yaml source: jdbcUrl: jdbc:dm://host:5236?useUnicodetruecharacterEncodingutf-8 username: FLINK_USER password: flink_123 archiveDir: /dm/arch tables: - SCHEMA1.TB_ORDER - SCHEMA1.TB_USER sink: kafkaBootstrap: kafka1:9092,kafka2:9092 topicPrefix: dm-cdcNacos 热更新只适合调整目标表过滤规则、topic 前缀这类动态参数。JDBC 连接串如果改了必须先重启作业。我见过有人把连接串也做成热更新结果达梦连接池里的旧连接一直挂着Flink 作业不断报错。所以我的习惯是Nacos 里只放可动态生效的配置数据库地址这类静态配置固化在作业启动参数里。5. 常见问题与避坑达梦日志同步翻车的 5 个现场5.1 有全量快照增量一直不触发现象Flink 作业跑着快照数据都出来了但 Kafka 里迟迟没有新增变更。原因LogMiner 没有识别到新的归档日志。快照后如果源库业务写入不频繁reod 日志可能长时间不切换归档文件不产生增量阶段自然没有数据。还有一种情况是作业启动时注册日志文件写死了路径后续新归档没有动态追加。解决增量循环里定期查询 V$ARCHIVE发现新归档就注册并重启 LogMiner。同时在源库主动切换日志验证链路达梦数据库常用命令里可以这样触发ALTER SYSTEM ARCHIVE LOG CURRENT;这个命令能强制当前 redo 切换到归档日常调试和故障验证特别有用。5.2 下游看到事务中间态数据顺序和源库不一致现象一行数据先被 UPDATE 成中间值再被 UPDATE 成最终值日志同步后下游短暂看到了中间值或者出现“修改前的值比修改后的值晚到”的情况。原因V$LOGMNR_CONTENTS 默认会输出事务内的所有操作包括未提交的中间状态。如果 Flink 端按行直接发射没有按 SCN 和事务 ID 缓冲就会出现乱序和中间态。解决解析时增加 SUPPRESS 中间重做记录的选项或者在 Flink 端按表把同一主键的变更按 SCN 排序后输出。对大多数 DML 场景过滤未提交事务是首选性能开销也小。5.3 作业恢复后丢数据提示找不到归档文件现象Flink 作业因为网络或资源原因挂了恢复后从 checkpoint 继续但 LogMiner 报错找不到对应的归档日志。原因达梦归档文件被清理了。归档保留时间短Flink checkpoint 里的 lastScn 已经对应到被清理的日志无法再解析。解决把 ARCH_NUM 和 ARCH_SPACE_LIMIT 调大确保归档保留时间大于“作业最长故障恢复时间 checkpoint 间隔”。同时监控 Flink 任务的延迟和 checkpoint 失败次数超过阈值立刻告警。我一般把归档空间按一天业务日志量的 1.5 倍预留。另外用达梦查询当前归档时间SELECT NAME, SEQUENCE#, TO_CHAR(FIRST_TIME,YYYY-MM-DD HH24:MI:SS) AS FIRST_TIME FROM V$ARCHIVE ORDER BY SEQUENCE# DESC;如果发现最旧归档时间离当前时间只有几分钟说明归档空间配置太紧需要立刻调整。5.4 DDL 变更导致作业停止现象业务对源表加了列Flink 作业报 Schema 不兼容直接失败。原因自定义 Source 里 buildRow 的字段结构是启动时固定好的达梦日志里出现了新列字段解析对不上。解决最简单的方案是 DDL 后手动重启作业全量快照重新拉一遍。如果表大这个成本太高。进阶做法是让 Flink 端解析 SQL_REDO 时动态感知字段变化或者定期从达梦查询最新的表结构并刷新 schema。但无论如何运维上要定义 DDL 发布流程CDC 作业升级步骤要和业务变更绑定不能只让 DBA 执行 SQL。5.5 中文乱码同步下来的字符串全是问号现象源端表里是中文Flink 解析后输出到 Kafka 变成乱码或者日志里出现“字符集不识别”的报错。原因达梦服务端字符集和 JDBC URL 指定的字符集不一致LogMiner 的 SQL_REDO 输出编码错乱。解决JDBC URL 显式指定编码并在注册 LogMiner 时打开正确的字符集选项String jdbcUrl jdbc:dm://host:5236?useUnicodetruecharacterEncodingutf-8;同时在启动 LogMiner 时如果达梦环境支持指定对应字符集 ID。如果配置都正确还乱码用 Navicat 连接达梦数据库执行SELECT * FROM V$NLS_PARAMETERS确认数据库实际字符集再反向调整 JDBC URL。6. 延迟验证与进阶从单表到整库的收尾手段6.1 用时间差做端到端延迟验证同步链路跑起来只是第一步延迟是否达标要量化。我一般在源表插入一条记录时写入一个业务时间戳下游消费时用当前时间减去业务时间戳画成一个直方图。这个差值包含了 LogMiner 轮询周期、Flink 处理时间、Kafka 传输时间。如果差值稳定在 2 秒以内说明链路健康如果持续上扬优先看 LogMiner 那边是否积压了大量未解析归档。6.2 从单表到整库按业务域拆分任务表多之后不建议一个任务解析全部库。单个任务注册过多日志文件LogMiner 解析内存吃紧而且某一张大表的突然变更会拖累其他表。常见做法是把业务域拆成多个 Flink 作业每个作业负责一个 schema 下的若干表Kafka topic 也按业务域命名。达梦本身可以并行解析多个日志文件集拆任务不会造成日志重复只要各任务的 lastScn 各自管理好即可。整库同步时还要考虑外键依赖先同步主表再同步从表或者在 Flink 端做字段补全。6.3 数据一致性回归验证日志同步最怕无声无息地丢数据。除了日常监控我每周会跑一次一致性比对从达梦导出每张表的主键和校验和再从 Kafka 按主键 replay 一份结果两边比对。这个比对作业用 Flink Batch 就能跑。做这件事还有一个额外收益能发现 LogMiner 解析选项配置错误导致的字段级偏移比如某些类型精度被截断。同步这套事情跑通容易持续不翻车难。我每次调完 LogMiner 参数都会主动触发一次日志切换观察增量是否在 2 秒内到达这个习惯让我少踩了很多坑。希望帮到你。本文还有配套的精品资源点击获取