达梦数据库实时同步实战:基于FlinkCDC日志捕获的完整指南
简介基于FlinkCDC组件的达梦数据库日志实时同步方案包面向大数据开发工程师与实时计算学习者帮助解决从国产达梦数据库中捕获增删改数据并实时接入流处理框架的难题。压缩包内共包含315个文件以263个依赖库文件为主体另含Java源文件、编译后的类文件、XML与属性配置文件、结构化查询语句脚本以及说明文档资源总大小约为三百四十一MB。方案同时提供基于Java编程接口和基于结构化查询语言的两种同步方式包含FlinkDMCDC、DMFlinkSQL等核心类以及自定义反序列化实现可直接作为工程模板导入开发环境。目前已有七百五十九人浏览学习适合需要快速搭建达梦数据库实时同步管道的开发者通过完整的代码、配置和脚本减少重复排错工作。 做达梦数据库的实时同步最忌讳上来就写代码。我先把自己踩过的弯路说清楚再给你一条能落地的正路。我在做某个国产化替换项目时上游是达梦下游要接一套实时数仓。当时时间紧第一反应是直接用FlinkCDC一把梭。结果发现网上关于达梦的FlinkCDC资料少得可怜而且很多方案其实是“假CDC”的轮询模式延迟根本压不下来。后来把达梦的日志机制、FlinkCDC新版本的连接器选项、还有整条链路的容错方式都摸了一遍才算真正跑通。这篇就把完整思路、配置和排查记录整理出来给同样被达梦同步折腾的人一个参考。1. 为什么“基于日志”才是达梦实时同步的正解先说结论达梦的实时同步底层必须靠日志解析不能靠JDBC轮询。1.1 同步方式的本质区别我见过很多人把“定时SELECT”和“实时同步”混为一谈。实际上达梦同步有三条路差别非常大方案实现方式延迟对源库影响能否感知删除/DDLJDBC定时轮询定时执行SELECT用时间戳或自增ID找增量秒到分钟级取决于轮询频率高频繁查询业务表很难删除和结构调整基本无能为力日志解析解析达梦归档日志/联机日志中的事务记录毫秒到秒级低不碰业务表可以完整感知Insert/Update/Delete/DDL厂商工具/数据集成软件达梦自身或第三方ETL工具的CDC能力秒级较低一般可以达梦的JDBC驱动只提供了普通查询能力没有像MySQL binlog那样方便的回放接口。如果只靠SELECT * FROM T WHERE UPDATE_TIME ?拉增量两个致命问题一是频繁扫描业务表高峰期会把源库CPU打满二是物理删除的数据直接消失下游永远不知道。走日志解析本质上是让达梦自己把数据变更“吐”出来。你开启归档模式后所有事务都会按顺序写进联机日志和归档日志FlinkCDC这类工具读取并翻译这些日志就得到了完整、有序、低延迟的变更流。这跟MySQL binlog、Oracle Redo日志的同步原理是一个思路只是达梦的日志格式和读取接口跟它们不一样。1.2 为什么选FlinkCDC而不是自己写解析器自己写一个达梦日志解析器我劝你打消这个念头。达梦的日志内部结构复杂要处理事务边界、回滚段、字段类型映射、DDL变更没几个月搞不定。FlinkCDC的价值在于它已经把“连接数据库、读取日志、解析变更、输出Changelog流”这些脏活封装好了你要做的只剩下配置连接参数、定义表结构、写下游写入逻辑。而且FlinkCDC不是简简单单一个采集工具它背靠Flink的分布式计算能力。同步过程中可以做过滤、转换、多表合并、维表关联甚至把一张大表拆成多路下游这些能力是普通同步工具不具备的。所以选FlinkCDC做达梦日志同步不是一个临时凑合的方案而是真正能在生产环境长期跑的架构选择。2. 同步前的达梦数据库准备归档、权限、测试环境日志同步的第一步不是写Flink代码而是把达梦调成“可被日志采集”的状态。很多人在这一步就卡住了。2.1 开启归档模式必须做达梦默认不开启归档模式不开启就意味着日志很快被覆盖或循环写入FlinkCDC根本来不及读。我用disql操作时大致是这样-- 以SYSDBA登录disql ALTER DATABASE MOUNT; ALTER DATABASE ARCHIVELOG; ALTER DATABASE ADD ARCHIVELOG DEST/dmdata/arch,TYPElocal,FILE_SIZE256M,SPACE_LIMIT4096M; ALTER DATABASE OPEN;执行完可以用视图确认SELECT * FROM V$DM_ARCH_INI;如果看到归档状态是“有效”或类似状态说明已经开启。注意ALTER DATABASE MOUNT/OPEN这种操作会短暂切换数据库状态。如果是在生产库上操作务必先确认维护窗口或者让DBA配合。我自己的习惯是用一套Docker测试库把命令跑熟了再上生产。归档目录的空间也要算好。达梦的归档日志文件可以设置单个文件大小和空间上限。我一般设单文件256MB总上限根据业务写入量估算至少保留够FlinkCDC一个消费周期比如3天的量。否则日志被清理后同步链路只能重新全量初始化。2.2 账号权限准备日志解析需要能读到归档日志的账号。我用SYSDBA账号最省事但出于安全考虑生产环境建议单独建一个最小权限账号。达梦里需要保证账号具备连接数据库权限CREATE SESSION读取业务表的权限SELECT ANY TABLE或对应表的SELECT读取日志相关视图/包的权限这个一般DBA角色才有如果你用业务账号同步失败报“权限不足”或“无法读取日志”先别查Flink配置回数据库把角色权限补上。这不是Flink的问题。2.3 Docker快速搭一套达梦测试库很多时候不是不想测是手边没有达梦环境。达梦官方提供了Docker镜像我用的就是社区版单机镜像拉起来就是一套能用的达梦8docker run -d --name dm8 \ -p 5236:5236 \ -e SYSDBA_PWDYourPassword123 \ dm8_single:dm8_20230808端口默认5236用户名SYSDBA密码是启动时设置的那个。起来之后用达梦的disql或者DBeaver都能连。DBeaver里选达梦驱动填上地址端口就能进这比我一开始用命令行高效多了。这套环境跑通之后后面所有FlinkCDC配置都可以先在Docker里验证不会污染生产数据。3. 用FlinkCDC实现达梦日志实时同步的核心链路准备工作做完接下来进入正题怎么把达梦的日志变更流变成Flink能用的数据流。3.1 版本选择与依赖引入这里有个关键点FlinkCDC对达梦的官方支持是逐步完善的。我最初用的是FlinkCDC 2.x老版本它对达梦支持不友好折腾了很久。后来切到FlinkCDC 3.x发现社区已经将达梦纳入了连接器支持矩阵才把日志同步这条路走通。如果你用的是Maven工程大致引入这些依赖以你实际使用的稳定版本为准dependency groupIdorg.apache.flink/groupId artifactIdflink-connector-dm-cdc/artifactId version3.x.x/version /dependency dependency groupIdcom.dameng/groupId artifactIdDmJdbcDriver18/artifactId version8.1.x/version /dependency注意达梦的JDBC驱动需要你自己下载或从本地Maven仓库引入很多公共仓库里没有。别小看这一步驱动版本和FlinkCDC连接器不匹配会导致连不上或日志解析异常。3.2 Flink SQL方式最直接的同步写法如果你只是想尽快把表同步到下游用Flink SQL创建source表是最快的方式。大致是这样的结构CREATE TABLE dm_order ( id INT, order_no VARCHAR(64), user_id INT, amount DECIMAL(10, 2), status INT, create_time TIMESTAMP(3), PRIMARY KEY (id) NOT ENFORCED ) WITH ( connector dm, hostname 192.168.10.20, port 5236, username SYSDBA, password your_password, database-name DMDB, schema-name DMUSER, table-name T_ORDER, scan.startup.mode initial, scan.startup.timestamp-millis 0 );这里有几个SQL参数要特别留意。scan.startup.mode initial表示任务启动时先做一次全量快照再自动切换到日志增量模式。这是最常用的模式解决了“历史数据实时增量”衔接的问题。如果只想从某个时间点开始同步可以用timestamp-millis指定。PRIMARY KEY NOT ENFORCED不是摆设。日志模式下Flink需要知道主键来生成正确的UPDATE/DELETE事件如果表本身没主键这里至少要声明一个逻辑主键否则下游无法正确处理更新。建好source表之后再建sink表一条INSERT INTO就完成了同步任务INSERT INTO sync_order SELECT id, order_no, user_id, amount, status, create_time FROM dm_order;3.3 YAML Pipeline方式多表同步更省心如果你要同步几十张表一张张写SQL会累死。FlinkCDC 3.x的YAML Pipeline方式更适合多表场景。大致是这样的配置source: type: dm hostname: 192.168.10.20 port: 5236 username: SYSDBA password: your_password database-name: DMDB schema-name: DMUSER table-name: T_ORDER,T_USER,T_PAY startup-mode: initial sink: type: doris fenodes: 192.168.10.30:8030 username: root password: your_password database-name: dwd table-name: sync_order实际字段以你使用的FlinkCDC版本官方配置为准但思路是一致的source定义从哪里读sink定义往哪里写。pipeline方式的好处是处理多表、字段路由、schema变更时更有全局视角我生产环境更推荐这种方式。3.4 备选架构达梦日志先进KafkaFlink再从Kafka消费如果你的下游不止Flink一套系统或者你需要在同步前做多份分发那更稳的架构是“达梦日志 → Kafka → Flink”。这个方案里达梦日志的采集可以由达梦数据同步工具或基于Debezium扩展的采集组件完成把日志解析出来的变更记录以JSON或Debezium格式写入Kafka指定Topic。Flink侧只需要用kafkaconnector消费Topic再做清洗、关联、写入下游。整体链路达梦数据库 → 日志采集端 → Kafka Topic → Flink SQL/DataStream → 下游Doris/ClickHouse/MySQL等这个方案比Flink直接连达梦多了一层Kafka但解耦性更好。采集端挂了不影响Flink任务Flink重启也能从Kafka的offset续跑。如果对可用性要求高我强烈建议加上这层缓冲。毕竟Flink任务重启时如果达梦日志已经被清理那就要整表重刷了成本非常高。4. 实际同步过程中的问题排查与避坑清单这一节是我最想写的因为配置本身并不难真正折磨人的是那些“看起来没问题但就是跑不起来”的隐藏坑。4.1 常见问题速查表现象根因解决办法任务启动报“权限不足”同步账号缺少读取日志或表的权限用SYSDBA账号或给账号授予DBA角色和对应表SELECT权限全量阶段正常增量阶段无数据归档模式未正确开启或日志已被清理检查V$DM_ARCH_INI确认归档目录空间足够时间字段相差8小时时区配置不一致Flink和达梦侧统一设置server-time-zone指定为Asia/Shanghai下游主键冲突/更新变插入Changelog丢失主键约束source表声明PRIMARY KEY NOT ENFORCED下游表也必须有主键任务卡住不动Checkpoint一直失败下游写入较慢或大事务日志堆积降低单批写入量调大checkpoint间隔给下游加索引同步出的DECIMAL字段精度丢失JDBC驱动或连接器类型映射问题检查达梦驱动版本必要时在SQL中显式CAST为DECIMAL(p,s)DDL变更后任务异常表结构变化未同步到下游下游表提前加好列任务重启前先核对schema4.2 几个容易忽视的细节物理删除与日志保留日志模式虽然能捕获DELETE但前提是归档日志还保留着对应的事务记录。如果达梦的归档日志被清理得太快FlinkCDC还没来得及消费删除事件就丢了。这个问题的唯一解法是给归档目录留足空间或者把FlinkCDC的checkpoint间隔调小尽量缩短消费延迟。全量阶段的大表性能initial模式下首次全量扫描会消耗较多IO。我同步过一张2亿行的订单表全量阶段跑了将近40分钟。这期间源库的IO有明显升高。如果源库是生产核心库建议先在业务低峰期启动同步任务或者先用数据导出工具把历史数据灌到下游再用timestamp-millis从某个时间点开始增量。达梦没有主键的表日志模式对无主键表的处理很不友好。Flink无法判断一条UPDATE对应哪一行。我遇到这种表一般会跟业务方沟通加一个唯一字段或者在达梦侧先做一层物化视图/中间表让同步对象有明确主键。时区问题达梦数据库的SESSIONTIMEZONE和Flink任务的TimeZone不一致时时间字段会偏移。我在同步交易流水时遇到过下游数据比业务时间早8小时的情况排查了半天发现是JDBC连接的时区参数没设。统一设置为Asia/Shanghai能规避绝大多数时间乱掉的问题。下游写入的幂等性日志模式下Flink可能会因为Checkpoint或事务回放对下游重复写入同一行数据。下游存储必须具备幂等更新能力比如Doris的Unique模型、MySQL的主键UPSERT。如果下游是普通的append-only Kafka Topic那你看到的可能就是重复消息这不是Flink的bug而是架构设计上要提前考虑的点。4.3 我的调试顺序建议如果你按照上面的步骤还是跑不通我建议按这个顺序排查不要东一下西一下先确认达梦归档模式是否开启用disql查看归档配置。用DBeaver或disql手工执行几条SQL确认同步账号能正常连库、能SELECT目标表。启动FlinkCDC任务先观察日志里是否出现“Connected to ... Database”之类连接成功的信息。手工在达梦里对目标表做一次INSERT和UPDATE观察Flink日志是否打印变更记录。最后才检查下游的写入是否成功、字段映射是否正确。这个顺序是从“源是否OK”到“连接是否OK”再到“变更是否被解析”再到“写入是否OK”的逐层排查。直接跳到下游找问题往往会浪费很多时间。5. 一个我正在用的生产链路经验最后分享一个我现在比较推荐的生产配置调整思路这个是在多次压测和故障复盘之后总结出来的。FlinkCDC消费达梦日志时Checkpoint间隔不要只图快。我一开始把Checkpoint间隔设成1秒想着实时性拉满结果频繁触发状态快照反而把任务拖慢而且达梦日志读取位置频繁记录一旦网络抖动恢复时容易重复消费。后来调整到5秒一个Checkpoint下游延迟从原来的几百毫秒增加到了1秒左右但整体稳定性好了很多。关键阈值参考可以这样设execution: checkpointing: interval: 5s mode: EXACTLY_ONCE timeout: 60s如果业务能接受秒级延迟5秒检查点是非常稳的配置如果必须亚秒级延迟那就要接受频繁快照带来的性能开销并对下游写入做更精细的限流。再一个体会是监控。日志同步链路最怕“看起来正常实际上日志早就断了”。我在生产环境中给同步任务额外加了告警规则如果连续10分钟消费的变更记录数为0但源库同时段的写入量明确不为0说明同步链路可能已经断掉需要立即人工介入。别等下游报表出现数据缺口再去查那时候补数据已经很难受了。这个项目整体跑下来我的核心感受是达梦的日志实时同步能不能做成关键在于前期对达梦日志机制和权限模型的准备是否到位。FlinkCDC本身只是中间一环真正决定成败的反而是环境配置、版本选型和容错设计。希望这篇记录能让你少走一些弯路。本文还有配套的精品资源点击获取