资讯详情

FlinkCDC 实时同步达梦数据库:日志级增量采集与 Kafka 链路实践

📅 2026/10/10 21:43:53 | 华诺云谱 👁 阅读
FlinkCDC 实时同步达梦数据库:日志级增量采集与 Kafka 链路实践
简介本资源面向大数据开发工程师与实时数仓建设者聚焦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_TYPEARCH_TYPE为LOCAL时归档在本机为REMOTE时走远程。如果归档没开需要用ALTER DATABASE MOUNT后ALTER DATABASE ARCHIVELOG再OPEN这一步会短暂停库务必在维护窗口做。2.2 FlinkCDC 怎么接达梦连接器选型与版本匹配FlinkCDC 官方连接器覆盖 MySQL、PostgreSQL、Oracle、SQL Server 等达梦并不在默认列表里。所以接达梦有两条路一是用达梦官方或社区提供的 CDC 工具把变更推到 KafkaFlink 再从 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.jarflink --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, TYPELOCAL, FILE_SIZE1024, SPACE_LIMIT102400; -- 5. 打开数据库 ALTER DATABASE OPEN;第 2 步切 MOUNT 是必须的归档配置只能在 MOUNT 下改。第 4 步的FILE_SIZE单位是 MBSPACE_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 # 日志缓冲区大小单位 MBENABLE_LOGIC_LOG置 1 后达梦会额外维护逻辑日志结构捕获组件读的就是它。LOG_BUFFER_SIZE适当调大能减少日志刷盘频率但会占内存64MB 是个折中值高并发写入场景可以加到 128。3.2 捕获组件到 KafkaTopic 设计与消息格式捕获组件把达梦变更写成消息推到 KafkaTopic 设计直接影响下游消费效率。我一般按「库.表」粒度建 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 只有 afterDELETE 只有 before。下游 Flink 根据op决定是插入、更新还是删除。Kafka Topic 的分区数按峰值吞吐定单分区写入能力大概几 MB/s如果变更量在 10MB/s 以上至少开 6 个分区。副本数生产环境设 2 或 3测试环境 1 即可。# 创建 Kafka Topic6 分区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-OnceFlink 侧用 SQL 最省事。先建 Kafka 源表再建下游结果表中间用 INSERT INTO 串起来。-- 建 Kafka 源表解析 JSON 变更消息 CREATE TABLE dm_orders_cdc ( op STRING, table STRING, ts BIGINT, before_row ROWID BIGINT, AMOUNT DECIMAL(10,2), STATUS STRING, after_row ROWID 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 和 UPDATEINSERT 时 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/checkpointsexecution.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 走 upsertDELETE 要么过滤要么走 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(...))下游用同样逻辑两边一比就知道有没有漏。这套链路我前后调了大概两周最大的教训是别信「配好就不管」——归档空间、位点记录、分区匹配这三件事任何一件偷懒都会在某个凌晨变成电话。把监控和告警做在前面比事后排查省心得多。希望帮到你。本文还有配套的精品资源点击获取
📝

华诺云谱内容团队

资深建站顾问 · 行业研究员

10年+企业数字化服务经验,专注智能建站、SEO优化与品牌营销,持续输出建站技巧、行业洞察与营销干货,已帮助5000+企业实现数字化增长。

你可能需要的服务

订阅华诺云谱资讯周报

每周一封,精选建站技巧、SEO与营销干货,直达邮箱。已有 8,000+ 企业主订阅,助你少走弯路。

↑