# dm-cdc **Repository Path**: web-tiny/dm-cdc ## Basic Information - **Project Name**: dm-cdc - **Description**: 为信创达梦数据库(DM8)提供实时数据变更监听(CDC)能力 —— Canal/Debezium 在国产化数据库上的替代方案 Real-time Change Data Capture for Dameng DM8 — Canal / Debezium alternative for 信创 databases - **Primary Language**: Unknown - **License**: Not specified - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2026-08-11 - **Last Updated**: 2026-08-11 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # dm-cdc > **为信创达梦数据库(DM8)提供实时数据变更监听(CDC)能力** —— Canal/Debezium 在国产化数据库上的替代方案 > Real-time **Change Data Capture** for **Dameng DM8** — Canal / Debezium alternative for 信创 databases ## ✨ 核心功能 - **实时监听 DM8 数据变更** —— 通过 DM `DBMS_LOGMNR` 包从归档日志读取 INSERT/UPDATE/DELETE 事件,绕过 ADD_LOGFILE -2849 已知 BUG - **填补国产化数据库的 CDC 空白** —— Canal、Debezium 等主流开源 CDC 工具仅支持 MySQL/PostgreSQL/Oracle,对达梦(DM8)等信创数据库无开箱即用方案 - **输出 Debezium 兼容 JSON** —— 对齐 `op=c/u/d`、`source` 子对象、`commit_scn/ts_ms` 等 Debezium 协议字段,下游 Kafka Connect / Materialize 等可直接消费 - **真实表名/列名输出** —— START_LOGMNR 用 named param + 默认时间窗口,让 DM 的 `DICT_FROM_ONLINE_CATALOG` 直接拼写真实表名/列名;`DmDictionary`(SYSOBJECTS/SYSCOLUMNS 字典)保留为兜底,应对 DM 对已删除表/虚拟列不匹配表回退内部 ID(`"OBJ# N"`/`"COL N"`)的情况 - **at-least-once 位点保证** —— 单文件位点持久化(`tmp` + `ATOMIC_MOVE`),失败重试不丢事件;下游需幂等 - **可插拔输出通道** —— 通过 `ChangeEventSink` 接口,支持日志(默认)/ Kafka(占位),后续可对接 RocketMQ 等信创消息队列 - **配置即策略** —— 表名匹配逻辑放在 `CdcConfig.Dm.matches()` 内嵌方法,过滤规则不污染编排层 ## 🆚 与同类工具对比 | 工具 | DM 8 支持 | 部署重量 | 国产化 | 二次开发 | |---|---|---|---|---| | **dm-cdc** | ✅ 原生 | 单 fat jar (~4MB) | ✅ | 灵活 | | Canal | ❌ 仅 MySQL | 需要 Kafka | ❌ | 受限 | | Debezium | ❌ 不支持 DM | 需要 Kafka Connect | ❌ | 受限 | | Flink CDC | ⚠️ 社区版 | 需要 Flink 集群 | ⚠️ | 复杂 | > ✅ **当前状态(2026-08-11)**:CDC 链路端到端跑通且稳定,单轮可处理 7000+ 条事件 / ~1 秒(单核 ~10000 events/s),已通过 SYS_DEPT 等 30+ 字段大表的批量导入实测。详见 [`doc/KNOWN_ISSUES.md`](doc/KNOWN_ISSUES.md)。 --- ## 📖 文档导航 | 想看什么 | 章节 | |---|---| | 项目整体架构 | [§一 架构](#一架构) | | 怎么用、怎么配、怎么启动 | [§四~§八](#四dm-侧前置配置一次) | | 常见问题排查 | [§十一 常见问题](#十一常见问题) | | **为什么这么设计** | [`doc/DESIGN.md`](doc/DESIGN.md)(独立文档) | | 项目历史与未来规划 | [`doc/KNOWN_ISSUES.md`](doc/KNOWN_ISSUES.md) | --- ## 一、架构 ``` DM8 (归档日志 + DBMS_LOGMNR) ──▶ dm-cdc (纯 Java) ──▶ 日志 (JSON 行) │ ├── JDBC 调 DBMS_LOGMNR.START_LOGMNR ├── 轮询 V$LOGMNR_CONTENTS ├── SQL 解析还原 before/after └── SCN 位点持久化(data/scn.txt 单文件) ``` **纯 Java 组件结构**(按 bootstrap 顺序组装): | 层 | 类 | 职责 | |---|---|---| | 入口 | `DmCdcApplication` | main 5 行,委托给 `CdcBootstrap` | | 编排 | `CdcBootstrap` | 组合根:装配、注册 Lifecycle、启动/停止 | | 编排 | `CdcPollingLoop` | 主轮询循环:拉取 → 解析 → 过滤 → 投递 → 位点 | | 调度 | `PollingScheduler` | `ScheduledExecutorService` 包装(fixedDelay) | | 配置 | `ConfigLoader` + `CdcConfig` + `ConfigBinder` | yml 多源回退 + 反射绑定到嵌套 POJO | | 数据库 | `DmConnectionManager` | JDBC 连接工厂 | | 数据库 | `LogMinerSession` | DBMS_LOGMNR 封装(已绕过 ADD_LOGFILE -2849,lastConsumedScn 统一位点) | | 数据库 | `DmDictionary` | ID → 真实表名/列名字典(SYSOBJECTS/SYSCOLUMNS),DM 回退内部 ID 时兜底 | | 解析 | `SqlRedoParser` | SQL_REDO/UNDO 字符串 → before/after Map | | 位点 | `CheckpointStore` | 单文件原子写入(data/scn.txt) | | 输出 | `ChangeEventSink` + `LogChangeEventSink` + `KafkaChangeEventSink` + `ChangeEventSinkFactory` | Sink 抽象 + 日志实现 + Kafka 占位 + 工厂路由 | | 模型 | `ChangeEvent` | Debezium 兼容的 JSON 输出模型 | | 抽象 | `Lifecycle` + `ServiceRegistry` | 启动顺序 start()、停止逆序 stop() | | 工具 | `LogThrottler` | 同 key 错误日志 30s 限速 | --- ## 二、目录结构 ``` dm-cdc/ ├── pom.xml ├── README.md ├── config/ # 生产配置目录(可选,外部 application.yml) ├── data/ # 位点文件目录(scn.txt 单文件) ├── doc/ # 设计文档、已知问题 ├── logs/ # 运行日志(含事件 JSON) ├── src/main/ │ ├── java/com/neo/cdc/ │ │ ├── DmCdcApplication.java # main 入口(5 行,委托给 CdcBootstrap) │ │ ├── bootstrap/{CdcBootstrap,CdcPollingLoop}.java # 编排层 │ │ ├── config/{CdcConfig,ConfigLoader,ConfigBinder}.java │ │ ├── dm/{DmConnectionManager,LogMinerSession,DmDictionary}.java │ │ ├── parser/SqlRedoParser.java │ │ ├── checkpoint/CheckpointStore.java │ │ ├── sink/{ChangeEventSink,LogChangeEventSink,KafkaChangeEventSink,ChangeEventSinkFactory}.java │ │ ├── event/ChangeEvent.java # 输出模型 │ │ ├── lifecycle/{Lifecycle,ServiceRegistry}.java │ │ ├── scheduler/PollingScheduler.java │ │ └── support/LogThrottler.java │ └── resources/ │ ├── application.yml │ └── logback.xml ├── src/test/ │ └── java/com/neo/cdc/{checkpoint,parser}/*Test.java └── target/ └── dm-cdc.jar # fat jar(构建产物,单文件 ~4MB) ``` --- ## 三、关键设计决策(速览) > 完整设计推理见 [`doc/DESIGN.md`](doc/DESIGN.md)。本节是**快速摘要**。 - **位点推进语义:at-least-once** ——`lastConsumedScn` 仅在 `sink.write(ev)` 成功后推进,失败不推进→下次重试(详见 [§十一 关键设计决策](#十一关键设计决策)) - **Sink 抽象:同步契约** ——`write()` 返回 = 已投递(位点可推进);抛异常 = 未投递(位点不推进) - **事件过滤:配置即策略** ——`CdcConfig.Dm.matches()` 是 POJO 内嵌方法,过滤逻辑不下放到编排层 - **配置加载:反射 + 多源回退** ——yml → snakeyaml → `ConfigBinder` 反射写入嵌套 POJO,4 级回退 - **生命周期:有序启停** ——`Lifecycle` + `ServiceRegistry`,注册顺序:字典 → 位点 → Sink → LOGMNR会话 → 调度 --- ## 四、DM 侧前置配置(一次) 修改 `dm.ini`: ```ini ARCH_INI = 1 # 开启归档(必须) RLOG_APPEND_LOGIC = 2 # 开物理逻辑日志(必须;=1 时 UPDATE/DELETE 的 before 只有主键) LOGMNR_PARSE_LOB = 1 # 如需 LOB 同步 ``` 修改 `dmarch.ini`: ```ini [ARCHIVE_LOCAL1] ARCH_TYPE = LOCAL # REDO 日志归档类型,LOCAL 表示本地归档,REMOTE 表示远程 ARCH_DEST = /dmdata/dmarch # REDO 日志归档目标,LOCAL 对应归档文件存放路径;REMOTE 对应远程目标节点实例名 ARCH_FILE_SIZE = 1024 # 单个 REDO 日志归档文件大小,取值范围 64~2048,单位 MB,缺省值为 1024MB,即 1GB ARCH_SPACE_LIMIT = 2048 # REDO 日志归档空间限制,当所有本地归档文件达到限制值时,系统自动删除最老的归档文件。0 表示无空间限制,取值范围 1024~2147483647,单位 MB,缺省值为 0 ARCH_FLUSH_BUF_SIZE = 0 # 归档合并刷盘缓存大小,取值范围 0~128,单位 MB,缺省值为 2,0 表示不使用归档合并刷盘 ARCH_HANG_FLAG = 0 # 本地归档写入失败时系统是否挂起。取值 0 或 1。0 不挂起;1 挂起。缺省为 1。第一路本地归档系统内固定设为 1,设 0 实际也不起作用 ``` DM 版本要求:**≥ 8.1.3.100**。不支持 DSC 集群。 启动前确认: ```sql -- 检查归档 SELECT NAME, STATUS FROM V$ARCHIVED_LOG WHERE STATUS = 'A' AND DELETED = 'NO'; -- 检查 RLOG_APPEND_LOGIC SELECT PARA_VALUE FROM V$DM_INI WHERE PARA_NAME = 'RLOG_APPEND_LOGIC'; -- 检查用户权限(SYSDBA 默认有) ``` ### DM 端手动验证 LOGMNR(**排错手册**——确认 DM LOGMNR 自身能跑通后再启动 dm-cdc) 如果 dm-cdc 启动时报 `-2849` 或拉不到事件,**先在 DM 端手动验证**: ```sql -- 1. 看归档列表(确认有可用归档) SELECT NAME, FIRST_TIME, NEXT_TIME, FIRST_CHANGE#, NEXT_CHANGE# FROM V$ARCHIVED_LOG WHERE DELETED = 'NO' AND STATUS = 'A' ORDER BY FIRST_CHANGE#; -- 2. 手动添加一个归档文件(用 1-arg 形式——避开 -2849 BUG) CALL DBMS_LOGMNR.ADD_LOGFILE('/dmdata/dmarch/ARCHIVE_LOCAL1_0x5464D700_EP0_2026-08-05_17-07-58.log'); -- ↑ 用真实归档路径;EP0 是 dm 内部魔数后缀 -- 如果报 -2849:换不同归档 / 检查归档状态 / 确认归档没被锁 -- 3. 看 LOGMNR 已加载的日志(确认 ADD_LOGFILE 成功) SELECT LOW_SCN, NEXT_SCN, LOW_TIME, HIGH_TIME, LOG_ID, FILENAME FROM V$LOGMNR_LOGS; -- 4. 启动 LOGMNR -- 关键:用 named param + 默认时间窗口(非 NULL)触发 DICT_FROM_ONLINE_CATALOG, -- 让 DM 自己拼写真实表名/列名(SEG_OWNER/TABLE_NAME/SQL_REDO 全部真实名)。 -- StartScn=0 表示从所有可用归档开始;EndScn=0 表示无结束限制;Options=2128 CALL DBMS_LOGMNR.START_LOGMNR( Options => 2128, StartScn => 0, EndScn => 0, StartTime => TO_DATE('1988/1/1','YYYY-MM-DD'), EndTime => TO_DATE('2110/12/31','YYYY-MM-DD')); -- ↑ dm-cdc 内部调用:StartScn=realStartScn(位点), Options=2128,时间窗口同上 -- ↑ 若改回 positional + NULL 时间(0,0,NULL,NULL,2128)→ SEG_OWNER 变数字 ID、 -- TABLE_NAME 变 NULL、SQL_REDO 变 "OBJ# N"/"COL N"(需 DmDictionary 解析) -- 5. 看 LOGMNR 解析出的事件(确认有 INSERT/UPDATE/DELETE) SELECT OPERATION_CODE, SCN, SQL_REDO, SQL_UNDO, TIMESTAMP, SEG_OWNER, TABLE_NAME, SESSION_INFO FROM V$LOGMNR_CONTENTS WHERE "TABLE_NAME" = 'TEST_DATA'; -- ↑ 当前形态下 TABLE_NAME 是真实表名(如 TEST_DATA),SEG_OWNER 是真实 schema 名 -- ↑ 如果 OPERATION_CODE=1/2/3 都是 DML -- 6. 清理(必须 END,否则下次 START_LOGMNR 会失败) CALL DBMS_LOGMNR.END_LOGMNR(); ``` **排错指引**: | 现象 | 排查 | |---|---| | ADD_LOGFILE 报 `-2849` | 用单参数形式(去掉 option);检查归档状态;换不同归档文件 | | V$LOGMNR_CONTENTS 返回 0 行 | `RLOG_APPEND_LOGIC` 设 2;确认归档内确实有 DML(用 `ALTER SYSTEM ARCHIVE LOG CURRENT` 强制切归档) | | 解析不到目标表 | 当前 START_LOGMNR 形态下 `TABLE_NAME`/`SEG_OWNER`/`SQL_REDO` 都是真实名;若仍遇到 `"OBJ# N"`/`"COL N"`(DM 对已删除表等回退内部 ID),`DmDictionary` 兜底解析 | --- ## 五、构建 ```bash cd dm-cdc mvn clean package # 产物: # target/dm-cdc.jar 单个 fat jar(约 4MB,含所有依赖) ``` 构建出**单个 fat jar**,所有依赖(DM JDBC、fastjson、logback、snakeyaml、commons-lang3、slf4j-api)都合并进了 `dm-cdc.jar`,启动时无需挂 lib/ 目录。 --- ## 六、配置 `application.yml` 完整带注释的配置文件参考 [`src/main/resources/application.yml`](src/main/resources/application.yml)。**核心字段概览**: ```yaml dm-cdc: # ===== 源数据库(达梦 DM8)===== dm: host: 192.168.117.8 # 必填:DM 地址 port: 5236 # 必填:DM 端口(默认 5236) user: TEST_LISTENER # 必填:建议 SYSDBA 或带 LOGMNR 权限的用户 password: TEST_LISTENER database: # 可选:JDBC URL ?schema= 参数 schemaList: [TEST_LISTENER] # 必填:监听 schema tableList: [TEST_LISTENER.TEST_DATA, TEST_LISTENER.SYS_DEPT] # 可选:留空=所有表 startScn: -1 # 首次 SCN;-1=从最新归档;之后从 scn.txt 恢复 # ===== LOGMNR ===== logMiner: fetchSize: 1000 # JDBC fetch size(默认 1000) sleepDefaultMs: 1000 # 轮询间隔(默认 1000ms;建议 ≥200ms) sessionMaxMs: 60000 # 会话最大时长,到期重建 archiveDest: "" # 归档目的地过滤;留空=全选 # sleepMinMs/sleepMaxMs/sleepIncrementMs/offlineCatalog 当前未生效 # ===== 位点 ===== checkpoint: dir: ./data # scn.txt 所在路径 # ===== 输出通道 ===== sink: type: log # log | kafka(占位) # kafka 专用:topic/bootstrapServers/clientId/acks # ===== 杂项 ===== misc: statsIntervalSec: 10 # STATS 输出间隔 includeRawSql: false # 是否在 JSON 中追加 sql_redo/sql_undo ``` > **生产建议**:把 `application.yml` 复制到 `config/` 目录,用 `-Ddm-cdc.config=config/application.yml` 启动,避免覆盖 jar 内的默认配置。 > > **配置加载优先级**:`1. -Ddm-cdc.config=` > `2. ./config/application.yml` > `3. classpath:/application.yml` > `4. CdcConfig` 内嵌默认值 > > **完整字段含义、调优建议、坑**全部带注释,参考 [`application.yml`](src/main/resources/application.yml) 即可。 --- ## 七、运行 ```bash # 使用 jar 外的配置文件(推荐生产方式) java -Ddm-cdc.config=config/application.yml -jar target/dm-cdc.jar # 或直接用 jar 内的 application.yml java -jar target/dm-cdc.jar ``` 优雅停机:`Ctrl+C` 或 `kill `,会自动关闭 LOGMNR + 关闭位点。 --- ## 八、日志输出格式 事件走独立的 logger `dm-cdc.event`,每条事件一行,格式 `监听到数据:{json}`(前缀便于在控制台/文件里区分,JSON 在行尾,下游解析不受影响): ### 默认输出(`misc.includeRawSql=false`) ```json { "op": "c", // c=insert, u=update, d=delete "before": null, // UPDATE/DELETE 有值 "after": { "ID": 31, "DATA_STR": "8768687687" }, // INSERT/UPDATE 有值 "source": { "schema": "TEST_LISTENER", "table": "TEST_DATA", "scn": 666211203, "version": "1.0.0", "connector": "dm", "name": "dm8", "ts_ms": 1786007306413 // DM 端事务时间 }, "ts_ms": 1786006972339 // 本地拉取时间(差值 = 采集延迟) } ``` ### 调试模式(`misc.includeRawSql=true`) 在事件末尾追加 DM 原始 SQL: ```json { "op": "c", "before": null, "after": { "ID": 31, "DATA_STR": "8768687687" }, "source": { "schema": "TEST_LISTENER", "table": "TEST_DATA", "scn": 666211203, ... }, "ts_ms": 1786006972339, "sql_redo": "insert into \"TEST_LISTENER\".\"TEST_DATA\"(\"ID\",\"DATA_STR\") values('31', '8768687687');", "sql_undo": null } ``` `logback.xml` 中 `dm-cdc.event` logger 配到独立的 EventFile(`./logs/dm-cdc-events.log`,按日期+50MB 滚动,gz 压缩,保留 30 份)+ Console。若需对接下游(例如 filebeat→ES),直接消费该事件文件即可。 > 日志文件一览(`logs/`): > - `dm-cdc.log` —— 主运行日志(含 DEBUG/INFO/ERROR) > - `dm-cdc-error.log` —— 仅 ERROR > - `dm-cdc-events.log` —— 变更事件 JSON 行(供下游消费) > **格式对齐**:事件 JSON 严格对齐 Debezium 协议(`op=c/u/d`、字段命名 `commit_scn/ts_ms`、子对象 `source`),下游 Kafka Connect / Materialize 等 Debezium 兼容工具无需写自定义反序列化器。 --- ## 九、对接方消费示例 ### 简单 grep ```bash tail -F ./logs/dm-cdc-events.log | grep '"op":"c"' ``` ### 解析日志为结构化数据(Python) ```python import json, re pat = re.compile(r'(\{.*\})$') # JSON 在行尾 with open('./logs/dm-cdc-events.log', encoding='utf-8') as f: for line in f: m = pat.search(line) if not m: continue ev = json.loads(m.group(1)) # 业务处理(ev.op / ev.before / ev.after / ev.source.*) ``` --- ## 十、限制(MVP 版本) - **RLOG_APPEND_LOGIC=1 时拿不到完整 before 旧值**(UPDATE/DELETE 的 before 只有主键):需 DM 端设 `RLOG_APPEND_LOGIC=2` 才有 SQL_UNDO - **只处理 DML(INSERT/UPDATE/DELETE)**,DDL/LOB/事务元数据暂不下发(第二版扩展) - **同步表必须有主键**(DM LOGMNR 解析要求) - **不支持 DSC 集群** - **单进程**,无 HA;HA 需自己加主备 + 共享位点 - **位点存储**用本地单文件(data/scn.txt),多机不共享;如需共享改 Redis/ZK/etcd - **Kafka 输出为占位**:当前 `sink.type=kafka` 仅记录日志不真正投递,**生产请用 `log` 模式 + filebeat 消费 events.log** --- ## 十一、常见问题 | 现象 | 排查 | |---|---| | 启动报错 "没有可用的归档日志" | DM 没开归档;执行 `ALTER DATABASE MOUNT;` + 检查 `ARCH_INI` | | `SQL_REDO` 拿到的字段不全 | `RLOG_APPEND_LOGIC` 没开或设为 1,改为 2 | | UPDATE 拿不到 before 旧值 | 同上,必须 2 | | 会话到期后数据断层 | 检查 `logMiner.sessionMaxMs`,并确认 SCN 位点已落盘 | | 日志文件没事件 | 检查 logger `dm-cdc.event` 是否被覆盖、确认 `logMiner.fetchSize` 不为 0 | | 日志量太大 | 把 `dm-cdc.event` 调为 `warn`(停止打印事件,但位点仍推进,**会丢下游**) | | 换了数据库查不到事件 | 检查 `data/scn.txt` 是否残留旧位点(指向已不存在的归档),必要时 `rm scn.txt` 重启 | | 看到 `lastConsumedScn=0` 卡住 | 程序启动后第一次 poll 还没拉到事件(正常),下一轮自动重试 | | 位点跨归档边界后查不到新事件 | 这是 0.731s 内 7000+ 事件能正确处理的边界场景;归档查询已修复为 `WHERE NEXT_CHANGE# > lastConsumedScn`,确保跨边界归档被加载。如旧版本仍卡住,**删除 `data/scn.txt` 重启**让程序重新加载所有归档 | --- ## 十二、参考与已知问题 - **必读**:[`doc/KNOWN_ISSUES.md`](doc/KNOWN_ISSUES.md) — 项目状态、历史 bug、未来规划、性能基线 - **设计文档**:[`doc/DESIGN.md`](doc/DESIGN.md) — 设计决策推理过程、达梦 CDC 原理、模块结构详解 - [`doc/DM8备份与还原.pdf`](doc/DM8备份与还原.pdf) — DM 8 归档配置机制(参考用) - 底层原理:DM `DBMS_LOGMNR` 包 + 归档日志 + 物理逻辑日志(类似 Oracle LogMiner 模型) --- ## 文档职责分工 | 文档 | 定位 | 目标读者 | |---|---|---| | `README.md` | **用户手册**——怎么用、怎么配、怎么启动、常见问题 | 用户 / 运维 | | `doc/DESIGN.md` | **设计决策**——为什么这么选、架构细节、模块职责 | 开发者 | | `doc/KNOWN_ISSUES.md` | **项目档案**——历史 bug、未来规划、性能基线、踩坑指南 | 维护者 / 接手人 |