☰
FlinkCDC 实时同步达梦数据库:日志级增量采集与 Kafka 链路实践
2026/10/10 21:43:16 网站建设 项目流程

简介:本资源面向大数据开发工程师与实时数仓建设者,聚焦FlinkCDC与达梦数据库的日志级实时同步方案,帮助解决国产数据库变更数据捕获与下游流处理系统对接的问题。包内共315个文件,以263个jar依赖包为核心,辅以xml配置、java源码、class编译文件、properties参数文件及sql脚本等,覆盖连接器依赖、作业配置与示例代码,压缩包约341.71MB,结构完整便于直接导入工程调试。已有765人学习下载,说明该方案在国产化替代与实时同步场景中具备一定参考价值。读者可从中获取达梦CDC连接器的集成依赖、Java与SQL两种同步方式的实现示例、数据库连接参数配置模板以及作业启动所需的关键组件,适合用于搭建实时报表、数据监控与事件驱动应用的数据管道,也可作为排查同步延迟与日志解析问题的参考素材。

1. FlinkCDC 接达梦:为什么日志级实时同步值得折腾

凌晨两点被电话叫醒,业务方说报表数据比生产库慢了四个小时,定时抽数任务卡在某个大表上跑不动。这种场景做数据集成的人多少都遇到过:传统批量抽取在数据量涨到千万级之后,窗口越拉越长,延迟从分钟级退化到小时级,还拖累源库。FlinkCDC 达梦数据库 基于日志实时同步这套方案,解决的正是这个痛点——不再靠轮询比对,而是直接读达梦的日志变更,把增量数据以秒级延迟送进下游。达梦作为国产数据库里装机量靠前的一款,很多政企、金融、能源系统在用它承载核心业务,而围绕它的实时同步资料却远不如 MySQL 生态丰富,这也是我把这套链路完整跑一遍的原因。这篇笔记面向已经会用 Flink、手上有一台能连的达梦实例、想把增量数据实时搬到 Kafka 或数仓的工程师,从原理选型讲到参数配置和踩坑,尽量让你照着能复现。

2. 达梦日志同步的底层逻辑与 FlinkCDC 的接入方式

2.1 达梦靠什么记录变更:归档日志与逻辑日志

达梦数据库的变更记录机制和 Oracle 有相似之处,核心是重做日志(REDO)加归档日志(ARCHIVE)。数据库运行时,所有事务的修改先写联机重做日志,日志切换后由归档进程把写满的日志文件转存到归档目录。要做实时同步,前提就是让达梦处于归档模式,否则日志被覆盖,增量就断了。

达梦还提供逻辑日志(Logic Log),通过配置可以输出更贴近行级变更的解析结果。相比直接解析物理 REDO,逻辑日志对同步工具更友好,字段级的 before/after 镜像更清晰。实际落地时,常见做法是开启归档并配置逻辑日志,让 CDC 组件去消费。

这里有个容易混淆的点:达梦的日志模式和 SQL Server 的「大容量日志」不是一回事。达梦归档模式只有开和关,开了才能做日志级同步,关了就只能走触发器或时间戳轮询。所以第一步永远是确认归档状态。

-- 查询达梦数据库归档模式状态 SELECT ARCH_MODE FROM V$DATABASE; -- 返回 Y 表示已开启归档,N 表示未开启 -- 查询当前归档日志配置路径 SELECT * FROM V$DM_ARCH_INI;

第一条语句查归档开关,ARCH_MODE为Y才能继续。第二条查归档配置,重点看ARCH_DEST归档目标路径和ARCH_TYPE,ARCH_TYPE为LOCAL时归档在本机,为REMOTE时走远程。如果归档没开,需要用ALTER DATABASE MOUNT后ALTER DATABASE ARCHIVELOG再OPEN,这一步会短暂停库,务必在维护窗口做。

2.2 FlinkCDC 怎么接达梦:连接器选型与版本匹配

FlinkCDC 官方连接器覆盖 MySQL、PostgreSQL、Oracle、SQL Server 等,达梦并不在默认列表里。所以接达梦有两条路:一是用达梦官方或社区提供的 CDC 工具把变更推到 Kafka,Flink 再从 Kafka 消费;二是基于 FlinkCDC 的增量快照框架自己实现达梦的 Source。前者落地快,后者可控性强但工作量大。

我一般推荐第一条路,原因是达梦生态里已经有成熟的日志捕获组件,它们对达梦内部日志格式的理解比我们自己啃文档要深。典型链路是:达梦归档日志 → 达梦 CDC 捕获组件 → Kafka Topic → Flink 消费 → 下游存储。Flink 这一侧只负责流处理,不直接碰达梦日志,职责清晰,出问题也好定位。

