MQ专题:延迟消息的通用方案

2024-08-30 09:44

本文主要是介绍MQ专题:延迟消息的通用方案,希望对大家解决编程问题提供一定的参考价值,需要的开发者们随着小编来一起学习吧!

一、主要内容

本文将实现一个MQ延迟消息的通用方案。

方案不依赖于MQ中间件,依靠MySQL和DelayQueue解决,不管大家用的是什么MQ,具体是RocketMQ、RabbitMQ还是kafka,本文这个方案你都可以拿去直接使用,可以轻松实现任意时间的延迟消息投递。

二、涉及技术点

  1. SpringBoot2.7
  2. MyBatisPlus
  3. MySQL
  4. 线程池
  5. java中的延迟队列:DelayQueue
  6. 分布式锁

三、延迟消息常见的使用场景

  1. 订单超时处理

比如下单后15分钟,未支付,则自动取消订单,回退库存。

可以采用延迟队列实现:下单的时候可以投递一条延迟15分钟的消息,15分钟后消息将被消费。

  1. 消息消费失败重试

比如MQ消息消费失败后,可以延迟一段时间再次消费。

可以采用延迟消息实现:消费失败,可以投递一条延迟消息,触发再次消费

  1. 其他任意需要延迟处理的业务

业务中需要延迟处理的场景,都可以使用延迟消息来搞定。

四、延迟消息常见的实现方案

方案1:MySQL + job定时轮询

由于延迟消息的时间不确定,若要达到实时性很高的效果,也就是说消息的延迟时间是不知道的,那就需要轮询每一秒才能确保消息在指定的延迟时间被处理,这就要求job需要每秒查询一次db中待投递的消息。

这种方案访问db的频率比较高,对数据库造成了一定的压力。

方案2:RabbitMQ 中的TTL+死信队列

rabbitmq中可以使用TTL消息 + 死信队列实现,也可以安装延时插件。

此方案对中间件有依赖,不同的MQ实现是不一样的,若换成其他的MQ,方案要重新实现

方案3:MySQL + job定时轮询 + DelayQueue

可以对方案1进行改进,引入java中的 DelayQueue。

job可以采用1分钟执行一次,每次拉取未来2分钟内需要投递的消息,将其丢到java自带的 DelayQueue 这个延迟队列工具类中去处理,这样便能做到实时性很高的投递效果,且对db的压力也降低了很多。

这种方案对db也没什么压力,实时性非常高,且对MQ没有依赖,这样不管切换什么MQ,这种方案都不需要改动。

本文将落地该方案。

需要一张本地消息表(t_msg)

这张表用来暂存事务消息和延迟消息

create table if not exists t_msg
(id               varchar(32) not null primary key comment '消息id',body_json        text        not null comment '消息体,json格式',status           smallint    not null default 0 comment '消息状态,0:待投递到mq,1:投递成功,2:投递失败',expect_send_time datetime    not null comment '消息期望投递时间,大于当前时间,则为延迟消息,否则会立即投递',actual_send_time datetime comment '消息实际投递时间',create_time      datetime comment '创建时间',fail_msg         text comment 'status=2 时,记录消息投递失败的原因',fail_count       int         not null default 0 comment '已投递失败次数',send_retry       smallint    not null default 1 comment '投递MQ失败了,是否还需要重试?1:是,0:否',next_retry_time  datetime comment '投递失败后,下次重试时间',update_time      datetime comment '最近更新时间',key idx_status (status)
) comment '本地消息表';

五、代码落地

将延迟消息保存到t_msg

在这里插入图片描述

投递事务消息 or 2分钟内的延时消息

在这里插入图片描述

