一:一般的数据缓存一致性解决方案
1.旁路缓存:
先更新数据库,然后直接删除缓存(优点:实现简单,命中率高,适合并发不高且数据一致性要求不高的场景)
2.延时双删:
删除缓存,更新数据库,sleep(时间要很短),再删缓存。或者先更新数据库再删缓存sleep之后再删一次缓存(优点:适合写操作频繁的场景强一致性,但sleep时间不好控制,sleep过长会拉长响应时间)
3.订阅 Binlog 更新缓存:
使用canal等类似的工具,数据库除select之外的变更通过canal监听binlog日志,发送到消息队列。消费者拿到消息再去删除或更新缓存(优点:跟项目解耦,适合高并发场景,实现简单,即使更新缓存失败也可以通过重试机制最终保持一致,但需要引入mq,canal等组件,同时延时稍高(毫秒级到秒级))
4.Write Behind Caching(写回):
写请求只更新缓存,立即返回成功。然后异步批量地将缓存中的变更刷新到数据库。(写入性能极高(内存操作),适合日志、计数等允许丢失少量数据的场景,所以适合数据一致性弱的场景,同时会有较大的缓存不一致的窗口期)
5.Read/Write Through(读穿/写穿):
把缓存作为主要数据源,应用只和缓存交互,缓存负责同步读写数据库。一般由缓存中间件(如 Apache Ignite、Hazelcast)或代理层实现。(应用层完全不用关心数据库,逻辑简单,但缓存中间件需要支持该特性;写操作性能受限于数据库。)
二:选型与准备工作
为演示canal的使用选择方案三
1.canal的安装
因为canal是使用java语言写的,所以需要java环境,java8,java11都可以,不演示java环境的安装了,下载canal安装包,推荐使用docker拉取,这里使用windows安装1.1.8稳定版本,地址:https://github.com/alibaba/canal/releases
解压之后你会看到:
canal.deployer-1.1.8/ ├── bin/ ← 启动和停止的脚本都在这里 ├── conf/ ← 所有配置文件都在这里 │ ├── canal.properties ← Canal 服务总配置 │ └── example/ ← 这个目录代表一个监听实例 │ └── instance.properties ← 实例专属配置 └── logs/ ← 日志文件存放位置(启动后会自动生成)
2.mysql配置
需要开启binlog(Canal 全靠它才能监听到数据变化,一般是默认开启的)mysql中使用这俩个命令检查
— 检查配置是否生效,log_bin 的值应为 ON SHOW VARIABLES LIKE 'log_bin';
— 检查格式是否为 ROW,确保 ROW 模式已生效 SHOW VARIABLES LIKE 'binlog_format';
如果没开启:在 MySQL 配置文件(Windows 的 my.ini 或 Linux 的 /etc/my.cnf)中添加以下内容:
[mysqld] log-bin=mysql-bin # 开启 binlog binlog-format=ROW # 必须设置为 ROW 模式 server-id=1 # 服务器唯一 ID,不能为 0
添加后,重启 MySQL 服务。然后在 MySQL 命令行中执行上面两条命令,确认 binlog 已开启
3.修改配置文件(重点)
我们这里使用rabbitmq来做消息队列,简单安装一下,我往期文章有讲安装,感兴趣可以看一下
1.canal.properties
# 选择消息中间件类型
canal.serverMode = rabbitMQ #文件第29行 canal.mq.flatMessage = true #文件第127行
#下拉找到rabbitmq栏目 ,配置最好跟我一样
rabbitmq.host = 192.168.200.*** #自己配置就好
rabbitmq.virtual.host =
rabbitmq.exchange = canal.exchange
rabbitmq.username = admin #记得创建用户 guest是默认不让
rabbitmq.password = 123456
rabbitmq.queue = #不用填,再example实例中填写
rabbitmq.routingKey =
rabbitmq.deliveryMode = 2 #必填
2. instance.properties(具体表的MQ配置)
# username/password
canal.instance.dbUsername=root
canal.instance.dbPassword=1234567
canal.instance.connectionCharset = UTF-8
# enable druid Decrypt database password
canal.instance.enableDruid=false
# table regex 用的是哪个数据库中哪张表,这里用黑马点评这个项目中的表吧,写法有具体规则,ai一下
canal.instance.filter.regex=hmdp\\\\..*
#这里我们用固定路由
canal.mq.topic=canal.routing.key
# canal.mq.dynamicTopic= #注意不配置这个,固定路由与自动路由同时配置,以自动路由为准
4.rabbitmq中添加交换机以及队列:
1.交换机

2.添加队列

3.绑定:

4.测试
4.1选择刚才建好的交换机发布消息:

4.2选择刚才新建的队列点击getmessage看是否接收消息

5.启动canal测试
1.启动
canal的bin文件夹下双击startup.bat,查看终端或者 log文件下的canal文件下的日志文件是否有报错

2.测试
去监听的数据库随便修改一条数据,回到rabbitmq中看是否监听到数据,

