简介:这份资源面向需要在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, TYPE=LOCAL, FILE_SIZE=1024, SPACE_LIMIT=102400'; 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 "Cdc@2024"; 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,保证断点续传 IncrementalSource<String> source = IncrementalSource.<String>builder() .hostname("10.0.0.21") .port(5236) // 达梦默认端口 .database("BIZDB") .username("CDC_USER") .password("Cdc@2024") .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.<String>builder() .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_NAME='ENABLE_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 连续失败两次就告警,逻辑日志位点落后超过阈值也告警。别等业务方打电话来说数据不对了才去查,那时候往往已经丢了一大段。希望帮到你。
本文还有配套的精品资源,点击获取