Canal + RocketMQ监听MySQL环境搭建
我们在项目中会引入缓存来提升系统性能。当一个请求进来,如果能够命中缓存,那么就直接返回缓存中的结果,不会去查询数据库。但引入缓存的同时,也就引入了与数据库之间不一致问题。 比如一个数据“AAA”存入了缓存,设置30分钟之后失效。那么这30分钟之内,查询到这条数据的请求拿到的就是“AAA”。但如果在第15分钟时,用户更新了这条数据,将它变成了“BBB”,这就导致了有15分钟的时间缓存与数据库之间不一致。 有一个方案能够解决这个问题:Canal监听MySQL的binlog,有数据变动就通过MQ发送数据更新的消息。后端接收到数据更新的消息,就可以及时的更新或者删除缓存,这样就避免了长时间缓存与数据库不一致。
版本说明
MySQL 8.0.35 Canal 1.1.8 RocketMQ 4.9.7
搭建步骤
这里牵扯到三方面内容,MySQL、Canal和RocketMQ。如果你是第一次搭建这个环境,那么我建议按照下面的步骤进行:
1 调整MySQL
2 安装并测试Canal的tcp模式
上述两步参考这篇文章:聚合搜索平台-SpringBoot整合Canal。
3 安装并测试RocketMQ
这一步网上的教程比较多,这里不再赘述。 这里要注意这样一个问题:RocketMQ是基于JDK8写的,并且Canal给MQ发消息时也需要调用MQ的jar包。所以如果你用的是高版本的JDK,那么可以修改Canal/RocketMQ启动文件中有关java的命令。这里以Canal启动文件为例。
▼text复制代码Windows: "D:\xxxxx\jdk1.8\bin\java.exe" %JAVA_OPTS% -classpath "%CLASSPATH%"……
▼text复制代码Linux: ## set java path if [ -z "$JAVA" ] ; then JAVA=/.../jdk8/java/bin/java fi ALIBABA_JAVA=/.../jdk8/java/bin/java TAOBAO_JAVA=/.../jdk8/java/bin/java
Canal与RocketMQ进行整合
修改Canal配置文件
canal.properties文件
▼text复制代码# 注意这个有的版本是RocketMQ,有的版本是rocketMQ,仔细看一下注释 canal.serverMode = rocketMQ rocketmq.producer.group = 你的groupName rocketmq.namesrv.addr = ip:端口 canal.mq.flatMessage = false canal.instance.filter.transaction.entry = true
example\instance.properties文件
▼text复制代码canal.mq.topic=你的topic # 这里稍后解释 canal.instance.filter.regex=数据库名\.表名1.*,数据库名\.表名.* canal.mq.dynamicTopic=true
关于canal.instance.filter.regex参数,有个大坑。我在搭这个环境的时候,tcp模式是没有问题的,但是切换到MQ模式之后,Canal就不发消息了,也不记录任何日志。这个bug我改了三天。 我在今天,先是把canal.mq.flatMessage改成了false。然后发现Canal发消息给MQ了,但是只发送事务开启和结束的消息EntryType.TRANSACTIONBEGIN和EntryType.TRANSACTIONEND。没有任何与数据变动有关的消息。 后来我也是在网上查到一篇文章,写的是canal.instance.filter.regex参数设置的不对的话,是会把所有表的数据变动全部过滤掉的(没错上面的bug基本上就是这个原因)。 那么tcp模式为什么是好的呢?如果你仔细观察Canal官方的示例代码,你会发现:
▼text复制代码connector.subscribe(".*\\..*");
我对于这个的理解就是,配置文件里的正则表达式,读取到程序里,可能就不是那个样子了。但是代码里面写的就不会有问题。tcp模式里有上面这行代码,就把问题给掩盖了。 这个正则表达式(.\..)的意思就是监听所有库的所有表。但是在instance.properties这样写就不对,反而会忽略掉所有表,只发事务开启和关闭的消息。我也试过(...)和(.*),都不行。一个“\”是因为网上说多一个\是为了转义,但试过还是不行。
我建议就是(数据库名.表名1.,数据库名.表名.),这样写就正常
后端代码测试
▼text复制代码import com.alibaba.otter.canal.client.CanalMessageDeserializer; import com.alibaba.otter.canal.protocol.CanalEntry; import com.alibaba.otter.canal.protocol.Message; import org.apache.rocketmq.common.message.MessageExt; import org.apache.rocketmq.spring.annotation.RocketMQMessageListener; import org.apache.rocketmq.spring.core.RocketMQListener; import org.springframework.stereotype.Component; import java.util.List; @Component @RocketMQMessageListener(topic = "topic名", consumerGroup = "group名") public class MessageConsumer implements RocketMQListener<MessageExt> { @Override public void onMessage(MessageExt message){ Message msg = CanalMessageDeserializer.deserializer(message.getBody()); List<CanalEntry.Entry> entries = msg.getEntries(); printEntry(entries); } private static void printEntry(List<CanalEntry.Entry> entrys) { for (CanalEntry.Entry entry : entrys) { if (entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONBEGIN || entry.getEntryType() == CanalEntry.EntryType.TRANSACTIONEND) { continue; } CanalEntry.RowChange rowChage = null; try { rowChage = CanalEntry.RowChange.parseFrom(entry.getStoreValue()); } catch (Exception e) { throw new RuntimeException("ERROR ## parser of eromanga-event has an error , data:" + entry.toString(), e); } CanalEntry.EventType eventType = rowChage.getEventType(); System.out.println(String.format("================> binlog[%s:%s] , name[%s,%s] , eventType : %s", entry.getHeader().getLogfileName(), entry.getHeader().getLogfileOffset(), entry.getHeader().getSchemaName(), entry.getHeader().getTableName(), eventType)); for (CanalEntry.RowData rowData : rowChage.getRowDatasList()) { if (eventType == CanalEntry.EventType.DELETE) { printColumn(rowData.getBeforeColumnsList()); } else if (eventType == CanalEntry.EventType.INSERT) { printColumn(rowData.getAfterColumnsList()); } else { System.out.println("-------> before"); List<CanalEntry.Column> beforeColumnsList = rowData.getBeforeColumnsList(); printColumn(beforeColumnsList); System.out.println("-------> after"); printColumn(rowData.getAfterColumnsList()); } } } } private static void printColumn(List<CanalEntry.Column> columns) { for (CanalEntry.Column column : columns) { System.out.println(column.getName() + " : " + column.getValue() + " update=" + column.getUpdated()); } } }
当你看到类似下面这样的输出,就表示成功了
▼text复制代码================> binlog[xxx-bin.xxx:xxx] , name[数据库名,表名] , eventType : UPDATE (或者 INSERT 或者 DELETE) -------> before id : xxx update=false createTime : 2025-09-02 17:40:13 update=false editTime : 2025-09-02 17:40:13 update=false updateTime : 2025-09-13 01:12:21 update=false isDelete : 0 update=false -------> after id : xxx update=false createTime : 2025-09-02 17:40:13 update=false editTime : 2025-09-02 17:40:13 update=false updateTime : 2025-09-13 02:39:58 update=true isDelete : 0 update=false
