ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

FlinkCDC同步达梦数据库:基于redo/archive日志的实时同步实践

2026/10/1 9:13:11 拓冰建站 浏览量
FlinkCDC同步达梦数据库:基于redo/archive日志的实时同步实践 简介FlinkCDC与达梦数据库结合的实时同步方案面向需要构建实时数据管道的数据工程师、Flink开发者和运维人员。方案基于达梦数据库日志解析实现变更数据捕获无需侵入业务表可稳定获取插入、更新、删除等操作并以低延迟方式汇入Flink流处理任务适用于实时数仓、数据同步、监控告警等场景。资料包共5个文件约35.48MB包含达梦CDC连接器jar包、JDBC驱动jar、SQL初始化脚本、参考程序压缩包和一份用户手册docx从环境配置到实际调用均有覆盖。已有2083人学习实践价值得到验证。通过该资源可掌握Flink CDC连接达梦数据库的配置方式、Java与SQL两种任务写法并借助手册快速排查常见问题缩短从零到可运行同步任务的周期适合具备一定Flink基础、希望将达梦纳入实时数据体系的开发者。1. FlinkCDC 与达梦数据库的组合正在成为日志级实时同步的标准答案业务库跑在达梦上下游数仓要实时拿变更上游又没有 MySQL 那样的现成 binlog 组件——这是很多数据工程师第一次搜“FlinkCDC 达梦数据库”时的真实处境。FlinkCDC 本身是一套基于日志抓取变更数据的计算框架达梦数据库同样靠 redo/archive 日志记录每一次数据修改两者结合就能在不侵入业务表、不增加应用代码的前提下把增删改操作以近乎实时的速度同步到 Kafka、其他数据库或数据湖。适合的人群很明确数据同步开发、数仓维护、以及正在做达梦到异构目标库迁移的团队。这篇文章不聊概念只讲怎么把这条链路搭起来以及哪些参数和坑决定了你到底是上线还是返工。2. 日志级同步为什么是达梦实时同步的最优解redo/archive 与三种方案的取舍2.1 达梦日志是怎么组织的redo、archive 与用户态解析接口达梦数据库在内存里维护一组 redo 日志缓冲区事务提交时把变更记录刷到联机 redo 日志文件。联机日志写满后会切换如果要支持任意时间点恢复或日志级解析就必须开启归档模式把历史 redo 内容持续输出到 archive log 文件。FlinkCDC 要做的实时同步本质上是“替下游读这堆日志并把 insert/update/delete 还原成结构化事件”。达梦不像 MySQL 那样公开 binlog 的解析协议但提供了类 Oracle 的日志挖掘包DBMS_LOGMNR。通过这个包可以加载 redo/archive 日志文件然后从V$LOGMNR_CONTENTS动态视图里读出每条操作的 schema、表名、操作类型、前后镜像和 row_id。FlinkCDC 的达梦连接器虽然不像 MySQL 连接器那样开箱即用但按这条路径自己封装一个 source 并不复杂而且完全是数据库原生接口不需要旁路抓包或 hook 存储引擎。要注意的是达梦还有一类基于触发器或时间戳的同步方案但那些都不走日志。触发器和应用双写都会侵入业务时间戳轮询只能发现“改了”发现不了“删了”也无法拿到变更前镜像。所以只要目标是“日志级实时同步”核心工作就集中在两件事上确认归档日志连续可读然后消费V$LOGMNR_CONTENTS。2.2 三种常见同步方案对比轮询、触发器、日志解析我把在项目里见过的达梦同步路子捋一遍按可靠度和改造成本排序也方便你判断为什么最终要落在 FlinkCDC 上。方案延迟对源库侵入删除识别断点续传适合场景应用双写/定时任务分钟级高改业务代码一般无临时报表触发器中间表秒级中每张表加触发器能弱小表、低频JDBC 轮询时间戳秒级低加索引不能弱只增不删DBMS_LOGMNR 日志解析秒级极低开归档能强生产级实时同步日志解析方案里FlinkCDC 扮演的角色有两个层面。第一层是“调度和容错”Flink 负责周期性拉取日志、记录消费位点、通过 checkpoint 保证不丢数据第二层是“事件流处理”把日志内容变换成统一的记录结构再交给下游 sink。相比自己写一个 Java 定时任务去读V$LOGMNR_CONTENTSFlink 的分布式快照机制解决了“日志消费到哪了”这个最头疼的问题任务重启后能从最近一次 checkpoint 恢复不会重复扫一大段历史日志。另外达梦支持生成列、压缩表、水平分区表等特性日志里记录的物理信息比逻辑信息多。FlinkCDC 接入时通常还会做一个“日志记录 → 逻辑表名 → 目标库表名”的映射层。这个映射层在方案选型时就该设计好否则后面接 Kafka 或异构库时会非常被动。2.3 FlinkCDC 对达梦的适用边界版本、归档模式与权限FlinkCDC 官方连接器列表里没有达梦这不代表不能做而是意味着你要敢于用社区方案或自研扩展。我一般建议先确认三件事再做技术选型。第一版本与日志格式。达梦 8 是当前主流版本生产环境大多也是 8.x达梦 7 的DBMS_LOGMNR在部分小版本上输出字段有差异SQL 解析逻辑要单独适配。建议开发前先在测试实例上跑一遍日志挖掘确认V$LOGMNR_CONTENTS里能否完整读到字段类型和长度。第二归档模式。必须开启归档否则日志文件被覆盖FlinkCDC 是抓不到任何历史数据的。第三数据库账号权限。同步账号至少需要SELECT ANY TABLE的读权限以及执行DBMS_LOGMNR包和访问V$LOGMNR_CONTENTS视图的权限最省事的做法是给 DBA 角色但生产环境更推荐按包授权。边界还有一个容易被忽略的点达梦日志里 DDL 操作的解析并不完整DBMS_LOGMNR更擅长还原 DML。如果你指望通过 FlinkCDC 同步ALTER TABLE之类的结构变更多半会失望。常规做法是 DDL 另走运维通道下发FlinkCDC 只负责数据变更。这个预期管理很重要省得项目验收时被一条“为什么加了个字段没同步”的需求拖死。3. 把 FlinkCDC 跑在达梦上从开启归档到最小同步任务3.1 达梦端三步准备开启归档、创建账号、授权第一步是在达梦实例上开启归档。不同版本的具体命令略有差异以达梦 8 为例在 disql 里执行下面这段-- 1. 配置归档目录和文件大小目录要先创建好 ALTER DATABASE ADD ARCHIVELOG DEST/dm8/arch, TYPElocal, FILE_SIZE512, SPACE_LIMIT10240; -- 2. 修改为归档模式需要数据库处于 mount 状态 ALTER DATABASE SET ARCHIVELOG; -- 3. 打开数据库 ALTER DATABASE OPEN;逻辑说明第一条 SQL 把归档日志写到/dm8/archFILE_SIZE512表示单个归档文件 512MBSPACE_LIMIT10240表示归档目录最多占 10GB超过后达梦会按策略清理旧归档。FILE_SIZE调小会增加归档文件数量调大则单文件恢复时间变长512MB 是折中值。SPACE_LIMIT决定了日志能回溯多久必须结合同步链路允许的最大中断时间一起定。第二步是创建专用同步账号不要直接拿 DBA 账号跑同步任务-- 创建账号并指定默认表空间 CREATE USER FLINKCDC IDENTIFIED BY flinkcdc#2024 DEFAULT TABLESPACE MAIN; -- 授权读任意表 执行日志挖掘包 访问视图 GRANT SELECT ANY TABLE TO FLINKCDC; GRANT EXECUTE ON DBMS_LOGMNR TO FLINKCDC; GRANT SELECT ON V$LOGMNR_CONTENTS TO FLINKCDC; GRANT SELECT ON V$ARCHIVED_LOG TO FLINKCDC;逻辑说明DBMS_LOGMNR是达梦的日志挖掘包需要建在 SYS 模式下普通用户默认没有执行权V$LOGMNR_CONTENTS是动态视图单独授权。这里给SELECT ANY TABLE是因为日志挖掘的结果会涉及所有被同步的业务表如果你只想同步特定模式可采用更细粒度的授权但后续每加一张表都得补权限运维成本高。3.2 Flink 工程依赖与项目结构FlinkCDC 连接达梦时Flink 侧要引入的依赖比用 MySQL 连接器多一些。因为官方没有直接给达梦连接器所以要自己引入达梦 JDBC 驱动并添加 Flink CDC 的基础依赖来复用它的状态管理和 SourceFunction 框架。以下是一个可编译的 Maven 依赖配置dependencies !-- Flink CDC 核心依赖 -- dependency groupIdcom.ververica/groupId artifactIdflink-cdc-base/artifactId version2.3.0/version /dependency !-- Flink DataStream 依赖 -- dependency groupIdorg.apache.flink/groupId artifactIdflink-streaming-java_2.12/artifactId version1.16.0/version /dependency !-- 达梦 JDBC 驱动本地仓库安装不依赖远程仓库 -- dependency groupIdcom.dameng/groupId artifactIdDmJdbcDriver18/artifactId version8.1.1.90/version /dependency /dependencies参数说明flink-cdc-base提供了SourceFunction的生命周期管理和检查点协调版本建议和 Flink 主版本匹配达梦 JDBC 驱动一般不在中央仓库需要从达梦安装目录里拿DmJdbcDriver18.jar手动安装到本地 Maven 仓库。注意驱动包名和版本号在不同达梦版本上有差异实际以你安装目录下为准别直接照抄。目录结构上我习惯分三层source日志抓取→ transform日志记录转逻辑记录→ sink写入目标。transform 层单独抽出来很有必要因为日志里记录的列顺序和类型与业务表不一定一致而且达梦日志对大字段如 CLOB、BLOB的还原方式比较特殊这一层的解析迟早要改。3.3 写一个能跑的最小同步任务下面这段代码是一个最小可运行的 Flink 任务骨架核心是自定义SourceFunction去调用达梦日志挖掘接口然后打印结果。生产级代码会复杂得多但理解这个骨架后替换成你需要的解析逻辑就顺理成章了。public class DmLogCdcJob { public static void main(String[] args) throws Exception { // 1. 创建 Flink 执行环境并开启 checkpoint StreamExecutionEnvironment env StreamExecutionEnvironment.getExecutionEnvironment(); env.enableCheckpointing(10000, CheckpointingMode.EXACTLY_ONCE); // 2. 添加自定义达梦日志 Source DataStreamString cdcStream env.addSource(new DmLogMinerSource( jdbc:dm://192.168.1.100:5236, FLINKCDC, flinkcdc#2024)); // 3. 打印到控制台验证链路 cdcStream.print(); env.execute(dm-log-cdc-sync); } }public class DmLogMinerSource extends RichSourceFunctionString { private final String jdbcUrl; private final String username; private final String password; private volatile boolean running true; private Connection connection; private long lastLogSequence 0L; public DmLogMinerSource(String jdbcUrl, String username, String password) { this.jdbcUrl jdbcUrl; this.username username; this.password password; } Override public void run(SourceContextString ctx) throws Exception { Class.forName(dm.jdbc.driver.DmDriver); connection DriverManager.getConnection(jdbcUrl, username, password); connection.setAutoCommit(false); // 增量轮询日志每次从上次位点之后开始挖掘 while (running) { ListString changes mineLogSince(lastLogSequence); for (String change : changes) { // 每条变更记录输出为一个 JSON 字符串 ctx.collect(change); } Thread.sleep(1000); } } private ListString mineLogSince(long lsn) { // 调用 DBMS_LOGMNR.START_LOGMNR然后查 V$LOGMNR_CONTENTS // 根据 SAFE_LSN 字段判断是否为新日志拼装 JSON 返回 return new ArrayList(); } Override public void cancel() { running false; } }逻辑说明DmLogMinerSource的核心在run方法里。它维护了一个lastLogSequence变量作为日志位点每次循环从该位点之后开始挖掘日志挖掘出的每条变更以 JSON 字符串形式通过ctx.collect发往下游。Thread.sleep(1000)控制轮询频率为 1 秒一次你可以按延迟要求调成 100 毫秒但注意频繁启动日志挖掘会消耗数据库 CPU。这里有个关键点达梦日志挖掘不是纯增量接口它每次要指定一个日志范围。为避免重复读取必须在mineLogSince里过滤掉已消费的日志序列号这个过滤逻辑是整个链路正确性的基础。生产环境中不建议自己维护lsn位点而是把位点交给 Flink 的 checkpoint 机制管理后面第 4 章会展开。3.4 架构选型直连目标库还是先走 Kafka最小任务验证通过后下一步是决定同步架构。常见做法有两种Flink 任务直接写目标库或者 Flink 先写 Kafka再由下游消费。我一般推荐先走 Kafka理由有三个。第一解耦源库和目标库目标库不可用时不会反过来阻塞日志消费第二多下游复用一份变更数据数仓、搜索、缓存都能订阅第三Kafka 可以设较长保留期相当于给日志同步加了第二层缓冲源端归档被清理后还能从 Kafka 补数。如果同步的目标库也是达梦比如做达梦到达梦的数据分发那直连是可行的但要注意目标库的提交频率。Flink 的 sink 默认是攒批提交如果目标库开了很高的 redo 生成速率反而会加剧源和目标两边的日志压力。直连场景下我一般把 sink 的批量参数调到flushInterval3000ms、batchSize1000避免每一两条变更就提交一次事务。4. 参数调优三件事并行度、checkpoint 与日志保留策略4.1 必调参数表先记这三张表的配置FlinkCDC 接达梦后的性能瓶颈通常不在 Flink 本身而在日志解析效率和下游写入速度。下面是三个层面的关键参数按先后顺序调。参数所属组件默认值建议值作用env.enableCheckpointingFlink不开启10000ms决定故障恢复时的数据回溯范围env.setStateBackendFlinkmemoryrocksdb保存日志位点和状态防 OOMparallelismFlink Source11-2达梦日志挖掘不宜高并行parallelismFlink Sink14-8写入目标库的并发度DBMS_LOGMNR.CONTINUOUS达梦-1连续模式挖掘减少重复加载日志文件关于 Source 并行度很多人第一反应是并行度拉高来提升吞吐这是一个翻车点。达梦日志挖掘是按时间线串行推进的并行度超过 2 并不会加快日志读取反而因为多个 TaskManager 同时触发START_LOGMNR导致日志文件被重复加载实例 CPU 和 IO 双双飙升。我踩过的经验是Source 并行度固定为 1把吞吐压力放在下游 sink 并行度上数据库写入并行度 8 对达梦来说是比较安全的起点。如果单 Topic 消费有瓶颈再考虑按表名做 key 拆分区。4.2 checkpoint 和状态后端决定你能追回多少数据日志同步任务跑久了最怕的是 Flink 任务因为网络抖动或下游写失败被重启重启后从哪接着读。这里的关键就是 checkpoint。把 checkpoint 间隔设置为 10 秒意味着最坏情况会重复消费最近 10 秒的日志目标库可能出现少量重复数据如果目标库不幂等就需要 sink 侧做去重。状态后端我固定用 RocksDBMemoryStateBackend 在状态稍大时极容易 OOM。FlinkCDC 的位点信息虽然很小但如果同步任务里还缓存了表结构映射信息、以及每个表最近一段时间的操作记录用于排序状态量会涨得比预期快用 RocksDB 后checkpoint 和恢复速度都稳定很多。# flink-conf.yaml 关键配置 state.backend: rocksdb state.checkpoints.dir: hdfs:///flink-checkpoints execution.checkpointing.interval: 10s execution.checkpointing.tolerate-failed-checkpoints: true参数说明tolerate-failed-checkpoints允许连续几次 checkpoint 失败不立刻杀任务给日志挖掘的抖动留点余地但如果日志源真的读不到了这个参数只能短暂续命不能当作长期容错机制。实际生产里可以在监控上盯numberOfFailedCheckpoints这个指标连续 3 次失败就告警介入。4.3 延迟的真实来源日志解析不是最慢的环节很多人以为日志同步的延迟主要花在挖掘日志上实际上达梦的V$LOGMNR_CONTENTS查询速度跟增量日志量成正比在数据量不大时纯解析延迟在毫秒级。真实延迟的大头通常有三个。第一是目标库的写延迟。达梦这类传统数据库在并发写入时存在锁等待如果目标表和源表结构不完全一致每条记录还要做类型转换写入更慢。第二是 Flink sink 的攒批等待。为了减少目标库频繁提交sink 默认攒够一批才写入如果你设置的batchSize过大而实时变更流量又不大很多数据会在缓冲里等凑数秒级延迟就是这么来的。第三是 Kafka 生产端如果开了 acksall单条延迟会受 Broker 刷盘速度影响在日志同步链路里把 acks 设为 all 的同时也应该把 batch.size 调大用吞吐换延迟均衡。判断延迟瓶颈的方法是看 Flink Web UI 里每个算子之间的Records Received和Lag指标。Source 输出到 transform 的节奏是稳定的但 transform 到 sink 之间如果出现积压优先检查 sink 的目标库写入耗时和批量参数而不是怀疑日志挖掘。5. 达梦日志同步避坑 5 例从归档丢失到重复数据5.1 任务启动报“日志不存在”归档日志被清掉了现象Flink 任务重启后立刻报log file not found日志挖掘无法继续目标库开始出现数据延迟。原因达梦按SPACE_LIMIT自动清理旧的归档日志而 Flink 的 checkpoint 位点停留在更早的日志位置。任务停了两天这两天里归档目录转满旧日志被数据库自动删除新写的日志又和旧位点之间断了档。解决把归档目录的SPACE_LIMIT从 10GB 调大到 100GB同时清理周期改为保留 7 天以上。另一个预防措施是给 Flink 任务加“日志断点告警”当消费位点对应的归档文件已被删除时立刻发告警而不是等任务失败。恢复时如果目标库允许可以采用重建全量快照的方式补齐缺口只靠日志已经救不回来。5.2 同步账号权限不足日志挖出来全是空的现象任务不报错V$LOGMNR_CONTENTS查询结果没有内容控制台输出的数据量为 0。原因账号没有SELECT ANY TABLE权限达梦日志挖掘时把无权限的表记录直接过滤掉了但不会报权限错误反而给调试制造了很大迷惑。解决验证权限的方式很简单用同步账号单独查一张业务表能查说明表权限没问题再调DBMS_LOGMNR.START_LOGMNR后看 v$logmnr_contents 里有没有该表的操作记录。如果还没有检查是不是开了行级安全或列级加密这两种情况下日志默认不做物理还原。5.3 大事务导致延迟尖刺监控图上出现锯齿现象每天固定时间点延迟从秒级飙升到分钟级过后又自己恢复日志挖掘任务没有报错。原因达梦日志挖掘按事务提交顺序输出如果一个事务一次性更新了几十万行V$LOGMNR_CONTENTS会一次性返回超大结果集解析和序列化都要耗时后续事务全部排队。解决对超大事务无能为力只能接受尖刺但可以让尖刺不传导到下游。具体做法是把 Source 的输出拆成“事务头 数据块”数据块逐条发送到下游同时调大 Sink 的缓冲队列让尖刺被缓冲吸收掉。如果大事务频繁出现在特定时间点建议跟业务方确认是否能拆分成小批量提交这对达梦自身的日志压力和实例性能也有好处。5.4 目标库是达梦时模式名和大小写对不上导致写不进去现象源端是达梦目标也是达梦日志挖掘出的表名在目标库执行时提示模式不存在或者表找不到。原因达梦的模式名默认和用户名一致而且表名在创建时如果不加双引号会统一存成大写。日志挖掘返回的模式名和表名大小写语义在不同版本上表现不一致直接拼 SQL 时如果没做大小写归一化目标端就会找不到对象。解决在做 transform 层时对 schema 名和表名统一做toUpperCase()处理同时禁止在表结构里使用双引号包裹的小写表名。另外一点达梦的模式错误提示在很多场景下其实是连接串里的 schema 参数没配对比如 JDBC URL 里指定了schemaFLINKCDC但实际数据在MAIN模式下也会在写入时报模式错误。5.5 任务重启后目标库出现重复数据EXACTLY_ONCE 不是万能的现象Flink 任务因为网络抖动重启过重启后目标库出现了重复记录count 对不上。原因Flink 的EXACTLY_ONCE只保证 Flink 内部状态不丢真正要做到端到端精确一次目标库的写入必须支持幂等。直接向达梦执行INSERT的 sink天然没有幂等性重复消费时就会重复插入。解决目标库是达梦时在 sink 层把写入改成“先按主键查再决定 insert 或 update”或者利用目标表的唯一索引做主键冲突时更新。达梦支持MERGE INTO语句用它替代INSERT配合 checkpoint 的重复消费就能把重复数据挡在外面。代价是写入吞吐会下降约两成这是精确一次语义的合理成本。6. 更稳的姿势用一张校验表验证同步没丢数再谈优化日志同步上线前我建议先花半天时间做一次“对账”而不是直接接业务流量。对账思路很简单在源端选一张有自增主键和更新时间字段的表同步跑一段时间后分别从源端和目标端按天聚合COUNT、MAX(主键)、SUM(校验字段)做对比。下面这个 SQL 就是常用的对账脚本原型-- 源端统计 SELECT COUNT(*) AS row_cnt, MAX(ID) AS max_id, SUM(AMOUNT) AS sum_amount FROM SRC_ORDER WHERE TRADE_DATE 2024-11-20; -- 目标端统计表名、字段名经过 transform 映射 SELECT COUNT(*) AS row_cnt, MAX(ID) AS max_id, SUM(AMOUNT) AS sum_amount FROM DWS_ORDER WHERE TRADE_DATE 2024-11-20;两个查询结果逐一比对row_cnt对不上说明有丢数row_cnt一致但max_id不一致说明中间某条主键发生跳变sum_amount对不上说明存在重复或字段值被改写。这套方法不需要额外引入数据校验平台一条 SQL 就能完成对上线前的验收足够用。要是源表没有自增字段就用MIN(ROWID)和MAX(ROWID)配合COUNT做范围对比逻辑是一样的。延迟验证则在 Flink Web UI 上观测 Source 到 Sink 的Lag指标通常会稳定在两秒以内如果超过了五秒优先查目标库写入耗时和 Kafka 生产端积压。我个人的习惯是即使日志同步已经跑稳了每天凌晨仍会跑一次这个对账输出一份只有 OK 和 NG 的日报。同步这行干久了就会发现最贵的事不是调优而是在业务方告诉你数据不对时拿不出“到底哪条不对、什么时候开始不对”的实据。有一个对账机制兜底排障压力和踩坑次数都会小很多。希望帮到你。本文还有配套的精品资源点击获取