简介: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 事件写入 Kafka,Flink CDC 侧负责消费 Kafka、做清洗转换、再分发到目标端。你不需要再纠结“Flink CDC 怎么直连达梦”,那不是重点,重点是事件格式和消费位点怎么对齐。
2.3 上线前的环境准备:归档、权限、网络与时钟
环境准备直接影响日志采集端能不能启动。第一件事是确认目标达梦库已经开启归档,并且归档目录有独立磁盘。常见做法是在 DIsql 里以 SYSDBA 执行下面的命令,注意不同版本命令写法有差异,以你当前环境的《达梦数据库管理员手册》为准:
-- 设置本地归档目录,目录要先创建好,且不能和数据库文件放同一个磁盘 ALTER DATABASE ADD ARCHIVELOG 'DEST=/dm8/arch/TYPE=LOCAL'; -- 开启归档模式,需要重启数据库实例生效 ALTER DATABASE ARCHIVELOG;代码的逻辑说明:第一句是将归档日志输出到/dm8/arch目录,TYPE=LOCAL表示本地归档,这是最常用的方式;第二句是把实例切换到归档模式。很多同步任务启动失败、日志采集端收不到增量,查到最后就是归档没开或者开完没重启实例。改完归档后建议再查一下实例状态,确认确实切换成功:
SELECT NAME, STATUS$, ARCH_MODE FROM V$DATABASE;注意:不同版本的V$DATABASE视图字段名不完全一样,有的版本叫ARCH_MODE,有的叫ARCHIVELOG。如果执行报“列不存在”,就去 DIsql 里执行DESC V$DATABASE看实际字段,别照抄。归档配置完成后,还要检查两件事:给日志采集端单独建一个同步账号,不要直接用 SYSDBA 跑业务任务,权限按“能读日志、能查目标表”的最小集授;另外确认所有节点时钟一致,最好有 NTP 同步。时钟不一致会让事件时间戳错乱,后面做延迟监控时你会看到延迟忽正忽负,非常难受。
3. 跑通第一条同步链路:从达梦日志到 Flink SQL 源表
3.1 链路设计:日志采集端写入 Kafka,Flink 只做消费与计算
我建议的最小可运行链路是:达梦归档日志 → 官方日志采集组件 → Kafka Topic → Flink SQL。日志采集端在达梦服务器上跑,它把翻译后的逻辑事件写到 Kafka,Flink 侧不关心达梦日志格式,只消费 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、DELETE;row_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 ROW<id INT, name STRING>, after_row ROW<id 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-offsets;json.ignore-parse-errors建议打开,因为达梦某些类型在极端情况下会产出异常字符串,宁可丢一条坏消息也不要让整个作业卡死;WATERMARK是给event_time定义乱序容忍度,如果你对数据顺序不敏感,可以删掉这行,用处理时间就行。这里ROW<id 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 标准写法,端口默认是 5236,sync_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 个区,源端并行度设 6,Sink 并行度设 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_TAB,Flink 作业也不报错,但下游目标库收到数据后写入失败,报“模式不存在”或“无效的模式名”;更隐蔽的情况是同步过来的数据有值,但落到了错误的模式或表里。
原因:达梦数据库里模式和用户名强绑定,默认情况下一个用户对应一个同名模式,比如用户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 ROW<id INT, name STRING>, after_row ROW<id 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和当前时间,记录到一张运维笔记里。这套动作做多了,哪些配置改动会影响延迟、哪些不会,心里就有数了。希望帮到你。
本文还有配套的精品资源,点击获取