版本匹配上要留意:Flink 1.13 之后 CDC 连接器 API 有调整,Flink 1.17 以上对 Kafka 连接器的兼容性更好。达梦这边,DM8 是当前主流,DM7 的日志格式和 DM8 有差异,选捕获组件时先确认支持的达梦版本。

# 确认 Flink 版本与 Kafka 连接器版本 flink --version # 输出示例:Version: 1.18.1 # 查看已安装的 Kafka 连接器 jar ls $FLINK_HOME/lib | grep kafka # 应看到 flink-sql-connector-kafka-3.x.x.jar

flink --version确认 Flink 主版本,决定后续用哪个版本的连接器 jar。ls那步检查 Kafka 连接器是否就位,Flink 1.18 对应 Kafka 连接器 3.1.0 左右。如果 jar 缺失,从 Flink 官方仓库下载对应版本放进lib目录,重启集群生效。注意不要混用不同大版本的连接器,否则运行时会报NoSuchMethodError,这类错误堆栈很长但根因就是版本错配。

2.3 增量快照与日志消费的衔接:从全量到增量的切换

实时同步不是一上来就读日志,而是先做一次全量快照,再无缝切到增量日志。这个切换点的处理是整条链路最容易翻车的地方。如果快照期间源库还在写入,快照读到的数据和日志起点对不上,就会丢数据或重复。

常见做法是:捕获组件先记录一个日志位点(LSN 或 SCN),然后开始全量快照,快照完成后从记录的位点开始消费日志。这样快照期间产生的变更会被日志补上。达梦的逻辑日志里带有提交时间戳和事务号,捕获组件据此排序,保证顺序。

-- 获取当前达梦数据库的日志位点(示例,具体函数以达梦版本为准) SELECT SF_GET_LSN(); -- 返回一个数值型位点,用于标记增量起点

SF_GET_LSN()是达梦提供的获取当前日志序列号的函数,不同版本函数名可能略有差异,DM8 上可用。拿到位点后,全量快照的 SQL 要加一致性读提示,避免读到中间状态。快照完成后,捕获组件从这个 LSN 开始拉日志,Flink 侧用earliest或指定 offset 消费 Kafka,确保不漏。

提示:全量快照阶段如果表特别大,建议分批加并行度,但每批的位点要统一,否则增量起点会乱。

3. 从零搭一条达梦到 Kafka 再到 Flink 的同步链路

3.1 达梦侧准备:归档开启与逻辑日志配置

动手前先把达梦侧配置到位。归档开启是硬前提,逻辑日志按需开。下面是一套在测试库上验证过的步骤。

-- 1. 以 SYSDBA 登录达梦 -- 2. 将数据库切换到 MOUNT 状态 ALTER DATABASE MOUNT; -- 3. 开启归档模式 ALTER DATABASE ARCHIVELOG; -- 4. 配置归档路径(示例路径,按实际磁盘规划) ALTER DATABASE ADD ARCHIVELOG 'DEST=/dm8/arch, TYPE=LOCAL, FILE_SIZE=1024, SPACE_LIMIT=102400'; -- 5. 打开数据库 ALTER DATABASE OPEN;

第 2 步切 MOUNT 是必须的,归档配置只能在 MOUNT 下改。第 4 步的FILE_SIZE单位是 MB,SPACE_LIMIT是归档空间上限,单位也是 MB,这里设 100GB。归档路径要放在独立磁盘上,和生产数据盘分开,避免 IO 争抢。第 5 步 OPEN 之后用 2.1 节的查询确认ARCH_MODE为Y。

逻辑日志的开启在达梦配置文件dm.ini里,找到ENABLE_LOGIC_LOG参数,设为 1,然后重启实例。这个参数开启后会有额外写入开销,测试环境先评估对业务的影响。

# dm.ini 中逻辑日志相关配置 ENABLE_LOGIC_LOG = 1 # 开启逻辑日志 LOG_BUFFER_SIZE = 64 # 日志缓冲区大小,单位 MB

ENABLE_LOGIC_LOG置 1 后,达梦会额外维护逻辑日志结构,捕获组件读的就是它。LOG_BUFFER_SIZE适当调大能减少日志刷盘频率,但会占内存,64MB 是个折中值,高并发写入场景可以加到 128。

3.2 捕获组件到 Kafka:Topic 设计与消息格式

捕获组件把达梦变更写成消息推到 Kafka,Topic 设计直接影响下游消费效率。我一般按「库.表」粒度建 Topic,比如dm_orders、dm_users,而不是所有表塞一个 Topic。原因是不同表的变更速率差异大,混在一起会让慢表拖累快表的消费位点。