每分钟执行一次Job,查出2分钟内应该被投递的延迟消息 和 2分钟内应该重新投递的上次投递失败的事务消息,再放到延时队列里

    /*** 每分钟执行一次*/@Scheduled(fixedDelay = 1, timeUnit = TimeUnit.MINUTES)public void sendRetry() {/*** 查询出需要重试的消息(状态为0 and 期望投递时间 <= 当前时间 + 2分钟) || (投递失败的 and 需要重试 and 下次重试时间小于等于当前时间 + 2分钟)* select * from t_msg where ((status = 0 and expect_send_time<=当前时间+2分钟) or (status = 2 and send_retry = 1 and next_retry_time <= 当前时间 + 2分钟))*/LocalDateTime time = LocalDateTime.now().plusMinutes(2);LambdaQueryWrapper<MsgPO> qw = Wrappers.lambdaQuery(MsgPO.class).and(query -> query.and(lq ->lq.eq(MsgPO::getStatus, MsgStatusEnum.INIT.getStatus()).le(MsgPO::getExpectSendTime, time)).or(lq -> lq.eq(MsgPO::getStatus, MsgStatusEnum.FAIL.getStatus()).eq(MsgPO::getSendRetry, 1).le(MsgPO::getNextRetryTime, time)));qw.orderByAsc(MsgPO::getId);//先获取最小的一条记录的idMsgPO minMsgPo = this.msgService.findOne(qw);if (minMsgPo == null) {return;}this.msgSender.sendRetry(minMsgPo);String minMsgId = minMsgPo.getId();//循环中继续向后找出id>minMsgId的所有记录,然后投递重试while (true) {//select * from t_msg where ((status = 0 and expect_send_time<=当前时间+2分钟) or (status = 2 and send_retry = 1 and next_retry_time <= 当前时间 + 2分钟)) and id>#{minMsgId}qw = Wrappers.lambdaQuery(MsgPO.class).and(query -> query.and(lq ->lq.eq(MsgPO::getStatus, MsgStatusEnum.INIT.getStatus()).le(MsgPO::getExpectSendTime, time)).or(lq -> lq.eq(MsgPO::getStatus, MsgStatusEnum.FAIL.getStatus()).eq(MsgPO::getSendRetry, 1).le(MsgPO::getNextRetryTime, time)));qw.gt(MsgPO::getId, minMsgId);qw.orderByAsc(MsgPO::getId);Page<MsgPO> page = new Page<>();page.setCurrent(1);page.setSize(500);this.msgService.page(page, qw);//如果查询出来的为空 || 当前服务要停止了(stop=true),则退出循环if (CollUtils.isEmpty(page.getRecords()) || stop) {break;}//投递重试for (MsgPO msgPO : page.getRecords()) {this.msgSender.sendRetry(msgPO);}// minMsgId = 当前列表最后一条消息的idminMsgId = page.getRecords().get(page.getRecords().size() - 1).getId();}}

在这里插入图片描述

这篇关于MQ专题:延迟消息的通用方案的文章就介绍到这儿,希望我们推荐的文章对编程师们有所帮助!



http://www.chinasem.cn/article/1120462

相关文章

Redis客户端连接机制的实现方案

《Redis客户端连接机制的实现方案》本文主要介绍了Redis客户端连接机制的实现方案,包括事件驱动模型、非阻塞I/O处理、连接池应用及配置优化,具有一定的参考价值,感兴趣的可以了解一下... 目录1. Redis连接模型概述2. 连接建立过程详解2.1 连php接初始化流程2.2 关键配置参数3. 最大连

RabbitMQ消息总线方式刷新配置服务全过程

《RabbitMQ消息总线方式刷新配置服务全过程》SpringCloudBus通过消息总线与MQ实现微服务配置统一刷新,结合GitWebhooks自动触发更新,避免手动重启,提升效率与可靠性,适用于配... 目录前言介绍环境准备代码示例测试验证总结前言介绍在微服务架构中,为了更方便的向微服务实例广播消息,

MyBatis-Plus通用中等、大量数据分批查询和处理方法

《MyBatis-Plus通用中等、大量数据分批查询和处理方法》文章介绍MyBatis-Plus分页查询处理,通过函数式接口与Lambda表达式实现通用逻辑,方法抽象但功能强大,建议扩展分批处理及流式... 目录函数式接口获取分页数据接口数据处理接口通用逻辑工具类使用方法简单查询自定义查询方法总结函数式接口

MySQL 迁移至 Doris 最佳实践方案(最新整理)

《MySQL迁移至Doris最佳实践方案(最新整理)》本文将深入剖析三种经过实践验证的MySQL迁移至Doris的最佳方案,涵盖全量迁移、增量同步、混合迁移以及基于CDC(ChangeData... 目录一、China编程JDBC Catalog 联邦查询方案(适合跨库实时查询)1. 方案概述2. 环境要求3.

SpringBoot3.X 整合 MinIO 存储原生方案

《SpringBoot3.X整合MinIO存储原生方案》本文详细介绍了SpringBoot3.X整合MinIO的原生方案,从环境搭建到核心功能实现,涵盖了文件上传、下载、删除等常用操作,并补充了... 目录SpringBoot3.X整合MinIO存储原生方案:从环境搭建到实战开发一、前言:为什么选择MinI

Python通用唯一标识符模块uuid使用案例详解

《Python通用唯一标识符模块uuid使用案例详解》Pythonuuid模块用于生成128位全局唯一标识符,支持UUID1-5版本,适用于分布式系统、数据库主键等场景,需注意隐私、碰撞概率及存储优... 目录简介核心功能1. UUID版本2. UUID属性3. 命名空间使用场景1. 生成唯一标识符2. 数

Knife4j+Axios+Redis前后端分离架构下的 API 管理与会话方案(最新推荐)

《Knife4j+Axios+Redis前后端分离架构下的API管理与会话方案(最新推荐)》本文主要介绍了Swagger与Knife4j的配置要点、前后端对接方法以及分布式Session实现原理,... 目录一、Swagger 与 Knife4j 的深度理解及配置要点Knife4j 配置关键要点1.Spri

SQLite3 在嵌入式C环境中存储音频/视频文件的最优方案

《SQLite3在嵌入式C环境中存储音频/视频文件的最优方案》本文探讨了SQLite3在嵌入式C环境中存储音视频文件的优化方案,推荐采用文件路径存储结合元数据管理,兼顾效率与资源限制,小文件可使用B... 目录SQLite3 在嵌入式C环境中存储音频/视频文件的专业方案一、存储策略选择1. 直接存储 vs

java向微信服务号发送消息的完整步骤实例

《java向微信服务号发送消息的完整步骤实例》:本文主要介绍java向微信服务号发送消息的相关资料,包括申请测试号获取appID/appsecret、关注公众号获取openID、配置消息模板及代码... 目录步骤1. 申请测试系统2. 公众号账号信息3. 关注测试号二维码4. 消息模板接口5. Java测试

SpringBoot服务获取Pod当前IP的两种方案

《SpringBoot服务获取Pod当前IP的两种方案》在Kubernetes集群中,SpringBoot服务获取Pod当前IP的方案主要有两种,通过环境变量注入或通过Java代码动态获取网络接口IP... 目录方案一:通过 Kubernetes Downward API 注入环境变量原理步骤方案二:通过