ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

达梦数据库实时同步:Flink CDC日志接入与五大避坑实践

达梦数据库实时同步:Flink CDC日志接入与五大避坑实践 简介FlinkCDC与达梦数据库结合的实时同步方案资料面向需要构建实时数仓、数据同步及事件驱动应用的Java开发者或数据工程师。该压缩包共315个文件约341.71MB以263个jar依赖库为核心辅以xml配置、class编译产物、java源码、sql脚本等覆盖Flink作业开发、连接器配置与SQL同步两种落地方式。已有764人学习下载。资料内含可直接运行的FlinkDMCDC示例、自定义反序列化实现、基于Flink SQL的同步逻辑及配套统计工具类并附有工程配置文件与开发文档可帮助读者快速理解达梦数据库日志捕获原理并基于实际代码改造适配自身业务场景。资源适合具备一定Flink基础、希望低成本上手国产数据库实时同步的中高级开发者。1. 做实时同步最怕的不是延迟是方案一开始就选错最近好几个团队在问同一件事FlinkCDC 接达梦数据库基于日志做实时同步到底能不能落地。我的结论是能但绝大多数人第一步就找错了入口——Flink CDC 官方连接器列表里根本没有达梦。这不是说标题里的方案不成立而是 Flink CDC 本身不是单个连接器它是一套“日志接入 流式计算 结果分发”的处理框架达梦侧要靠官方日志采集组件把 redo/归档日志解析成结构化事件再喂给 Flink。真正的实时同步链路是“达梦日志 → 消息管道 → Flink 消费计算”而不是把达梦直接塞进某个现成 Source。这套思路适合实时大屏、报表宽表、数据仓库增量入仓这些场景也适合刚从 MySQL 思维转过来、准备认真接达梦的团队。下面按我实际部署过的链路展开先说日志原理和环境准备再给最小可跑任务、关键参数最后列五个必踩的坑。2. 达梦的日志从哪里来先看清 redo、归档与同步日志的边界2.1 达梦日志体系的三层认知redo、归档日志和逻辑日志做基于日志的同步第一个要扭转的认知是不要拿 MySQL 的 binlog 思维去套达梦。达梦的日志体系更像 Oracle核心是重做日志redo log和归档日志archive log。redo log 记录的是物理变更比如数据页上哪个位置被改成了什么值它循环写入写满就切换归档日志则是 redo 切换后保留下来的副本用来做恢复和追日志。很多从 MySQL 转过来的同事会习惯性问“binlog 日志可以删除吗、删了影响同步吗”在达梦这里你要关心的是“归档日志保留多久、目录会不会写满”这决定了你的同步链路能回溯多少。再往上一层才是“逻辑日志”。所谓基于日志的实时同步本质上是让日志采集组件去读 redo/归档把物理变更翻译成“哪张表哪一行在什么时间被 INSERT/UPDATE/DELETE”的逻辑事件。这个翻译过程非常消耗资源和心思因为要处理事务边界、回滚段、DDL 变更、类型映射所以基本不会有人自己从零写解析器。这也是后面选型时最重要的判断依据谁来做日志翻译决定了这个方案稳不稳。还有一类常见的误区是把“慢查询日志”或应用日志当作同步数据源。慢查询日志只能帮你事后排查 SQL应用日志是业务自己打的点它们都不是事务日志给不了准确的增删改前后镜像。我在项目里见过有人为了赶工期直接轮询一张“最后修改时间”字段的表来做增量那叫伪同步不是基于日志的 CDC事务内多行变更、物理删除都抓不住线上跑两个月必然对不上账。2.2 三条实现路径与选型对比别一上来就想造一个“原生连接器”既然 Flink CDC 官方没有达梦连接器从业者一般会走下面三条路。我把它们的差异直接放在一张表里你照着业务约束选就行。实现路径是否真基于日志开发量维护成本常见度达梦官方日志采集组件 → Kafka → Flink CDC是低主要是配置和格式约定中依赖官方组件版本最高生产环境首选自研 Flink Source 直连达梦日志解析接口是高需要熟悉内部日志视图与类型映射高达梦版本升级可能不兼容低只适合有专门团队的大厂应用双写业务代码同时写达梦和消息队列否中侵入业务高漏写一处就丢数据低仅试点我一般会直接选第一条。理由很实际达梦的 redo 日志格式和解析接口在不同版本之间有差异很多细节是黑匣子自己写 Source 意味着每个版本都要跟着适配而官方日志采集组件已经处理了事务拆分、归档切换、类型映射这些脏活我只需要把事件格式约定好剩下的交给 Flink。第二条路听起来很“硬核”但实际项目里你会在类型映射上耗掉大量时间比如 DM 的 NVARCHAR2、CLOB 在不同日志模式下读出来的前镜像可能带格式符这些坑没有文档可查只能拿真实数据一遍遍试。第三条路不做评价它连日志都没碰不在本文讨论范围。选第一条路之后架构就清晰了日志采集端负责把达梦日志翻译成 JSON 事件写入 KafkaFlink CDC 侧负责消费 Kafka、做清洗转换、再分发到目标端。你不需要再纠结“Flink CDC 怎么直连达梦”那不是重点重点是事件格式和消费位点怎么对齐。2.3 上线前的环境准备归档、权限、网络与时钟环境准备直接影响日志采集端能不能启动。第一件事是确认目标达梦库已经开启归档并且归档目录有独立磁盘。常见做法是在 DIsql 里以 SYSDBA 执行下面的命令注意不同版本命令写法有差异以你当前环境的《达梦数据库管理员手册》为准-- 设置本地归档目录目录要先创建好且不能和数据库文件放同一个磁盘 ALTER DATABASE ADD ARCHIVELOG DEST/dm8/arch/TYPELOCAL; -- 开启归档模式需要重启数据库实例生效 ALTER DATABASE ARCHIVELOG;代码的逻辑说明第一句是将归档日志输出到/dm8/arch目录TYPELOCAL表示本地归档这是最常用的方式第二句是把实例切换到归档模式。很多同步任务启动失败、日志采集端收不到增量查到最后就是归档没开或者开完没重启实例。改完归档后建议再查一下实例状态确认确实切换成功SELECT NAME, STATUS$, ARCH_MODE FROM V$DATABASE;注意不同版本的V$DATABASE视图字段名不完全一样有的版本叫ARCH_MODE有的叫ARCHIVELOG。如果执行报“列不存在”就去 DIsql 里执行DESC V$DATABASE看实际字段别照抄。归档配置完成后还要检查两件事给日志采集端单独建一个同步账号不要直接用 SYSDBA 跑业务任务权限按“能读日志、能查目标表”的最小集授另外确认所有节点时钟一致最好有 NTP 同步。时钟不一致会让事件时间戳错乱后面做延迟监控时你会看到延迟忽正忽负非常难受。3. 跑通第一条同步链路从达梦日志到 Flink SQL 源表3.1 链路设计日志采集端写入 KafkaFlink 只做消费与计算我建议的最小可运行链路是达梦归档日志 → 官方日志采集组件 → Kafka Topic → Flink SQL。日志采集端在达梦服务器上跑它把翻译后的逻辑事件写到 KafkaFlink 侧不关心达梦日志格式只消费 Kafka 里的 JSON。这样的好处是解耦日志采集端出问题不影响 Flink 作业Flink 重启也不需要回放达梦日志直接从 Kafka 位点恢复即可。日志采集端默认输出的 JSON 事件格式一般长这样具体字段由采集端的版本和配置决定但核心信息逃不开这几个{ op: UPDATE, source_table: SCOTT.USER_TAB, row_seq: 99123456, event_time: 2025-06-01T14:23:45.123, before: {id: 1001, name: 张三}, after: {id: 1001, name: 李四} }这里op是操作类型常见值是 INSERT、UPDATE、DELETErow_seq是日志序列号用来判断事件先后和做对账event_time是事务提交时间一般用 ISO-8601 格式带毫秒before和after是变更前后的完整行镜像。特别提醒如果采集端配置里把字段输出成了大写比如ID: 1001那后面 Flink 建表时字段名必须跟着大写走否则解析出来全是 NULL。这个细节在第 4 章和第 5 章还会反复踩到。3.2 在 Flink SQL Client 里建源表与结果表拿到 Kafka 里的样例 JSON 后就可以在 Flink SQL Client 里建源表了。下面的建表语句把 Kafka 的 JSON 事件映射成 Flink 表before_row和after_row用 ROW 类型承载整行镜像CREATE TABLE dm_log_event ( op STRING, source_table STRING, row_seq BIGINT, event_time TIMESTAMP(3), before_row ROWid INT, name STRING, after_row ROWid INT, name STRING, WATERMARK FOR event_time AS event_time - INTERVAL 5 SECOND ) WITH ( connector kafka, topic user_db.dm_log, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id flink-cdc-dm-group, scan.startup.mode earliest-offset, format json, json.ignore-parse-errors true );参数说明scan.startup.mode设为earliest-offset是为了第一次跑能从头看到完整数据流联调通过后可以改成group-offsetsjson.ignore-parse-errors建议打开因为达梦某些类型在极端情况下会产出异常字符串宁可丢一条坏消息也不要让整个作业卡死WATERMARK是给event_time定义乱序容忍度如果你对数据顺序不敏感可以删掉这行用处理时间就行。这里ROWid INT, name STRING的字段名必须和 JSON 里before、after的子字段完全一致大小写敏感不一致查出来的值全是 NULL。然后建结果表。这里我特别说明一下如果你的目标端是另一个达梦库而且目标表有主键、需要精确更新和删除那官方 JDBC Sink 的 UPSERT 语义需要在达梦上单独验证不能默认它一定会生成 MERGE。为了先跑通链路我通常先建一张追加写的结果表把转换后的明细落进去CREATE TABLE sync_result ( id INT, action_type STRING, sync_time TIMESTAMP(3), new_name STRING ) WITH ( connector jdbc, url jdbc:dm://192.168.10.20:5236, username sync_app, password ********, table-name SYNC_RESULT );这里url用的是达梦 JDBC 标准写法端口默认是 5236sync_app账号需要有对SYNC_RESULT表的 INSERT 权限。JDBC Sink 的连接器会自动加载驱动但前提是你把达梦驱动 jar 放到了 Flink 的lib目录下而且每个 TaskManager 节点都要有这个问题在第 5 章会详细展开。3.3 提交任务并核对日志事件建表完成后写一条最简单的转换 SQL把日志事件里的关键字段抽出来写入结果表。我用COALESCE处理 UPDATE 和 DELETE 的差异UPDATE 取 after 行DELETE 取 before 行INSERT INTO sync_result SELECT COALESCE(after_row.id, before_row.id) AS id, op AS action_type, event_time AS sync_time, COALESCE(after_row.name, before_row.name) AS new_name FROM dm_log_event;这段 SQL 的逻辑是读 Kafka 里的每一条日志事件把操作类型、事务提交时间、行主键和变更后的名字投影出来。COALESCE在这里很关键它处理了 DELETE 事件没有after_row、INSERT 事件没有before_row的情况。第一次联调不建议直接跑生产目标表先落一张调试用的结果表确认数据能持续写入再改目标。SQL 文件准备好之后提交方式用 SQL Client 最直接$FLINK_HOME/bin/sql-client.sh \ -f /opt/flink-jobs/dm_sync.sql这段命令会把dm_sync.sql里的建表和 INSERT 语句提交到 Flink 集群。联调阶段你可以在集群 Web UI 上看到作业的 Records Received 和 Records Sent 指标在涨如果想确认 Kafka 里确实有事件直接在 Kafka 机器上消费一下kafka-console-consumer.sh --bootstrap-server kafka-1:9092 \ --topic user_db.dm_log \ --from-beginning --max-messages 3能打印出三条 JSON 事件说明源端链路已经通了。此时去查目标库SYNC_RESULT表如果能看到对应行这条基于日志的实时同步链路就算真正跑通了。生产环境我一般不用 SQL Client而是把 SQL 转成 Flink CDC 的 YAML Pipeline 或者打成 JAR 提交但联调阶段用 SQL Client 验证问题最直观改起来也最快。4. 同步任务要稳先调好这四个参数并行度、checkpoint、时区与格式4.1 并行度设计Kafka 分区数、Flink 并行度与源端采集线程的关系很多人的第一反应是把 Flink 并行度调大觉得并行度越大吞吐越高。在达梦日志同步链路里这个直觉会翻车。并行度的上限不取决于 Flink 想跑多快而是取决于 Kafka 分区数和日志采集端的投递能力。如果 Kafka Topic 只有 6 个分区那 Flink 源端并发最多也就 6超过分区的并行度等于空转反过来如果 Kafka 有 6 个分区Flink 源端并行度只有 2那两个空闲分区会造成数据积压和乱序。我常用的起步设置是Flink 源端并行度等于 Kafka 分区数下游 JDBC Sink 并行度设为目标库写入压力的一半也就是不要超过源端并行度。比如 Topic 分了 6 个区源端并行度设 6Sink 并行度设 3。原因是 JDBC Sink 每个并发会占用一个数据库连接6 个并发同时写达梦加上达梦本身的会话开销很容易把目标库的连接数打满。在 SQL Client 里的设置方式如下注意新版在 Key 上加引号旧版不加SET parallelism.default 6; SET pipeline.operator-chaining true;operator-chaining打开后相邻的 Map、Filter 算子会合并到同一个线程里减少线程切换和序列化开销。如果你要精确控制 Sink 并发可以在建表时给 Sink 单独指定sink.parallelism 3源表指定scan.parallelism 6而不是只依赖全局默认值。4.2 checkpoint 与恢复点一下恢复按钮之前先想清楚位点Flink 作业跑起来简单真正让团队头疼的是重启以后数据对不对。基于日志的同步链路checkpoint 参数直接影响数据准确度我的起点配置是SET execution.checkpointing.interval 10s; SET execution.checkpointing.timeout 2min; SET execution.checkpointing.min-pause 5s; SET execution.checkpointing.max-concurrent-checkpoints 1; SET state.backend.type rocksdb;参数说明interval 10s意味着每 10 秒做一次快照太短会让 Kafka 和状态后端压力大太长会导致重启后回放的数据量大timeout 2min是单次 checkpoint 最长耗时超时说明下游写入有瓶颈min-pause 5s是两次 checkpoint 的最小间隔防止连续做快照把 CPU 打满max-concurrent-checkpoints必须为 1并发 checkpoint 会对 Kafka 位点和状态的一致性带来额外复杂度日常场景没必要开。rocksdb状态后端适合大状态的任务如果你的同步任务只做简单过滤投影用默认的 HashMap 状态后端也可以但 RocksDB 对内存更友好不容易因为状态增长把 TaskManager 搞 OOM。恢复作业时优先从 Savepoint 或最近一次 Checkpoint 恢复而不是从头消费 Kafka。命令行是flink run-application -s savepoint路径但这里有个关键细节恢复前确认作业拓扑没变过如果改了表名、加了字段或改了算子 ID恢复时会报状态不匹配。不要用--allowNonRestoredState去绕过那等于把状态对账的责任甩给了自己。第 5 章第 5 节会讲一个因为恢复方式不对导致丢数据的真实案例。4.3 时间、小数与大小写三种最容易翻车的格式细节第一是时区。Flink 的 JSON Format 解析带时区的时间戳时默认按 UTC 处理。如果你的 Flink 集群在本地时区事件时间显示出来会比数据库时间少 8 小时。解决方式是在 SQL Client 里设置集群本地时区SET table.local-time-zone Asia/Shanghai;第二是 DECIMAL 和浮点数精度。达梦的 DECIMAL(38,10) 这类高精度数值经过 JSON 序列化再解析如果 Flink 侧字段声明成 DOUBLE会有精度丢失更稳妥的做法是把小数位较多的字段在 JSON 里输出成字符串Flink 侧先按 STRING 接住在计算层再决定是转 DECIMAL 还是保留字符串透传。日志采集端一般有“数值类型转字符串”的开关能开就开。第三是大小写。达梦的默认行为是未加双引号的标识符统一按大写处理而日志采集端输出的 JSON 字段名常常是建表时的原始大小写。如果源表字段名是小写采集端也输出小写但你在 Flink 里写ID那映射就会变成 NULL。我习惯在 Flink SQL 里严格按 JSON 里的实际大小写写字段名并给目标表字段加双引号确保达梦不改变大小写语义INSERT INTO sync_result (id, action_type, sync_time) VALUES (1001, INSERT, TIMESTAMP 2025-06-01 14:23:45.123);这条 SQL 里的id、action_type都是按建表时定义的小写写的Flink 不会强行转换但目标库里的列如果实际是大写存储就需要在 JDBC Sink 的表名或字段映射里做对应处理。具体的模式名映射坑下一章会单独讲。5. 达梦日志同步避坑实录五条最常见的翻车现场5.1 归档没开或目录写满日志采集端一直看不到增量现象日志采集端启动后不报错但 Kafka 里就是没有新事件Flink 作业状态是 RUNNING数据延迟却越来越大。排查时发现源库的 redo 文件一直在切换但采集端日志显示“找不到可解析的归档日志”。原因达梦实例没有开启归档模式redo log 循环覆盖后无法追溯采集端拿不到完整的日志序列或者归档目录所在的磁盘写满了达梦实例直接暂停写入采集端锁死等待。解决按 2.3 节的命令开启归档并重启实例同时把归档目录放到独立磁盘容量按“每天日志量保留 3 天”估算。再有就是清理归档时不要直接rm /dm8/arch/*那是血泪教训——归档日志可能还在被采集端解析手删会导致解析中断和日志空洞应该用达梦自带的归档管理功能或系统函数清理并保留至少一个完整切换周期。5.2 驱动不匹配任务启动时报驱动类或方法不存在现象Flink 作业启动后报NoClassDefFoundError或NoSuchMethodError栈信息指向达梦 JDBC 驱动用 Navicat 或者 IDEA 连同一个达梦库都能连上说明库本身没问题。原因Flink 的 lib 目录里存在多个版本的达梦驱动 jar常见的是从达梦数据库安装目录拷出来的驱动和 Maven 仓库里的驱动版本不一致而 Flink 各个节点加载的类路径不同导致部分 TaskManager 加载到旧版驱动。另一个原因是某些团队的 Flink lib 里塞了全家桶驱动驱动类冲突。解决统一驱动版本只保留与达梦服务器小版本一致的驱动 jar并且确保它在每个 TaskManager 节点的FLINK_HOME/lib下都有一份不要只放在 JobManager 节点。如果是从某个下载渠道拿到的驱动注意看 jar 包内META-INF/MANIFEST.MF的版本信息和服务器版本匹配后再部署。5.3 模式名映射错乱能连库但同步后找不到表现象日志采集端正常输出SCOTT.USER_TABFlink 作业也不报错但下游目标库收到数据后写入失败报“模式不存在”或“无效的模式名”更隐蔽的情况是同步过来的数据有值但落到了错误的模式或表里。原因达梦数据库里模式和用户名强绑定默认情况下一个用户对应一个同名模式比如用户SCOTT默认访问模式SCOTT。日志采集端和 Flink 侧对表名的处理方式不同采集端输出可能是大写模式名也可能是小写而 Flink 里建表时如果写错了大小写到了 JDBC Sink 就会把SCOTT.USER_TAB解析成另一个模式导致“模式错误”。Navicat 能连上是因为它做了图形化适配不代表 JDBC 层面的 schema 映射一致。解决在日志采集端配置里固定 schema 大小写推荐统一用小写Flink 侧的source_table过滤条件写死成采集端实际输出的大小写不要靠人眼猜。然后在达梦里验证一下SELECT * FROM 你的模式名.USER_TAB是否能查到数据能查到说明 JDBC 连接串里的 schema 是对的再对齐 Flink 这边的写法。5.4 大事务把延迟从秒级拖成小时级现象业务侧跑了一个批量 UPDATE影响到几万行结果整个 Kafka Topic 的消费延迟从几秒涨到几十分钟其他表的同步也一起卡住Flink 作业看起来还活着但数据就是出不来。原因日志采集端默认把同一个事务里所有变更打包成一条大消息投递到 Kafka 单分区下游 Flink 必须顺序处理这一个分区几万行的变更塞进一条消息Flink 单并发解析 JSON 就要花很长时间后面的消息全堵住。解决分两层处理。第一层是从日志采集端配置“单个事务拆分阈值”把超过比如 5000 行的大事务按批拆成多条消息避免一条消息撑死一个消费者第二层是 Kafka Topic 的分区数要足够至少要大于等于 Flink 源端并行度让不同表的日志能散到不同分区并行消费。如果业务场景允许按表拆分 Topic 是更好的方式这样一张大表刷数不会堵住其他表的实时链路。5.5 重启作业后“丢尾”checkpoint 恢复了但记录少了现象Flink 作业因为发布或异常重启从 checkpoint 恢复后状态显示正常Kafka lag 也不大但和对账表一比少了一部分数据。原因日志采集端维护的日志位点就是我们前面事件里的row_seq和 Kafka 的 offset 不是同一个参照系。Flink checkpoint 恢复的是 Kafka offset而日志采集端如果按自己的位点回放它可能在重启时重新投递了部分消息也可能跳过了 Flink 还没消费的尾部消息。两个位点对不齐就出现“恢复成功但数据不一致”的诡异现象。解决让日志采集端把row_seq写进事件体里Flink 侧在启动后先做一个对账查询统计每个表最大row_seq和源库当前日志序列号的差距。如果差距在合理范围说明位点对齐如果对不上就要检查采集端的投递语义是不是“写入 Kafka 即返回”以及 Flink 的消费位点有没有在重启时被重置。我在生产上还会给 Flink 作业配一个延迟监控表每 5 分钟统计一次最新事件的event_time与当前时间的差值超过阈值就告警这样至少不会等业务方发现数据不对了才知道出了问题。6. 证明链路是“真实时”端到端延迟验证的三种做法6.1 三步验证造数、看消息、对比时间很多团队把链路跑通就以为完事了直到业务方问“到底实时到秒级还是分钟级”才答不上来。我的习惯是每套同步链路都做一次端到端延迟验证。先在源库插入一条带当前时间的探测数据INSERT INTO USER_TAB(ID, NAME, UPDATE_TIME) VALUES (999, probe, SYSTIMESTAMP);然后在 Kafka 端取这条数据到达的时间和源库插入时间对比kafka-console-consumer.sh --bootstrap-server kafka-1:9092 \ --topic user_db.dm_log \ --from-latest --timeout-ms 10000 | jq -r .event_time | tail -1 date %Y-%m-%d %H:%M:%S如果两条时间差在 10 秒内说明链路健康如果差到分钟级就要开始排查是日志采集端投递慢还是 Flink 源端并行度不够。这里要注意的是event_time是达梦事务提交时间不是采集端解析时间所以它衡量的是“事务提交到 Kafka”的传输延迟不包含达梦内部提交前的排队时间。6.2 用 Kafka 消息时间戳持续观测延迟生产环境不能靠人肉跑脚本我一般会在 Flink 建表时把 Kafka 消息时间戳暴露出来做实时延迟监控。Kafka 的timestamp元数据是消息写入 broker 时由服务端分配的或由生产者指定不容易被业务数据造假比事件体里的字段更适合做链路健康度判断CREATE TABLE dm_log_event ( op STRING, source_table STRING, event_time TIMESTAMP(3), kafka_ts TIMESTAMP_LTZ(3) METADATA FROM timestamp, before_row ROWid INT, name STRING, after_row ROWid INT, name STRING ) WITH ( connector kafka, topic user_db.dm_log, properties.bootstrap.servers kafka-1:9092,kafka-2:9092, properties.group.id flink-cdc-dm-group, scan.startup.mode latest-offset, format json );然后写一个简单的查询把kafka_ts和event_time的差值算出来超过阈值就通过告警通道通知值班人。kafka_ts使用的是 TIMESTAMP_LTZ 类型这是 Flink 处理带时区时间戳的标准类型显示时按集群本地时区换算。这样每条消息的“数据库提交时间”和“到达 Kafka 时间”都在同一张表里延迟是透明的。6.3 一个我保留到现在的检查习惯做达梦日志同步这几年我最大的变化是不再依赖“看日志觉得没问题”来判断链路健康。日志只能证明作业没挂不能证明数据没丢、没延迟。我现在每次改完采集端配置或 Flink SQL固定做三件事第一跑一遍造数脚本确认探测数据能穿透到最终结果表第二看 Flink Web UI 上这个作业的收到记录数和发出记录数是否持续增长第三对比 Kafka 最新消息的event_time和当前时间记录到一张运维笔记里。这套动作做多了哪些配置改动会影响延迟、哪些不会心里就有数了。希望帮到你。本文还有配套的精品资源点击获取
返回列表