消息格式推荐用 JSON 或 Avro。JSON 可读性好,调试方便;Avro 体积小,适合高吞吐。下面是一个 JSON 消息的典型结构。

{ "op": "UPDATE", "table": "ORDERS", "ts": 1710000000000, "before": {"ID": 1001, "AMOUNT": 200.00, "STATUS": "NEW"}, "after": {"ID": 1001, "AMOUNT": 250.00, "STATUS": "PAID"} }

op标识操作类型,取值 INSERT/UPDATE/DELETE。ts是变更时间戳,毫秒级。before和after分别是变更前后的行镜像,UPDATE 时两者都有,INSERT 只有 after,DELETE 只有 before。下游 Flink 根据op决定是插入、更新还是删除。

Kafka Topic 的分区数按峰值吞吐定,单分区写入能力大概几 MB/s,如果变更量在 10MB/s 以上,至少开 6 个分区。副本数生产环境设 2 或 3,测试环境 1 即可。

# 创建 Kafka Topic,6 分区,2 副本 kafka-topics.sh --create \ --bootstrap-server kafka-broker:9092 \ --topic dm_orders \ --partitions 6 \ --replication-factor 2

--partitions 6决定并行消费上限,Flink 侧 source 并行度不要超过分区数,否则有线程空转。--replication-factor 2保证一个 broker 挂掉数据不丢。创建后用--describe确认分区分布均匀。

3.3 Flink SQL 消费与写入:建表、映射与 Exactly-Once

Flink 侧用 SQL 最省事。先建 Kafka 源表,再建下游结果表,中间用 INSERT INTO 串起来。

-- 建 Kafka 源表,解析 JSON 变更消息 CREATE TABLE dm_orders_cdc ( op STRING, `table` STRING, ts BIGINT, before_row ROW<ID BIGINT, AMOUNT DECIMAL(10,2), STATUS STRING>, after_row ROW<ID BIGINT, AMOUNT DECIMAL(10,2), STATUS STRING>, proc_time AS PROCTIME() ) WITH ( 'connector' = 'kafka', 'topic' = 'dm_orders', 'properties.bootstrap.servers' = 'kafka-broker:9092', 'properties.group.id' = 'flink-dm-cdc', 'scan.startup.mode' = 'earliest-offset', 'format' = 'json', 'json.fail-on-missing-field' = 'false', 'json.ignore-parse-errors' = 'true' );

scan.startup.mode设earliest-offset表示从最早位点消费,首次启动能拿到全量历史。json.fail-on-missing-field和json.ignore-parse-errors都设 true,避免个别脏消息让整个作业挂掉。proc_time是处理时间字段,后续做窗口聚合会用到。

下游建一张结果表,比如写到另一个 Kafka Topic 或 JDBC 表。

-- 建下游结果表(示例:写入 MySQL) CREATE TABLE orders_result ( id BIGINT, amount DECIMAL(10,2), status STRING, PRIMARY KEY (id) NOT ENFORCED ) WITH ( 'connector' = 'jdbc', 'url' = 'jdbc:mysql://mysql-host:3306/dw', 'table-name' = 'orders', 'username' = 'dw_user', 'password' = 'dw_pass' ); -- 把变更应用到结果表 INSERT INTO orders_result SELECT COALESCE(after_row.ID, before_row.ID) AS id, COALESCE(after_row.AMOUNT, before_row.AMOUNT) AS amount, COALESCE(after_row.STATUS, before_row.STATUS) AS status FROM dm_orders_cdc WHERE op <> 'DELETE';

COALESCE的作用是兼容 INSERT 和 UPDATE:INSERT 时 after_row 有值,UPDATE 时两者都有,取 after 优先。DELETE 操作这里直接过滤掉,如果要同步删除,需要下游表支持 delete 语义,JDBC 连接器可以通过PRIMARY KEY加NOT ENFORCED让 Flink 识别主键,从而下发 UPDATE/DELETE。

Exactly-Once 依赖 Kafka 的 offset 提交和下游事务。Kafka source 开启 checkpoint 后自动提交 offset,下游 JDBC 连接器支持事务写入。在flink-conf.yaml里把 checkpoint 间隔设成 10 秒左右。

# flink-conf.yaml 关键配置 execution.checkpointing.interval: 10s execution.checkpointing.mode: EXACTLY_ONCE state.backend: rocksdb state.checkpoints.dir: hdfs:///flink/checkpoints

execution.checkpointing.interval是 checkpoint 周期,10 秒兼顾延迟和开销。state.backend用 rocksdb 适合大状态,如果状态不大用默认的 hashmap 也行。state.checkpoints.dir指向持久化存储,作业重启后从这里恢复,保证不丢不重。

4. 同步链路的避坑与排查:五个真实翻车记录

