资讯动态

SpringBoot整合Canal监听MySQL binlog,RabbitMQ异步同步缓存与数据

发布时间:2026/9/14 7:36:23 来源:尧图企业网站定制
有个做电商的朋友前几天找我诉苦系统里Redis缓存的商品库存又跟MySQL对不上了排查到最后发现是运营直接在数据库改了价格代码里负责发消息同步缓存的那条路径完全没感知。他问我要不要在业务代码里所有写库的地方都插一段“发送变更消息”的逻辑。我说千万别你管得住自己的代码管不住DBA的SQL、管不住运营后台上的一键批量更新、更管不住历史数据订正脚本。只要是数据库里的数据变动最干净的感知方式只有一种——从binlog里拿。于是就有了这套我用了很久的组合SpringBoot整合Canal监听MySQL的binlog变更再把变更事件投递到RabbitMQ由业务服务异步消费处理用来刷缓存、同步搜索引擎、推送到大数据平台都非常合适。这篇文章把我的完整落地过程、配置、代码和踩过的坑一次性放出来适合做微服务、做数据同步、做中间件方向的Java后端同学参考。1. 缓存又对不上了为什么偏要监听数据库变更1.1 一个典型的异步同步困境几乎所有互联网业务都会遇到同一个问题MySQL里的数据变了但MySQL以外的存储还停留在旧状态。最常见的就是Redis缓存其次是Elasticsearch索引再往后还有数仓里的宽表、大数据平台的日志、下游系统的订单快照。很多人第一个想到的方案就是在业务代码里做双写。比如用户修改了个人资料Service层先更新MySQL再手动删缓存、再发一条MQ消息告诉下游去更新。这个方案在系统规模小、团队边界清晰的时候没毛病但一旦业务复杂起来就变味了改了A服务B服务也在写这张表今天加了个C服务DBA临时跑了个批量UPDATE订正数据运营同学在管理后台改了个状态字段。这些变更全部绕过了你代码里精心埋好的“通知逻辑”于是缓存还是旧数据搜索还是旧数据。1.2 对比几个常见方案的取舍我把“感知数据变更”这件事拆开来看其实有四种主流思路方案变更感知完整性实时性对业务代码侵入维护成本业务代码双写/MQ通知低只覆盖业务入口高高中定时任务全量扫描高但只能扫到有更新时间字段的表低分钟级起步低中应用层拦截SQL解析高但不支持跨语言高中很高监听binlog高任何入口的变更都能感知高秒级几乎为零中相比之下基于binlog的监听有天然优势它不管你的变更来自Java代码、PHP脚本还是Navicat手工改数只要MySQL产生了binlog事件就能被捕获根本不关心上层业务是怎么写的。1.3 为什么中间要塞一个RabbitMQCanal本身是把binlog采集并解析好的组件它能把变更事件推给客户端。但生产环境不会让Canal直连业务服务因为Canal的位置点管理、订阅推送、故障恢复都需要和消费方解耦。在中间加一层RabbitMQ价值就出来了Canal只负责生产消息不用管消费者在不在线消费者可以自由扩容缩容哪怕业务服务重启、发布、升级消息也先在队列里待着不会丢以后要加新的消费方比如再加一个数据同步任务只需要再写一个消费者监听同一个队列或者单独绑一个队列原有链路完全不用动。这套方案的完整链路长这样MySQL binlog - Canal Server - RabbitMQ Exchange - Queue - SpringBoot Consumer接下来我一步步拆解从原理到环境再到代码和坑保证跟着做能直接跑通。2. 先搞懂Canal在干嘛伪装成从库读binlog2.1 MySQL主从复制的核心机制要理解Canal得先理解MySQL主从复制是怎么工作的。主库把每次数据变更写入binlog二进制日志从库启动一个IO线程跟主库建立连接请求“从某个binlog文件的某个位置开始给我推送日志”主库负责把binlog一段段推给从库从库拿到日志后通过SQL线程回放完成数据同步。整个机制的突破口在于从库跟主库通信使用的是一套公开的、基于MySQL协议的复制协议而协议本身并不校验对面的进程到底是不是一个真正的从库实例。Canal做的事情非常“钻空子”——它自己用Java实现了一个假的从库向主库发送dump请求主库就真把它当成从库来推日志了。Canal拿到binlog的原始字节流之后再把里面的表结构、字段值、变更类型解码成结构化数据。2.2 binlog三种格式里为什么必须用ROWbinlog有三种记录格式这个点直接影响Canal能不能干活。格式记录内容特点Canal可用性STATEMENT记录执行SQL日志量小但某些函数和不确定操作无法准确回放不可用拿不到具体字段变更ROW记录每行数据的前后镜像日志量大但最准确、信息最全必须使用MIXED根据SQL动态选择STATEMENT或ROW折中方案不稳定不推荐所以MySQL侧必须把binlog-format设置为ROW否则Canal解析出来的可能就是一堆SQL文本而不是精确的行级变更数据。2.3 flatMessage这个配置决定了消息好不好读Canal把binlog解析成结构化事件后对外发布消息有两种格式由canal.mq.flatMessage控制。false是默认值消息体是protobuf序列化后的二进制解析需要引入Canal的protobuf依赖优点是体积小true则输出为JSON文本字段名一目了然SpringBoot里用Fastjson或Jackson直接转成Map就能用。我在实际项目里一律用true牺牲一点网络带宽换来的是排查问题时的极度舒适。3. 环境落地MySQL开binlog、Canal服务端部署、RabbitMQ初始化3.1 MySQL开启binlog并新建Canal账号第一步是确保MySQL开启了binlog并且格式为ROW。打开MySQL配置文件一般位于/etc/my.cnf或/etc/mysql/mysql.conf.d/mysqld.cnf在[mysqld]区块加上以下配置[mysqld] server-id1 log-binmysql-bin binlog-formatROW binlog_row_imageFULL这里有个容易踩的坑server-id是主从复制环境里区分不同节点的关键参数如果你的MySQL上已经挂了别的从库这个server-id一定不要跟现有从库重复否则主库会报错。binlog_row_image从MySQL 5.6开始支持默认值是FULL在8.0里也是默认。如果被改成了MINIMALUPDATE事件里的旧值镜像会不完整导致Canal收到的事件里old字段缺失给后面做变更对比带来麻烦。我建议显式写出来防止环境迁移时被默认值坑到。改完配置要重启MySQL然后确认生效SHOW VARIABLES LIKE log_bin; SHOW VARIABLES LIKE binlog_format;接下来创建一个专门给Canal用的MySQL账号。这个账号不需要DBA权限只需要复制链路相关的权限即可CREATE USER canal% IDENTIFIED BY canal; GRANT SELECT, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO canal%; FLUSH PRIVILEGES;REPLICATION SLAVE是用于复制协议的REPLICATION CLIENT用于查询主库状态SELECT则确保Canal能读取表结构元数据。3.2 Canal服务端部署Docker一行命令Canal官方提供了Docker镜像我推荐直接用它部署省去本地安装JDK和配置环境的麻烦也方便版本切换。以下面的命令为例docker run --name canal-server -d \ -p 11111:11111 \ -e canal.instance.master.address127.0.0.1:3306 \ -e canal.instance.dbUsernamecanal \ -e canal.instance.dbPasswordcanal \ -e canal.instance.connectionCharsetUTF-8 \ -e canal.instance.filter.regextestdb\\..* \ -e canal.serverModeMQ \ -e canal.mq.servers127.0.0.1:5672 \ -e canal.mq.vhost/canal \ -e canal.mq.usernamecanal \ -e canal.mq.passwordcanal \ -e canal.mq.exchangecanal.exchange \ -e canal.mq.queuecanal.queue \ -e canal.mq.routing.keycanal.routing.key \ -e canal.mq.flatMessagetrue \ registry.cn-hangzhou.aliyuncs.com/canal/canal-server:v1.1.7简单解释几个关键环境变量canal.instance.master.addressMySQL的主机地址和端口。注意Docker容器内访问宿主机MySQL时不能用127.0.0.1要填宿主机的局域网IP或者用--network host启动容器。canal.instance.filter.regex要监听哪个库的哪张表格式是库名\\..*多个规则用逗号分隔。正则里的反斜杠在配置文件里需要转义写成testdb\\..*。canal.serverModeMQ让Canal以MQ模式启动而不是默认的TCP模式。canal.mq.serversRabbitMQ的AMQP服务端口是5672不是Web管理界面那个15672。canal.mq.flatMessagetrue消息体输出JSON方便SpringBoot直接解析。如果你需要监听多个库或者多个实例更推荐把Canal的配置文件挂载出来用传统方式改。上面的Docker环境变量方式适合单实例快速启动跑通链路后你再根据实际情况调整。3.3 RabbitMQ侧要提前创建好哪些资源Canal不会自动创建RabbitMQ里的vhost、exchange、queue和binding这些必须提前准备。这一步经常有人漏掉结果Canal日志里一直报投递失败但MySQL和SpringBoot看起来都没问题。我的做法是用RabbitMQ的管理界面人工创建一次顺便把账号权限也配好在Admin标签页新建用户canal密码canal。在Virtual Hosts标签页新建vhost/canal。给canal用户配置/canal这个vhost下的所有权限。在/canal下新建direct类型的exchange名称为canal.exchange。在/canal下新建queue名称为canal.queue。将canal.exchange和canal.queue绑定routing key为canal.routing.key。如果你习惯命令行用rabbitmqadmin也一样rabbitmqctl add_vhost /canal rabbitmqctl add_user canal canal rabbitmqctl set_permissions -p /canal canal .* .* .* rabbitmqadmin declare exchange namecanal.exchange typedirect vhost/canal durabletrue rabbitmqadmin declare queue namecanal.queue vhost/canal durabletrue rabbitmqadmin declare binding sourcecanal.exchange destinationcanal.queue destination_typequeue routing_keycanal.routing.key vhost/canalRabbitMQ有两个端口非常容易弄混5672是AMQP协议端口服务间收发消息走这个15672是Web管理界面端口浏览器访问用这个。如果你发现SpringBoot连不上RabbitMQ先确认连的是不是5672。至于本地RabbitMQ启动失败的问题我排查时基本固定两条思路先看Erlang和RabbitMQ的版本是否匹配再看启动日志里的具体报错。RabbitMQ对Erlang版本要求很苛刻版本对不上起不来非常常见日志一般在/var/log/rabbitmq/下里面有明确的error信息比瞎猜强得多。4. SpringBoot消费端代码从配置到事件处理4.1 项目依赖和基础配置新建一个SpringBoot工程引入三个核心依赖spring-boot-starter-amqpRabbitMQ、spring-boot-starter-data-redis示例里用来刷缓存、fastjson2解析消息体。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdcom.alibaba.fastjson2/groupId artifactIdfastjson2/artifactId version2.0.51/version /dependencyapplication.yml里的RabbitMQ配置如下spring: rabbitmq: host: 127.0.0.1 port: 5672 username: canal password: canal virtual-host: /canal listener: simple: acknowledge-mode: manual concurrency: 3 max-concurrency: 10 prefetch: 20这里有几个关键点要说明。acknowledge-mode: manual表示手动确认这是生产环境必须开的后面我会专门讲为什么。concurrency和max-concurrency控制消费线程数prefetch是每个消费者预取的消息数20表示消费者一次性最多拉20条到本地待处理处理完一条ack一条才能继续拿。如果prefetch不设置RabbitMQ默认会把消息全推给消费者处理不过来时消息积压在消费者本地队列界面显示unacked飙升服务重启还会造成大量重复消费。4.2 Canal发过来的消息长什么样在写代码之前先给你看一条Canal投递到RabbitMQ的JSON消息这比读文档更直观{ data: [ { id: 1, name: 张三, age: 20 } ], database: testdb, es: 1718680200000, id: 452, isDdl: false, mysqlType: { id: bigint, name: varchar(50), age: int }, old: [ { age: 18 } ], pkNames: [id], sql: , sqlType: { id: -5, name: 12, age: 4 }, table: user, ts: 1718680200123, type: UPDATE }这条消息对应的是表testdb.user里id1这一行发生了UPDATEage从18变成了20。核心字段就是这几个type变更类型INSERT、UPDATE、DELETE三种。database和table变更发生在哪个库的哪张表。data变更后的行数据列表一次变更涉及多行时这里就是多个对象。oldUPDATE时被修改字段的旧值列表INSERT和DELETE这里为null。isDdl是不是DDL语句比如ALTER TABLE、CREATE TABLE这类事件没有行级数据。sqlDDL语句原文普通DML事件为空字符串。tsCanal解析这条binlog的时间戳。pkNames主键字段名列表DELETE时特别有用因为data里通常只保留主键值。4.3 消费者核心代码Canal消息体直接就是JSON字符串所以定义一个对应的实体类来接收public class CanalMessage { private ListMapString, Object data; private String database; private Long es; private Long id; private Boolean isDdl; private MapString, String mysqlType; private ListMapString, Object old; private ListString pkNames; private String sql; private MapString, Integer sqlType; private String table; private Long ts; private String type; // 省略getter/setter }然后写一个消费者用RabbitListener监听canal.queueComponent public class CanalMessageConsumer { private static final Logger log LoggerFactory.getLogger(CanalMessageConsumer.class); Autowired private StringRedisTemplate redisTemplate; RabbitListener(queues canal.queue) public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag message.getMessageProperties().getDeliveryTag(); String body new String(message.getBody(), StandardCharsets.UTF_8); try { CanalMessage canalMessage JSON.parseObject(body, CanalMessage.class); if (Boolean.TRUE.equals(canalMessage.getIsDdl())) { log.info(收到DDL事件忽略处理: {}.{} - {}, canalMessage.getDatabase(), canalMessage.getTable(), canalMessage.getSql()); channel.basicAck(deliveryTag, false); return; } String type canalMessage.getType() null ? : canalMessage.getType().toUpperCase(); switch (type) { case INSERT: handleInsert(canalMessage); break; case UPDATE: handleUpdate(canalMessage); break; case DELETE: handleDelete(canalMessage); break; default: log.warn(未支持的变更类型: {}, type); } channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(消费Canal消息失败消息内容{}, body, e); channel.basicNack(deliveryTag, false, true); } } private void handleInsert(CanalMessage msg) { if (!user.equalsIgnoreCase(msg.getTable())) { return; } ListMapString, Object dataList msg.getData(); for (MapString, Object data : dataList) { String id String.valueOf(data.get(id)); redisTemplate.delete(cache:user: id); log.info(用户新增缓存已清理: id{}, id); } } private void handleUpdate(CanalMessage msg) { if (!user.equalsIgnoreCase(msg.getTable())) { return; } ListMapString, Object dataList msg.getData(); ListMapString, Object oldList msg.getOld(); for (int i 0; i dataList.size(); i) { MapString, Object data dataList.get(i); MapString, Object old (oldList ! null oldList.size() i) ? oldList.get(i) : null; String id String.valueOf(data.get(id)); redisTemplate.delete(cache:user: id); log.info(用户更新缓存已清理: id{}, 变更前数据{}, id, old); } } private void handleDelete(CanalMessage msg) { if (!user.equalsIgnoreCase(msg.getTable())) { return; } ListMapString, Object dataList msg.getData(); for (MapString, Object data : dataList) { String id String.valueOf(data.get(id)); redisTemplate.delete(cache:user: id); log.info(用户删除缓存已清理: id{}, id); } } }这段代码的逻辑非常简单拿到消息判断是不是DDL是就直接ack然后根据type分发到不同的处理方法。处理方式是删缓存这是因为缓存同步最稳妥的策略是“删缓存而非写缓存”——删掉之后下游查询就会回源MySQL重新加载天然保证最终一致而直接写缓存则容易遇到并发写和旧值覆盖的问题。4.4 一些值得注意的实现细节表名过滤我直接写在处理方法里用if (!user.equalsIgnoreCase(msg.getTable())) return;这种方式。如果以后监听表多了建议抽一个路由表或者用策略模式按表名分派到不同的handler代码会清爽很多。这里为了给你看核心逻辑保持了最简写法。关于msg.getOld()这个方法有一个细节容易误会UPDATE事件里old数组的元素数量和data数组是一致的但并不是每行都有非空old。比如一条批量UPDATE把10行数据的age都更新了如果某些行的age本来就和目标值相同MySQL在ROW格式下不会记录这些行的前后变化那它们对应的old元素就可能是null。所以遍历时做了oldList.size() i的判空防止数组越界和空指针。另外我在消费者开头没有对body做非空校验因为RabbitMQ不会投递空消息体。但如果Canal配置成flatMessagefalse消息体就是protobuf的二进制数据JSON.parseObject会直接解析失败。所以务必把canal.mq.flatMessage设成true。5. 联调验证一条UPDATE是怎样跑通全链路的5.1 准备一张测试表并接入监听假设要监听testdb库下的user表建表SQL长这样CREATE TABLE user ( id bigint(20) NOT NULL, name varchar(50) DEFAULT NULL, age int(11) DEFAULT NULL, PRIMARY KEY (id) ) ENGINEInnoDB DEFAULT CHARSETutf8mb4;Canal的监听正则我在前面配的是testdb\\..*也就是testdb库下所有表都能监听到。如果你想只监听user表可以把正则改成testdb\\.user。这个限定在业务上很有用——生产环境通常库很大全部监听会产生海量无效消息尽量把正则收窄到真正关心的表上。5.2 修改一条数据逐层看效果先把MySQL、Canal、RabbitMQ、SpringBoot四个服务全部启动然后在MySQL里执行一条UPDATEUPDATE user SET age 20 WHERE id 1;正常情况下几秒之内你就能在RabbitMQ管理后台的canal.queue页面看到消息数量从0变成1。这个现象本身就说明整条链路已经通了MySQL产生了binlogCanal解析并投递到了RabbitMQ消息正在队列里等着消费者处理。此时去RabbitMQ管理后台的Queue页面点击Get messages可以人工拉取一条消息看原文就长我前面贴的那个样。需要注意Get messages默认会把消息取出来如果选了ack_requeue_false消息被捞走之后就会从队列里消失。想看原文又不想消费消息的同学选ack_requeue_true即可。SpringBoot这边的日志会输出用户更新缓存已清理: id1, 变更前数据{age18}到这一步链路就算完全跑通了。5.3 换INSERT和DELETE再验证一遍我建议你联调时不要只测UPDATE把三种DML都过一遍。因为三种类型产生的消息结构差异很大INSERTdata中有完整的新行数据old为nulltype为INSERT。UPDATEdata中有新值old中有被修改字段的旧值。DELETEdata中通常只有主键字段的值old一般为null。这条消息的作用不是拿到被删行的完整数据而是告诉你“这个主键被删了”下游据此清理对应的缓存或索引。我当年第一次联调DELETE时还困惑过为什么data里的其他字段全是null后来想明白了binlog在ROW格式下DELETE事件只需要记录主键就能定位被删除的行其他字段对回放没有意义MySQL自然也懒得写进去。5.4 链路不通时怎么快速定位联调时如果发现消息没到按顺序排查这几个节点基本都能定位MySQL有没有产生binlog在MySQL里执行SHOW MASTER STATUS;看File和Position有没有变化。没变化说明binlog没开或者格式不对。Canal有没有连上MySQL查看Canal容器日志搜connected和binlog关键字。如果一直报连接拒绝检查canal.instance.master.address填的地址对不对以及MySQL账号密码是否配置正确。RabbitMQ里有没有消息管理后台看canal.queue的Ready数量。如果是0说明Canal没投递成功看Canal日志有没有rabbit相关的error。SpringBoot有没有消费看日志有没有打印“用户更新”之类的输出。没有的话检查消费者的queues名称和RabbitMQ里实际队列名是否完全一致。这套排查顺序我每次都会按因为从数据源到消费端是一条单向链路任何一个节点出问题后面的节点都拿不到数据。从前往后查效率最高。6. 生产化要提前想清楚的事幂等、顺序、丢消息6.1 RabbitMQ重复投递是一件必然的事先给你交个底这套链路里消息一定会重复。这不是RabbitMQ的bug而是分布式环境中ack确认机制导致的必然现象。比如消费者处理完业务逻辑还没来的及发送ack网络闪断RabbitMQ没收到确认就会认为消息没消费成功重新投递消费者重启时unacked的消息也会被重新入队。所以消费端的第一原则就是处理逻辑必须天然幂等。如果你在消费端是删缓存重复执行一次没有副作用天然幂等但如果你是往ES里写数据、给下游调用接口、往统计表里累加数字就必须自己保证幂等。一个简单的幂等方案是利用消息的唯一标识。Canal消息里的id字段代表binlog事件在某个binlog文件里的偏移位置全链路里是唯一的很适合做幂等键。消费时先尝试写入RedisBoolean ok redisTemplate.opsForValue() .setIfAbsent(canal:msg: canalMessage.getId(), 1, Duration.ofMinutes(5)); if (!Boolean.TRUE.equals(ok)) { log.info(重复消息已忽略: {}, canalMessage.getId()); return; }Redis里已经有这个key说明消息已经处理过直接跳过。setIfAbsent对应Redis的SETNX命令原子性有保障不会出现并发判断的缝隙。TTL设5分钟因为消息重复通常发生在短时间内超过这个时间基本不会再来一遍。6.2 顺序问题别让并发消费把数据搞乱binlog里的变更事件本身是有序的一个事务的多个变更Canal会按顺序推出来。但一旦进入RabbitMQ消费者的并发执行就会打破这种顺序。举个典型场景同一行的数据先收到一条UPDATE改成20又收到一条UPDATE改成30两个consumer各处理一条结果后处理的反而是改成20那条缓存最终就停在了20这个旧值上。解决顺序问题有几种思路适用场景各不相同方案具体做法优点缺点单线程消费concurrency1实现简单顺序绝对保证吞吐量受限按主键路由到不同队列用routing key区分同一类表不同主键绑定不同队列兼顾并发和单行顺序配置复杂扩缩容要小心处理端带上一次时间戳对比消费时比较消息时间戳和当前数据时间戳旧的不处理能缓解乱序影响需要额外存储和判断侵入性强我的建议是看业务容忍度。缓存同步、搜索索引这类场景允许秒级不一致只要最终一致就行一般不用为顺序做太多设计但如果你要同步的是账户余额这类强一致数据单线程消费就是最简单的保命方案先把吞吐放一边数据不能错。6.3 手动确认和异常处理的正确姿势我把acknowledge-mode设成manual手动确认。消费者在处理完业务后调用channel.basicAck(deliveryTag, false)让RabbitMQ知道消息已经处理完可以删除处理异常则调用channel.basicNack(deliveryTag, false, true)第三个参数true表示重新放回队列。这里有个陷阱basicNack的requeue参数如果设成true消息会无限重试如果这条消息有bug永远处理不了就会形成死循环把队列和CPU都打满。我在生产环境是这么处理的第一次重试用requeue但如果重试几次还是失败就投递到单独的dead-letter队列让定时任务去补偿或者人工介入。RabbitMQ支持通过DLX死信交换机配置把处理失败的消息自动转入死信队列这块在SpringBoot里可以靠RabbitListener的异常处理和额外的listener来实现属于进阶玩法但上线前一定要想好失败策略。6.4 三个环节的消息丢失风险消息可能丢失的环节主要有三个每个环节的防护手段不同环节风险点防护手段MySQL - CanalMySQL执行了变更但Canal未捕获Canal开启位点持久化重启后从上次位置继续但最好还是保证Canal高可用Canal - RabbitMQCanal投递消息时RabbitMQ不可用让RabbitMQ尽量稳定队列和exchange声明为durableRabbitMQ - 消费者消费者未ack就崩溃或持久化配置不当必须开手动ack持久化配置好消费者尽量不要长耗时其中MQ - 消费者这个环节是我们可以完全控制的。durabletrue保证exchange和queue在RabbitMQ重启后不会消失但消息本身是否持久化取决于发送端Canal是否设置了delivery_mode2。Canal在MQ模式投递RabbitMQ时是支持持久化的但为了保险最简单的兜底手段还是做定时对账比如每天凌晨跑一次批量比对把缓存和数据源不一致的记录捞出来重新刷新。系统设计里永远不要追求“不可能丢消息”而是要让“丢消息之后可以被发现和纠正”。6.5 日常监控队列积压和Canal延迟链路跑起来之后最需要关注的两个指标一是RabbitMQ队列的积压数量二是Canal从binlog解析到投递的延迟。RabbitMQ管理后台直接看canal.queue的Ready和Unacked数量就行。Ready长时间上涨说明消费者处理速度跟不上Unacked一直很高说明prefetch配置可能太大或者某个消费者卡住了。Canal的延迟可以用一个笨办法监控在数据库里建一张心跳表每10秒更新一次时间戳Canal监听这张表的变化消费者收到消息后把时间戳写入Redis。然后监控程序比对Redis里这条时间戳和当前时间之差超过一分钟就报警。这个方案不依赖任何第三方组件逻辑很直观我一直在用。它的本质是利用Canal自产自销一条低频消息通过观测这条消息的流转时延来判断整条链路是否健康。这套架构从原理到落地我已经完整走了一遍踩过最大的坑出现在RabbitMQ的队列资源没提前创建导致Canal投递失败上排障花了不少时间。生产环境建议把环境准备这一步固化成脚本或者自动化流程因为人肉操作越频繁漏配置的概率就越大。链路跑通之后我最大的体会是监听binlog这套方案真正解决的是“所有入口的数据变更都被感知”这个难题它带来的收益远大于那点运维成本。如果你的业务也在被缓存一致性、多系统数据同步折磨顺着这套思路先搭一个最小可用链路跑一段时间再逐步完善监控和异常处理你会发现很多老问题突然就不是问题了。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价