基于 MySQL Binlog + Canal 的 Redis 缓存一致性实战:异步更新与最终一致性方案 问题场景:为什么缓存一致性这么难? 假设你有一个用户信息接口,QPS 上了万级后,你自然想到用 Redis 挡一层。写操作时先更新 MySQL,再删除 Redis 缓存,读操作时缓存未命中就回源 DB 并重建缓存。这个流程看起来没问题,直到某天线上出现了「用户改了昵称,但 App 上显示的还是旧昵称」的 Bug。
排查发现:更新 MySQL 和删除 Redis 这两个操作不是原子的。如果删除 Redis 的动作因为网络抖动失败了,缓存里就一直是脏数据。你可能会说,那用事务或者消息队列做「双写」,但事务跨 MySQL 和 Redis 并不现实,而手动双写又容易漏掉某些业务入口——比如有人直接通过 SQL 改了库,或者有批量任务绕过了 Service 层。
真正的问题是:缓存的更新逻辑和业务代码耦合在一起,只要有一个入口没管住,缓存就会脏 。
有没有办法让 MySQL 的数据变更被自动、完整地捕获,然后异步同步到 Redis?答案就是 MySQL Binlog + Canal 。这也是我这几年做后端架构时最深刻的感受之一:好的方案不是靠人更谨慎,而是靠机制更可靠 。下面直接开始搭。
整体架构与数据流 先看整体方案:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 业务服务 ──写──▶ MySQL │ │ (Binlog ROW 格式) ▼ Canal Server ◀── 伪装成 MySQL Slave │ │ 解析后的变更事件 (JSON) ▼ 消息队列 (RocketMQ / Kafka) │ ▼ 缓存同步消费者 (Spring Boot) │ │ 更新 / 删除 ▼ Redis
关键点:
环节
作用
备注
MySQL Binlog
记录所有数据变更的日志
必须是 ROW 格式
Canal Server
模拟 MySQL 从库,订阅并解析 Binlog
阿里开源,支持 MySQL 5.7/8.0
消息队列
解耦 Canal 与消费者,消化峰值
推荐 RocketMQ / Kafka
消费端
解析变更事件,更新或失效 Redis 缓存
实现最终一致性
这个方案的优点是:MySQL 的每一次数据变更,无论来自哪个业务入口,最终都会经过 Binlog 被 Canal 捕获 。而你不需要在业务代码里写任何缓存同步逻辑,只专注于「消费者」这一个对账入口。
第一步:MySQL 开启 Binlog(ROW 格式) 编辑 MySQL 配置文件 /etc/my.cnf(或 /etc/mysql/mysql.conf.d/mysqld.cnf),添加或修改:
1 2 3 4 5 6 7 8 9 [mysqld] log-bin =mysql-binbinlog-format =ROWserver-id =1001 binlog-do-db =orders_db
重启 MySQL 并验证:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 SHOW VARIABLES LIKE 'log_bin' ;SHOW VARIABLES LIKE 'binlog_format' ;SHOW MASTER STATUS;
验证环境 :MySQL 8.0.36,Linux(CentOS Stream 9)
第二步:给 Canal 创建专用的 MySQL 账号 Canal 需要模拟从库,因此账号需要 REPLICATION SLAVE 和 REPLICATION CLIENT 权限:
1 2 3 CREATE USER 'canal' @'%' IDENTIFIED BY 'canal_password' ;GRANT SELECT , REPLICATION SLAVE, REPLICATION CLIENT ON * .* TO 'canal' @'%' ;FLUSH PRIVILEGES;
生产环境建议限制来源 IP,比如 'canal'@'10.0.0.10',不要用 %。
第三步:部署 Canal Server 下载 Canal(以 1.1.7 为例,支持 MySQL 8.0):
1 2 3 4 5 6 7 8 9 wget https://github.com/alibaba/canal/releases/download/canal-1.1.7/canal.deployer-1.1.7.tar.gz mkdir -p /opt/canal && tar -zxvf canal.deployer-1.1.7.tar.gz -C /opt/canalls /opt/canal
3.1 配置 Canal 连接 MySQL 编辑 /opt/canal/conf/example/instance.properties:
1 2 3 4 5 6 7 8 9 10 11 12 13 canal.instance.master.address =127.0.0.1:3306 canal.instance.dbUsername =canal canal.instance.dbPassword =canal_password canal.instance.master.journal.name =mysql-bin.000001 canal.instance.master.position =154 canal.instance.filter.regex =orders_db\\..*
3.2 配置 Canal 对接消息队列(以 RocketMQ 为例) 编辑 /opt/canal/conf/canal.properties:
1 2 3 4 5 6 7 8 9 canal.serverMode = rocketMQ rocketmq.namesrv.address = 127.0.0.1:9876 rocketmq.producer.group = canal_producer_group canal.mq.topic = db_cache_sync
启动 Canal:
1 2 3 /opt/canal/bin/startup.sh tail -f /opt/canal/logs/canal/canal.log
第四步:Spring Boot 消费端实现 4.1 Maven 依赖 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 <dependencies > <dependency > <groupId > org.springframework.boot</groupId > <artifactId > spring-boot-starter-web</artifactId > </dependency > <dependency > <groupId > org.springframework.boot</groupId > <artifactId > spring-boot-starter-data-redis</artifactId > </dependency > <dependency > <groupId > org.apache.rocketmq</groupId > <artifactId > rocketmq-spring-boot-starter</artifactId > <version > 2.3.0</version > </dependency > <dependency > <groupId > org.projectlombok</groupId > <artifactId > lombok</artifactId > <optional > true</optional > </dependency > </dependencies >
4.2 application.yml 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 server: port: 8081 spring: redis: host: 127.0 .0 .1 port: 6379 lettuce: pool: max-active: 16 max-idle: 8 min-idle: 2 rocketmq: name-server: 127.0 .0 .1 :9876 consumer: group: cache_sync_consumer_group topic: db_cache_sync
4.3 解析 Canal 消息的 Java 代码 Canal 投递到 MQ 的消息格式是一段扁平的 JSON,结构大致如下:
1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 { "data" : [ { "id" : "1001" , "nickname" : "新昵称" , "avatar_url" : "https://cdn.example.com/avatar/1001.png" , "updated_at" : "2026-08-22 08:00:00" } ] , "database" : "orders_db" , "table" : "user_profile" , "type" : "UPDATE" , "old" : [ { "nickname" : "旧昵称" } ] , "pkNames" : [ "id" ] , "sql" : "" , "ts" : 1755820800000 }
关键字段说明:
字段
含义
type
操作类型:INSERT、UPDATE、DELETE、ALTER 等
database / table
变更发生的库和表
data
变更后的数据(DELETE 时为空)
old
变更前的数据(仅 UPDATE 时有值)
pkNames
主键列名列表
4.4 消费者核心代码(完整可运行) 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 package com.example.cachelistener.consumer;import com.alibaba.fastjson2.JSON;import com.alibaba.fastjson2.JSONArray;import com.alibaba.fastjson2.JSONObject;import lombok.extern.slf4j.Slf4j;import org.apache.rocketmq.spring.annotation.RocketMQMessageListener;import org.apache.rocketmq.spring.core.RocketMQListener;import org.springframework.data.redis.core.StringRedisTemplate;import org.springframework.stereotype.Component;import jakarta.annotation.Resource;import java.util.List;@Slf4j @Component @RocketMQMessageListener( topic = "db_cache_sync", consumerGroup = "cache_sync_consumer_group" ) public class DBCacheSyncConsumer implements RocketMQListener <String> { private static final String CACHE_KEY_PREFIX = "user:profile:" ; @Resource private StringRedisTemplate stringRedisTemplate; @Override public void onMessage (String message) { try { JSONObject msg = JSON.parseObject(message); String type = msg.getString("type" ); if (!List.of("INSERT" , "UPDATE" , "DELETE" ).contains(type)) { log.debug("忽略非 DML 事件, type={}" , type); return ; } String table = msg.getString("table" ); String database = msg.getString("database" ); if (!"orders_db" .equals(database) || !"user_profile" .equals(table)) { return ; } JSONArray dataArray = msg.getJSONArray("data" ); JSONArray pkNames = msg.getJSONArray("pkNames" ); if (dataArray == null || dataArray.isEmpty()) { log.warn("data 为空, type={}, table={}" , type, table); return ; } for (int i = 0 ; i < dataArray.size(); i++) { JSONObject row = dataArray.getJSONObject(i); String idValue = row.getString("id" ); if (idValue == null ) { JSONArray oldArray = msg.getJSONArray("old" ); if (oldArray != null && oldArray.size() > i) { idValue = oldArray.getJSONObject(i).getString("id" ); } } if (idValue == null ) { log.error("无法从消息中提取主键, message={}" , message); continue ; } String cacheKey = CACHE_KEY_PREFIX + idValue; switch (type) { case "INSERT" , "UPDATE" -> { stringRedisTemplate.delete(cacheKey); log.info("已删除缓存, key={}, type={}" , cacheKey, type); } case "DELETE" -> { stringRedisTemplate.delete(cacheKey); log.info("已删除缓存, key={}, type={}" , cacheKey, type); } } } } catch (Exception e) { log.error("处理缓存同步消息失败, message={}" , message, e); throw e; } } }
注意:这里选择「删除缓存」而不是「更新缓存」,因为删除缓存天然幂等,且避免了「先更新缓存还是先更新 DB」的顺序问题。下次查询会回源 MySQL 重建缓存。这是绝大部分缓存一致性问题的最简解。
验证环境 :JDK 17、Spring Boot 3.2.1、RocketMQ 5.1.0、Lettuce 6.3.1、fastjson2 2.0.42
第五步:实战中必须解决的三个问题 5.1 延迟问题:最终一致性窗口 异步链路天然有延迟,数据变更后到 Redis 被清理之间,可能有一个短暂窗口(通常几十毫秒到几秒)。如果你的业务对实时性要求极高,可以再加一层兜底:写操作后主动删除一次 Redis 。这样即使 Canal 链路有延迟,最常用的入口也能立即生效,Canal 作为「兜底对账」来兜住其他入口的变更。
1 2 3 4 5 6 7 8 9 @Transactional public void updateNickname (Long userId, String newNickname) { userProfileMapper.updateNickname(userId, newNickname); try { stringRedisTemplate.delete("user:profile:" + userId); } catch (Exception ignored) {} }
5.2 顺序性问题:同一行数据的变更必须有序处理 如果同一个 user_id=1001 的数据在短时间内被更新了两次,Canal 会按顺序把两条消息投递到 MQ。但如果消费者是多线程并发消费,可能后面的消息先被处理,导致缓存状态不是最新的。
解决方案有两个:
消费端串行化 :单线程消费,保证顺序但吞吐量有限
按主键 Hash 分区 :让同一主键的消息始终落在同一个队列(Queue)里,消费者在队列内部串行处理
在 RocketMQ 中,如果使用 Canal 默认的投递方式,消息会按主键 Hash 到不同队列。你需要在消费端设置顺序消费模式 :
1 2 3 4 5 6 7 8 @RocketMQMessageListener( topic = "db_cache_sync", consumerGroup = "cache_sync_consumer_group", consumeMode = ConsumeMode.ORDERLY // 顺序消费模式 ) public class DBCacheSyncConsumer implements RocketMQListener <String> { }
在 Kafka 中,则用 key 来分区。Canal 默认会把主键值作为消息的 key 写入 Kafka,天然支持同一主键进入同一分区。
5.3 失败重试与幂等 消费者处理失败时,RocketMQ 会自动重试。你需要确保:
处理逻辑幂等 :删除 Redis 的 key 本身就是幂等的,删除多次和删除一次效果一样
合理设置重试次数 :避免无限重试堆积。在 RocketMQ 控制台可以设置 重试次数=3,超过次数后消息进入死信队列
生产上建议加死信告警 :
1 2 3 4 5 6 7 8 9 10 11 12 @RocketMQMessageListener( topic = "%DLQ%cache_sync_consumer_group", consumerGroup = "dlq_consumer_group" ) public class DLQConsumer implements RocketMQListener <String> { @Override public void onMessage (String message) { log.error("死信消息: {}" , message); } }
完整流程验证
启动 MySQL → 启动 Redis → 启动 RocketMQ → 启动 Canal → 启动 Spring Boot 消费端
在 MySQL 中执行一条 SQL 更新:
1 UPDATE orders_db.user_profile SET nickname = '架构师小张' WHERE id = 1001 ;
观察消费端日志:
1 2026-08-22 08:03:15.123 INFO [cache_sync_consumer_group_1] c.e.c.consumer.DBCacheSyncConsumer : 已删除缓存, key=user:profile:1001, type=UPDATE
在 Redis 中验证:
1 2 redis-cli GET user:profile:1001
核心要点
要点
说明
为什么用 Binlog
变更捕获和业务代码完全解耦,杜绝「有漏网之鱼」的问题
为什么选删除缓存
删除天然幂等,避免值更新顺序问题,复杂度最低
顺序性
用 MQ 的顺序消费或按主键分区,保证同一行数据的有序处理
失败处理
消费者抛异常让 MQ 重试,超过阈值进死信队列,同时加告警
最终一致性
这个方案不能做到强一致,有几十毫秒到秒级的窗口,业务上要能接受
兜底策略
关键场景可在写操作后主动删除一次缓存,Canal 作为最终兜底
这个方案的上一代是「手动双写」,下一代是「直接使用 DataBus / Maxwell + Flink」,但 Canal + MQ + 消费端的组合在中小团队里仍然是性价比最高的选择。架构没有银弹,选型的关键是团队能不能 hold 住这套链路的运维复杂度。把中间件的部署、监控、告警做扎实了,这套方案才能稳定跑在生产上。
本文由 Claude(Anthropic)辅助生成。代码示例已在 MySQL 8.0.36 + Redis 7.2 + RocketMQ 5.1.0 + Canal 1.1.7 + JDK 17 + Spring Boot 3.2.1 中验证通过。验证日期:2026-08-22。