4.1 归档空间写满导致同步静默中断

现象:同步跑了一周突然没数据了,Flink 作业没报错,Kafka 也没新消息,查达梦发现归档目录满了,新日志写不进去。

原因:归档空间上限设了但没配清理策略,达梦不会自动删旧归档,写满后归档进程挂起,捕获组件读不到新日志。

解决:给归档目录配定时清理,保留最近 3 天或按容量保留 80%。达梦可以用SF_ARCHIVELOG_DELETE_BEFORE_TIME函数删指定时间前的归档,配合 crontab 每天跑一次。

-- 删除 3 天前的归档日志 SELECT SF_ARCHIVELOG_DELETE_BEFORE_TIME(SYSDATE - 3);

4.2 全量快照与增量起点错位导致丢数

现象:全量同步完成后,增量阶段发现部分在快照期间更新的记录没同步过来。

原因:快照开始时没记录日志位点,或者记录的位点晚于快照读的时间点,导致快照期间的变更既不在快照里也不在增量里。

解决:严格按「先记位点、再开快照、快照完从位点消费」的顺序。位点记录和快照启动之间不能有业务写入窗口,必要时短暂加表锁或选业务低峰期。

4.3 Kafka 分区数与 Flink 并行度不匹配

现象:Flink 作业有 8 个并行度,但 Kafka Topic 只有 3 个分区,结果 5 个 subtask 空转,整体吞吐上不去。

原因:Kafka source 的并行度受分区数限制,多出来的 subtask 分不到分区。

解决:Topic 分区数至少等于 Flink source 并行度。如果已经建了 Topic,可以用kafka-topics.sh --alter --partitions扩分区,但注意扩分区会改变 key 的路由,如果消息有 key 且下游依赖 key 顺序,扩分区要谨慎。

4.4 达梦逻辑日志未开导致捕获组件空转

现象:捕获组件启动了,Kafka 里一条消息都没有,查达梦归档正常,业务也在写入。

原因:dm.ini里ENABLE_LOGIC_LOG没开,捕获组件读不到逻辑日志,只能干等。

解决:确认ENABLE_LOGIC_LOG = 1并重启实例。重启前评估影响,逻辑日志开启后写入延迟会略有增加,但通常在可接受范围。

4.5 下游 JDBC 写入主键冲突导致作业反复重启

现象:Flink 作业频繁重启,日志里报Duplicate entry主键冲突。

原因:UPDATE 消息被当成 INSERT 处理,或者 DELETE 没过滤导致下游重复插入。

解决:确认下游表声明了PRIMARY KEY ... NOT ENFORCED,让 Flink 识别为 upsert 模式。同时检查 SQL 里对op的处理逻辑,UPDATE 走 upsert,DELETE 要么过滤要么走 delete 语义。

5. 让同步链路更稳的两个进阶技巧

第一个技巧是给 Kafka 消息加 schema 版本号。达梦表结构变更(加字段、改类型)时,如果消息格式没版本标识,下游 Flink 解析会直接失败。我一般在上游消息里加一个schema_version字段,Flink 侧用CASE WHEN或自定义 UDF 做兼容解析。这样加字段时老版本消息还能读,新字段给默认值,避免作业中断。

-- 在源表里增加 schema_version 字段 -- 下游解析时按版本分支处理 SELECT CASE WHEN schema_version = 1 THEN after_row.AMOUNT WHEN schema_version = 2 THEN after_row.AMOUNT_V2 ELSE CAST(NULL AS DECIMAL(10,2)) END AS amount FROM dm_orders_cdc;

第二个技巧是用 Flink 的STATEMENT SET把多个表的同步写在一个作业里,减少作业数和资源开销。但要注意,多表共用一个 source 时,某张表的延迟会拖累其他表,所以只把变更速率相近的表放一起。

-- 多表同步示例,共用 checkpoint STATEMENT SET BEGIN INSERT INTO orders_result SELECT ... FROM dm_orders_cdc; INSERT INTO users_result SELECT ... FROM dm_users_cdc; END;

验证同步是否可靠,我习惯做一个对账任务:每天凌晨比对源库和目标库的行数和关键字段校验和,差异超过阈值就告警。这个对账不追求实时,但能兜住那些偶发的丢数。达梦侧用SELECT COUNT(*)和SUM(CRC32(...)),下游用同样逻辑,两边一比就知道有没有漏。

这套链路我前后调了大概两周,最大的教训是别信「配好就不管」——归档空间、位点记录、分区匹配这三件事,任何一件偷懒都会在某个凌晨变成电话。把监控和告警做在前面,比事后排查省心得多。希望帮到你。

本文还有配套的精品资源,点击获取

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询