欢迎光临
我们一直在努力

使用canal解决数据缓存不一致问题实操

一:一般的数据缓存一致性解决方案

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 落地缓存一致性的全部实操内容,实测可用,感兴趣可以点个收藏。

赞(0)
未经允许不得转载:171主机测评 » 使用canal解决数据缓存不一致问题实操
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址