
简介这份资源面向需要在Flink生态中实现达梦数据库实时同步的开发者与数据工程师聚焦基于日志解析的变更数据捕获场景可用于数据仓库同步、实时报表、数据监控与告警等事件驱动应用。包内共5个文件以jar连接器与驱动包为主另含zip示例工程、sql初始化脚本和docx用户手册压缩包约35.48MB覆盖从依赖引入到作业配置的完整链路。资源提供达梦CDC连接器及配套参考程序读者可据此快速搭建Flink CDC作业理解日志解析、插入更新删除捕获与低延迟同步的实现方式并对照手册完成连接参数配置与SQL或Java两种同步路径的落地。目前已有2083人学习下载适合希望降低源库压力、构建高可靠实时数据流的进阶开发者参考。1. FlinkCDC 接达梦为什么日志级实时同步值得做达梦数据库在国产化替代里出现得越来越频繁很多团队把 Oracle、MySQL 上的业务迁到 DM8 之后第一个撞上的问题就是原来那套基于 binlog 的实时同步链路断了。达梦没有 MySQL 那种开箱即用的 binlog 生态但它在 DM8 之后提供了逻辑日志Logic Log能力配合归档日志可以做到不侵入业务表、不靠触发器、不靠时间戳轮询的增量捕获。FlinkCDC 从 2.x 开始支持了通用的增量快照框架社区里也有人把达梦接进了这套体系。这篇笔记讲的就是怎么用 FlinkCDC 把达梦的日志级变更实时同步出去中间要开哪些库级开关、连接器参数怎么配、哪些坑我踩过。适合两类人看一类是正在做国产数据库实时数仓、CDC 入湖入仓的工程师另一类是手上已经有 Flink 集群想把达梦接进现有同步链路、又不想改业务代码的人。读完你应该能自己判断这套方案在你的环境里能不能落地以及落地时要先动哪几个配置。2. 达梦日志级 CDC 的前置条件归档、逻辑日志与权限2.1 达梦的日志体系和 MySQL binlog 不是一回事MySQL 的 binlog 是语句级或行级的逻辑日志直接就能解析。达梦的物理归档日志ARCHIVELOG记录的是页级变更不能直接拿来还原成 INSERT/UPDATE/DELETE。真正能用于 CDC 的是达梦的逻辑日志功能它需要在数据库实例上显式开启并且依赖归档模式。换句话说达梦做 CDC 有两道门第一道是归档模式必须开第二道是逻辑日志必须开。少一道连接器连上去也拿不到变更。常见做法是先在测试库上确认这两项状态再动生产。生产库开归档和逻辑日志通常需要重启实例这个窗口要提前和业务方对齐。我一般会先在备库或者测试环境把整条链路跑通再上生产。2.2 开启归档和逻辑日志的具体命令下面这些命令用达梦的 disql 或者管理工具执行都可以。注意路径要换成你自己的实际归档目录并且确保达梦实例的操作系统用户对该目录有写权限。-- 1. 开启归档模式需要 MOUNT 状态通常要重启实例 ALTER DATABASE MOUNT; ALTER DATABASE ARCHIVELOG; ALTER DATABASE ADD ARCHIVELOG DEST/dmdata/arch, TYPELOCAL, FILE_SIZE1024, SPACE_LIMIT102400; ALTER DATABASE OPEN; -- 2. 开启逻辑日志不同 DM8 小版本语法略有差异以实际版本为准 SP_SET_PARA_VALUE(1, ENABLE_LOGIC_LOG, 1); -- 3. 确认归档和逻辑日志状态 SELECT ARCH_MODE FROM V$DATABASE; SELECT PARA_NAME, PARA_VALUE FROM V$DM_INI WHERE PARA_NAME IN (ENABLE_LOGIC_LOG);逻辑说明ALTER DATABASE ARCHIVELOG把实例切到归档模式这是逻辑日志能持续落盘的前提。ADD ARCHIVELOG指定归档路径、单文件大小MB和空间上限MB空间上限设太小会导致归档写满后实例挂起这个参数我一般给到 100GB 以上。SP_SET_PARA_VALUE是达梦改参数的系统过程第一个参数 1 表示动态参数部分版本需要重启才生效改完务必用V$DM_INI查一次实际值。参数说明FILE_SIZE建议 1024MB 起步太小会频繁切文件SPACE_LIMIT按你每天归档增量乘以保留天数估算宁可给大。逻辑日志开启后对写入性能有轻微影响实测在 5% 以内但具体要看业务写入模式。2.3 给 CDC 单独建一个只读账号不要用 SYSDBA 去跑连接器。达梦的权限模型和 Oracle 接近CDC 账号需要能读系统视图、能读业务表、能访问逻辑日志。最小权限集大概是这些CREATE USER CDC_USER IDENTIFIED BY Cdc2024; GRANT SELECT ON V$DATABASE TO CDC_USER; GRANT SELECT ON V$DM_INI TO CDC_USER; GRANT SELECT ON 你的业务表 TO CDC_USER; -- 逻辑日志相关权限按实际版本授予部分版本需要 RESOURCE 角色 GRANT RESOURCE TO CDC_USER;逻辑说明V$DATABASE和V$DM_INI用来做启动时的状态自检业务表的 SELECT 权限是增量快照阶段全量读需要的。逻辑日志的读取权限在不同 DM8 版本里授予方式不完全一样有的版本需要额外角色建议先用这个账号手动连一次、跑一条查询验证。提示达梦对密码大小写和特殊字符敏感连接串里的密码如果含或#记得做 URL 编码否则连接器解析会出错。3. FlinkCDC 达梦连接器的选型与作业搭建3.1 连接器从哪来社区版还是自研FlinkCDC 官方连接器列表里达梦不在第一梯队。实际落地有两条路一是用社区里已经有人维护的达梦 CDC 连接器 jar二是基于 FlinkCDC 的 IncrementalSource 框架自己实现一个 Source。前者省事但版本兼容性和后续维护要自己评估后者可控但要投入人力。我一般会先看社区连接器支持的 DM8 版本和 Flink 版本是否和现有集群对得上。对不上就别硬凑自己包一层反而更快。下面给的配置以通用增量快照框架的写法为准具体类名和参数名以你拿到的连接器为准思路是通的。3.2 依赖和作业骨架把连接器 jar 放到 Flink 的 lib 目录或者用 Maven 打进 fat jar。下面是一个 DataStream 方式的作业骨架用 FlinkCDC 的 Source 构建器。// 依赖以实际连接器坐标为准这里示意结构 // flink-connector-dm-cdc // flink-connector-base import org.apache.flink.streaming.api.environment.StreamExecutionEnvironment; import org.apache.flink.api.common.eventtime.WatermarkStrategy; import com.ververica.cdc.connectors.base.source.IncrementalSource; import com.ververica.cdc.connectors.base.options.StartupOptions; public class DmCdcJob { public static void main(String[] args) throws Exception { StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(10000); // 10 秒一次 checkpoint保证断点续传 IncrementalSourceString source IncrementalSource.Stringbuilder() .hostname(10.0.0.21) .port(5236) // 达梦默认端口 .database(BIZDB) .username(CDC_USER) .password(Cdc2024) .tableList(BIZDB.ORDERS, BIZDB.ORDER_ITEM) .startupOptions(StartupOptions.initial()) // 先全量再增量 .deserializer(new DmChangeDeserializer()) // 自定义反序列化 .build(); env.fromSource(source, WatermarkStrategy.noWatermarks(), dm-cdc) .print(); env.execute(dm-cdc-sync); } }逻辑说明enableCheckpointing是 CDC 作业的命根子没有 checkpoint 就没有断点续传作业重启只能从头全量。StartupOptions.initial()表示先做一次全量快照再切到增量日志读取如果只想读增量用latest()。tableList要写「库名.表名」的全限定形式达梦对大小写敏感表名建的时候是大写就写大写。参数说明hostname和port是达梦实例地址默认 5236。deserializer负责把逻辑日志记录转成你要的结构这部分通常要自己实现把达梦的变更类型映射成 Flink 的 RowData 或 JSON 字符串。checkpoint 间隔我一般给 10 到 30 秒太短会增加达梦侧读取压力太长故障恢复时重放的数据多。3.3 全量加增量的切换逻辑FlinkCDC 的增量快照框架会把表按主键切成 chunk并行做全量读同时记录一个「快照位点」。全量读完后从位点对应的日志位置开始读增量。达梦这边位点通常对应逻辑日志的 LSN 或者归档日志的偏移。这里最容易翻车的地方是全量读期间发生的变更如果位点没对齐会丢或者重复。我的做法是全量阶段用连接器自带的 chunk 切分不要自己写 SELECT 全表增量起点严格用连接器返回的位点不要手动指定时间戳。达梦的逻辑日志位点和 Oracle SCN 类似手动换算很容易错。注意如果业务表没有主键增量快照框架没法切 chunk只能退化成单并发全量读大表会很慢。上 CDC 之前先确认目标表都有主键或唯一索引。4. 同步链路调优并发、位点与反序列化4.1 并发度怎么定FlinkCDC 的并行度分两层Source 并行度和全量阶段的 chunk 并行度。Source 并行度决定增量阶段同时读几个日志流达梦侧逻辑日志读取通常是单流的所以 Source 并行度给太高没意义一般 1 到 2 就够。全量阶段的 chunk 并行度可以给高按表数量和主键分布来定我一般给 4 到 8。// 全量阶段 chunk 切分和并行相关参数参数名以实际连接器为准 IncrementalSource.Stringbuilder() .splitSize(8096) // 每个 chunk 的行数上限 .splitMetaGroupSize(2048) // 元数据分组大小 .chunkKeyColumn(ID) // 切分用的主键列 // ...逻辑说明splitSize控制单个 chunk 的行数太小会导致 chunk 数量爆炸、元数据开销大太大则单 chunk 读得久、并行度上不去。chunkKeyColumn指定用哪一列切必须是数值型或可比较的主键字符串主键切分效率会差一些。参数说明splitSize我一般按单表总行数除以期望 chunk 数来估比如一亿行的表想切 100 个 chunk就给 100 万。splitMetaGroupSize影响元数据读取的批量大小默认值通常够用不用频繁调。4.2 位点管理和断点续传位点存在 Flink 的 checkpoint 和 savepoint 里。作业正常跑的时候每次 checkpoint 会把当前读到的日志位点存下来。作业失败重启从最近一次成功的 checkpoint 恢复位点之后的变更会重放。这里有个坑达梦的归档日志如果被清理了而 checkpoint 里的位点对应的归档文件已经不在作业就恢复不了只能重新全量。所以生产上要做两件事一是归档日志的保留时间要大于 checkpoint 的最大保留时间二是监控归档目录的使用率别等写满了才发现。我一般把归档保留设成 7 天checkpoint 保留设成 3 天留足余量。# 查看达梦归档目录使用情况在数据库服务器上执行 du -sh /dmdata/arch ls -lt /dmdata/arch | head -20逻辑说明du看总占用ls -lt按时间列出最近的归档文件确认最新的归档在持续生成。如果发现归档停止生成先查实例是不是卡在归档写满的状态。4.3 反序列化要处理的达梦特有类型达梦有些数据类型在逻辑日志里的表示和 MySQL 不一样反序列化时要特别处理。比如NUMBER精度、TIMESTAMP时区、CLOB/BLOB大字段。大字段在逻辑日志里可能是分段记录的反序列化时要拼装。// 反序列化里处理达梦时间戳和数值的示意 if (TIMESTAMP.equals(columnType)) { // 达梦时间戳可能带纳秒精度转成 Flink 的 TimestampData 时注意精度截断 long millis dmTimestamp.toMillis(); return TimestampData.fromEpochMillis(millis); } if (NUMBER.equals(columnType)) { // NUMBER 可能超出 double 精度用 BigDecimal 接 return DecimalData.fromBigDecimal(new BigDecimal(rawValue), precision, scale); }逻辑说明达梦的TIMESTAMP精度可能到纳秒Flink 的TimestampData是毫秒转换时会丢精度如果业务对纳秒敏感要在下游单独存原始值。NUMBER用BigDecimal接别用 double否则金额类字段会出现精度丢失这种问题上线后很难查。参数说明precision和scale从达梦的列元数据里取不要写死。大字段建议在反序列化阶段就决定是透传还是截断透传会占内存截断会丢数据按业务需求定。5. 避坑与排查达梦 CDC 常见的五类翻车5.1 现象作业启动报「逻辑日志未开启」原因ENABLE_LOGIC_LOG参数没生效或者实例没在归档模式。有些 DM8 版本改完参数需要重启只动态改不重启不生效。解决用SELECT PARA_NAME, PARA_VALUE FROM V$DM_INI WHERE PARA_NAMEENABLE_LOGIC_LOG确认实际值是 0 就重启实例。同时确认ARCH_MODE是 Y。5.2 现象全量阶段读得动一切到增量就没数据原因全量快照的位点和增量起点没对齐或者逻辑日志的读取权限不够。也有可能是业务表在全量期间没有新变更误以为没数据。解决先手动在业务表插一条数据看作业有没有输出。没有的话查 CDC 账号对逻辑日志的读取权限再看连接器日志里增量起点位点是不是比当前日志位点还大说明位点算错了。5.3 现象作业跑一段时间后 OOM原因反序列化时把大字段全量缓存在内存或者 checkpoint 状态太大。达梦的 CLOB 字段如果很大逐条缓存会撑爆内存。解决大字段改成流式处理或者截断checkpoint 状态后端换成 RocksDB并且调大托管内存。同时检查是不是有表没主键导致 chunk 元数据膨胀。5.4 现象归档目录写满实例挂起原因SPACE_LIMIT设小了或者归档清理策略没配。达梦归档写满后实例会挂起业务全停。解决紧急情况先扩容或者清理旧归档恢复实例。长期方案是把SPACE_LIMIT调大配一个定时清理脚本保留时间大于 checkpoint 保留时间。# 归档清理脚本示意删除 7 天前的归档文件 find /dmdata/arch -name *.log -mtime 7 -delete逻辑说明-mtime 7表示修改时间在 7 天前-delete直接删除。这个脚本要放在 crontab 里定时跑跑之前确认 checkpoint 保留时间小于 7 天否则恢复时会找不到归档。5.5 现象同步到下游的数据有重复原因作业失败重启后从 checkpoint 位点重放位点之后已经同步过的数据会再发一次。这是 at-least-once 语义的正常表现。解决下游做幂等用主键做 upsert或者在 Flink 作业里加去重算子。如果业务要求 exactly-once需要下游支持事务两阶段提交复杂度会高不少。我一般优先让下游幂等比在 Flink 里做 exactly-once 省事。6. 验证同步正确性的一个笨办法和一条经验验证 CDC 同步对不对最靠谱的不是看日志是对数据。我常用的笨办法是在业务库上开一个事务对目标表做一批有特征的增删改记下操作前后的行数和关键字段的校验和然后去下游查同样的校验和。特征数据要包含边界值比如数值型的最大值最小值、字符串的空串和超长串、时间的边界。-- 在达梦侧生成校验和示意按实际字段调整 SELECT COUNT(*) AS cnt, SUM(CRC32(CAST(ID AS VARCHAR) || NAME || CAST(AMOUNT AS VARCHAR))) AS chk FROM BIZDB.ORDERS WHERE UPDATE_TIME 2024-06-01 00:00:00;逻辑说明CRC32把多个字段拼起来算一个校验值两边对比这个值就能快速判断数据是否一致。达梦和下游数据库的CRC32实现可能不同如果对不上换成MD5或者直接在 Flink 侧算好了写下去。WHERE条件限定在测试数据的时间范围内避免全表扫描。参数说明校验字段要选能唯一标识一行并且变更时会变的列UPDATE_TIME如果有的话最好用上。没有UPDATE_TIME就用主键范围限定。一条经验达梦 CDC 这套东西配置本身不难难的是运维。归档空间、逻辑日志位点、checkpoint 保留这三样任何一个出问题都会导致同步中断甚至要重新全量。我现在的习惯是上线第一天就把这三个指标的监控告警配好归档使用率超过 70% 就告警checkpoint 连续失败两次就告警逻辑日志位点落后超过阈值也告警。别等业务方打电话来说数据不对了才去查那时候往往已经丢了一大段。希望帮到你。本文还有配套的精品资源点击获取