3.数据格式以及必要解释如下:
{
//修改的数据
"data": [
{
"id": "1",
"name": "999999",
"type_id": "1",
"images": "https://qcloud.dpfile.com/pc/jiclIsCKmOI2arxKN1Uf0Hx3PucIJH8q0QSz-Z8llzcN56-_QiKuOvyio1OOxsRtFoXqu0G3iT2T27qat3WhLVEuLYk00OmSS1IdNpm8K8sG4JN9RIm2mTKcbLtc2o2vfCF2ubeXzk49OsGrXt_KYDCngOyCwZK-s3fqawWswzk.jpg,https://qcloud.dpfile.com/pc/IOf6VX3qaBgFXFVgp75w-KKJmWZjFc8GXDU8g9bQC6YGCpAmG00QbfT4vCCBj7njuzFvxlbkWx5uwqY2qcjixFEuLYk00OmSS1IdNpm8K8sG4JN9RIm2mTKcbLtc2o2vmIU_8ZGOT1OjpJmLxG6urQ.jpg",
"area": "大关",
"address": "金华路锦昌文华苑29号",
"x": "120.149192",
"y": "30.316078",
"avg_price": "80",
"sold": "4215",
"comments": "3035",
"score": "37",
"open_hours": "10:00-22:00",
"create_time": "2021-12-22 18:10:39",
"update_time": "2026-06-04 14:45:40"
}
],
//数据库
"database": "hmdp",
//数据戳
"es": 1780555540000,
"gtid": "",
"id": 2,
"isDdl": false,
//表结构
"mysqlType": {
"id": "bigint unsigned",
"name": "varchar(128)",
"type_id": "bigint unsigned",
"images": "varchar(1024)",
"area": "varchar(128)",
"address": "varchar(255)",
"x": "double unsigned",
"y": "double unsigned",
"avg_price": "bigint unsigned",
"sold": "int(10) unsigned zerofill",
"comments": "int(10) unsigned zerofill",
"score": "int(2) unsigned zerofill",
"open_hours": "varchar(32)",
"create_time": "timestamp",
"update_time": "timestamp"
},
//原先的数据(修改前的)
"old": [
{
"name": "888888",
"update_time": "2026-05-29 12:05:42"
}
],
"pkNames": [
"id"
],
"sql": "",
"sqlType": {
"id": -5,
"name": 12,
"type_id": -5,
"images": 12,
"area": 12,
"address": 12,
"x": 8,
"y": 8,
"avg_price": -5,
"sold": 4,
"comments": 4,
"score": 4,
"open_hours": 12,
"create_time": 93,
"update_time": 93
},
"table": "tb_shop",
"ts": 1780555540701,
//操作类型
"type": "UPDATE"
}
三:项目中代码实现
1.实体类,定义一个数据结构来解析 Canal 推送的 JSON 消息也就是上面这个
package com.hmdp.dto;
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
import lombok.Data;
import java.util.List;
import java.util.Map;
@Data
@JsonIgnoreProperties(ignoreUnknown = true) //加上这个注解去掉没有配置属性的消息
public class CanalMessage {
private String database;
private String table;
private String type; // INSERT, UPDATE, DELETE
private Long ts;
private List<Map<String, Object>> data; // 变更后的数据
private List<Map<String, Object>> old; // 变更前的数据(UPDATE 时存在)
}
2.消费者
package com.hmdp.mq;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.hmdp.dto.CanalMessage;
import com.hmdp.utils.CacheUtils;
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.Exchange;
import org.springframework.amqp.rabbit.annotation.Queue;
import org.springframework.amqp.rabbit.annotation.QueueBinding;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.amqp.core.ExchangeTypes;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.stereotype.Component;
import java.util.Map;
@Slf4j
@Component
public class CanalCacheConsumer {
@Autowired
private CacheUtils cacheUtils;
private final ObjectMapper objectMapper = new ObjectMapper();
//定义交换机、队列、以及绑定
@RabbitListener(bindings = @QueueBinding(
value = @Queue(value = "canal.queue", durable = "true"),
exchange = @Exchange(value = "canal.exchange", type = ExchangeTypes.TOPIC),
key = "canal.routing.key"
))
public void handleCanalMessage(String messageJson) {
try {
// 1. 先打印原始消息,便于调试
log.info("收到原始消息: {}", messageJson);
// 2. 反序列化
CanalMessage message = objectMapper.readValue(messageJson, CanalMessage.class);
log.info("解析成功: database={}, table={}, type={}",
message.getDatabase(), message.getTable(), message.getType());
String table = message.getTable();
String type = message.getType();
if (!"INSERT".equals(type) && !"UPDATE".equals(type) && !"DELETE".equals(type)) {
return;
}
if ("tb_shop".equals(table)) {
for (Map<String, Object> row : message.getData()) {
Object idObj = row.get("id");
if (idObj != null) {
Long shopId = Long.valueOf(idObj.toString());
cacheUtils.evictShopCache(shopId); //修改或者删除缓存
}
}
}
} catch (Exception e) {
log.error("处理 Canal 消息失败,原始消息: {}", messageJson, e);
}
}
}
3.删除缓存代码:
这里代表性的删除本地caffeine缓存以及redis,具体看业务
package com.hmdp.utils;
import com.github.benmanes.caffeine.cache.Cache;
import com.hmdp.entity.Shop;
import lombok.extern.slf4j.Slf4j;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.data.redis.core.StringRedisTemplate;
import org.springframework.stereotype.Component;
import static com.hmdp.utils.RedisConstants.*;
@Slf4j
@Component
public class CacheUtils {
@Autowired
private StringRedisTemplate stringRedisTemplate;
@Autowired
private Cache<Long, Shop> caffeineCache; // Caffeine 实例
/**
* 清除店铺缓存(Redis + Caffeine)
*/
public void evictShopCache(Long shopId) {
String redisKey = CACHE_SHOP_KEY + shopId;
stringRedisTemplate.delete(redisKey);
caffeineCache.invalidate(redisKey);
log.info("通过canal清除店铺缓存,shopId: {}", shopId);
}
}
4.代码测试:
1.原数据库:

2.原redis中缓存:

3.修改过后:

4.控制台打印:

5.看redis缓存是否删除

四:总结
安装canal、rabbitmq,确保mysql开启binlog,配置canal.properties以及example中instance.properties文件,添加交换机队列及绑定关系,重启canal,项目中添加消费者消费。
以上就是 Canal 配合 RabbitMQ 落地缓存一致性的全部实操内容,实测可用,感兴趣可以点个收藏。

