SpringBoot集成阿里云RocketMQ 实现延时队列
mhr18 2024-12-09 12:15 17 浏览 0 评论
前言
消息队列RocketMQ版是阿里云基于Apache RocketMQ构建的低延迟、高并发、高可用、高可靠的分布式消息中间件。消息队列RocketMQ版既可为分布式应用系统提供异步解耦和削峰填谷的能力,同时也具备互联网应用所需的海量消息堆积、高吞吐、可靠重试等特性。具体详情参考 消息队列RocketMQ版 阿里云官网介绍。
在这一篇里,我们将介绍SpringBoot如何集成阿里云版RocketMQ,通过本篇的学习,你将收获:
- 什么是延时队列
- 延时队列在设班排课领域内的使用场景
- SpringBoot整合RocketMQ
什么是延时队列
顾名思义,延迟队列是指被延迟消费的队列。首先,队列意味着内部的元素是有序的,元素出队和入队是有方向性的。 普通的队列,消息一旦入队后就会被消费者马上消费。而延时队列,最重要的特性就体现在它的延时属性上,跟普通的队列不一样的是,延时队列中的元素则是希望被在指定时间取出和处理,可以说延时队列中的元素是都是带时间属性的。简单来说,延时队列就是用来存放需要在指定时间被处理的元素的队列。
延时队列使用场景
延时队列的运用场景比较多,在设班排课领域有如下场景:
- 班级某节课到达开课时间后,不允许插班,需要变更班级状态
- 预定未来教室后,到达预定教室使用时间之前,给预订者发提醒消息
- 商品过售卖期后自动下架
大家应该都发现了这些场景都有一个共同特点,需要在指定的时间点完成某一项任务。这个时候就有人会说了,这些功能可以使用定时任务轮询数据库,固定时间周期取出所需的数据进行处理。如果数据量比较少,确实可以这样做;如果对数据实时性不是有严格的限制,那么每小时执行下定时任务出来数据,确实也是一个可行的方案。但是对于数据量比较大,并且时效性较强的场景,这种轮询方式很难在一个周期内处理完业务数据,这就会给系统和数据库带来很大压力,无法满足业务要求而且性能低下。
这时候,延时队列就可以闪亮登场了,以上场景,正是延时队列的用武之地。延迟队列的实现方式比较多,比如DelayQueue延时队列、RabbitMQ队列、Redis键通知或Redis zset等等。
阿里云版RocketMQ 支持任意时间等级延时,简单易用并省去运维工作,所以本文将介绍如何用阿里云版RocketMQ来实现延时队列。
首先需要在阿里云开通RocketMQ服务,具体请参考阿里云官方文档
下面讲解SpringBoot整合RocketMQ来实现延时队列的,话不多说,上代码
- 引入依赖包
<!-- aliyun business version rocketmq-->
<dependency>
<groupId>com.aliyun.openservices</groupId>
<artifactId>ons-client</artifactId>
<version>1.8.4.Final</version>
</dependency>
- yml配置文件添加属性
nepenthe: #服务名,替换成自己的
rocketmq:
accessKey: #连接AccessKey id ,替换成自己的
secretKey: #连接AccessKey secret,替换成自己的
nameSrvAddr: #生产者ons接入域名,替换成自己的
delayTopic: #生产主题,替换成自己的
groupId: #生产者id(旧版本是生产者id,新版本是groupid),替换成自己的
timeoutMillis: #发送过期时间,自定义
- 封装RocketMQ配置文件
@Configuration
public class RocketMqConfig {
/**
*鉴权需要的AccessKey ID
*/
@Value("${nepenthe.rocketmq.accessKey}")
private String accessKey;
/**
*鉴权需要的AccessKey Secret
*/
@Value("${nepenthe.rocketmq.secretKey}")
private String secretKey;
/**
* 实例TCP 协议公网接入地址(实际项目,填写自己阿里云MQ的公网地址)
*/
@Value("${nepenthe.rocketmq.nameSrvAddr}")
private String nameSrvAddr;
/**
* 延时队列topic
*/
@Value("${nepenthe.rocketmq.delayTopic}")
private String delayTopic;
/**
* 延时队列group
*/
@Value("${nepenthe.rocketmq.groupId}")
private String groupId;
/**
* 设置发送超时时间
*/
@Value("${nepenthe.rocketmq.timeoutMillis}")
private String timeoutMillis;
public Properties getRocketMqProperty() {
Properties properties = new Properties();
properties.setProperty(PropertyKeyConst.GROUP_ID,this.getGroupId());
properties.setProperty(PropertyKeyConst.AccessKey, this.accessKey);
properties.setProperty(PropertyKeyConst.SecretKey, this.secretKey);
properties.setProperty(PropertyKeyConst.NAMESRV_ADDR, this.nameSrvAddr);
properties.setProperty(PropertyKeyConst.SendMsgTimeoutMillis, this.timeoutMillis);
return properties;
}
//TODO get set方法,自己实现
}
- 配置生产者
@Component
public class RocketMqProducerInit {
@Autowired
private RocketMqConfig mqConfig;
private static Producer producer;
@PostConstruct
public void init(){
System.out.println("启动RocketMq生产者!");
producer = ONSFactory.createProducer(mqConfig.getRocketMqProperty());
// 在发送消息前,初始化调用start方法来启动Producer,只需调用一次即可,当项目关闭时,自动shutdown
producer.start();
}
/**
* 初始化生产者
* @return
*/
public Producer getProducer(){
return producer;
}
}
- 生产者发送消息的方法
@Component
public class ProducerHandler {
private Logger logger = LoggerFactory.getLogger(ProducerHandler.class);
@Autowired
private RocketMqConfig config;
@Autowired
private RocketMqProducerInit producer;
/**
* 同步发送延时消息
* @param msgTag 标签,可用于消息小分类标注,对消息进行再归类
* @param messageBody 消息body内容,
* @param msgKey 消息key值,
* @param delayTime 服务端发送消息时间
* @return success:SendResult or error:null
*/
public SendResult sendTimeMsg(String msgTag, byte[] messageBody, String msgKey, long delayTime) {
Message msg = new Message(config.getDelayTopic(),msgTag,msgKey,messageBody);
msg.setStartDeliverTime(delayTime);
return this.send(msg);
}
/**
* 消息发送
* @param msg 消息
*/
private SendResult send(Message msg) {
try {
SendResult sendResult = producer.getProducer().send(msg);
//获取发送结果,不抛异常即发送成功
if (sendResult != null) {
logger.info("发送RocketMQ消息成功 -- Topic:{} ,msgId:{} , Key:{}, tag:{}, message:{}"
,msg.getTopic(),sendResult.getMessageId(),msg.getKey(),msg.getTag(),new String(msg.getBody()));
return sendResult;
}else {
logger.error("发送RocketMQ消息失败-- Topic:{}, Key:{}, tag:{}, message:{}"
,msg.getTopic(),msg.getKey(),msg.getTag(),new String(msg.getBody()));
return null;
}
} catch (Exception e) {
logger.error("发送RocketMQ消息失败-- Topic:{}, Key:{}, tag:{}, message:{},错误原因:{}"
,msg.getTopic(),msg.getKey(),msg.getTag(),new String(msg.getBody()),e.getMessage());
return null;
}
}
- 配置消费者
@Component
public class RocketmqConsumerInit {
@Autowired
private RocketMqConfig mqConfig;
private static Consumer consumer;
@PostConstruct
public void init(){
System.out.println("启动消费者!");
consumer = ONSFactory.createConsumer(mqConfig.getRocketMqProperty());
//监听第一个topic,new对应的监听器
consumer.subscribe(mqConfig.getDelayTopic(), "*", new DelayQueueMessageListener());
// 在发送消息前,必须调用start方法来启动consumer,只需调用一次即可,当项目关闭时,自动shutdown
consumer.start();
}
/**
* 初始化消费者
* @return
*/
public Consumer getConsumer(){
return consumer;
}
}
- 消费者消息监听
@Service
public class DelayQueueMessageListener implements MessageListener {
private Logger logger = LoggerFactory.getLogger(this.getClass());
@Override
public Action consume(Message message, ConsumeContext consumeContext) {
logger.info("接收到MQ消息 -- Topic:{}, tag:{},msgId:{} , Key:{}, body:{}",
message.getTopic(), message.getTag(), message.getMsgID(), message.getKey(), new String(message.getBody()));
try {
//TODO 处理业务
//消费成功,继续消费下一条消息
return Action.CommitMessage;
} catch (Exception e) {
logger.error("消费MQ消息失败! msgId:" + message.getMsgID() + "----ExceptionMsg:" + e.getMessage());
//消费失败,告知服务器稍后再投递这条消息,继续消费其他消息
return Action.ReconsumeLater;
}
}
}
- 运行结果
@Autowired
private ProducerHandler producerHandler;
@RequestMapping("delayMsg")
public void delayMsg2(String msg, Integer delayTime) {
JSONObject body = new JSONObject();
body.put("userId", "this is userId");
body.put("notice", "同步消息");
producerHandler.sendTimeMsg("userMessage", msg.getBytes(), "messageId", System.currentTimeMillis()+delayTime);
}
- 日志如下
2021-02-25 18:05:20.584 [http-nio-8080-exec-7] 1591467 INFO c.x.n.m.mq.rocketmq.ProducerUtil - TraceId[1614247520XAkUQkT0qv6] 发送MQ消息成功 -- Topic:auto_publish_delay_queue_topic ,msgId:0A4922E703D914DAD5DC7F7A4908000D , Key:messageId, tag:userMessage, body:我是延时队列延时10秒
2021-02-25 18:05:26.389 [http-nio-8080-exec-9] 1597272 INFO c.x.n.m.mq.rocketmq.ProducerUtil - TraceId[1614247526jymlRyaX9m6] 发送MQ消息成功 -- Topic:auto_publish_delay_queue_topic ,msgId:0A4922E703D914DAD5DC7F7A5F780010 , Key:messageId, tag:userMessage, body:我是延时队列延时12秒
2021-02-25 18:05:30.620 [ConsumeMessageThread_3] 1601503 INFO c.x.n.m.m.r.DelayQueueMessageListener - 接收到MQ消息 -- Topic:auto_publish_delay_queue_topic, tag:userMessage,msgId:bf58a99274cfadeddbec720c12d8607d , Key:messageId, body:我是延时队列延时10秒
2021-02-25 18:05:38.871 [ConsumeMessageThread_4] 1609754 INFO c.x.n.m.m.r.DelayQueueMessageListener - 接收到MQ消息 -- Topic:auto_publish_delay_queue_topic, tag:userMessage,msgId:d44270736cdfa2927194650c451584c3 , Key:messageId, body:我是延时队列延时12秒
两条消息的发送时间与消费时间,符合预期。至此,阿里云版RocketMQ实现延时队列的部分就完成了。如果考虑阿里云版RocketMQ费用,也可以考虑开源版RocketMQ,默认支持18个延时等级,或RabbitMQ加延时插件。
相关推荐
- 一文读懂Prometheus架构监控(prometheus监控哪些指标)
-
介绍Prometheus是一个系统监控和警报工具包。它是用Go编写的,由Soundcloud构建,并于2016年作为继Kubernetes之后的第二个托管项目加入云原生计算基金会(C...
- Spring Boot 3.x 新特性详解:从基础到高级实战
-
1.SpringBoot3.x简介与核心特性1.1SpringBoot3.x新特性概览SpringBoot3.x是建立在SpringFramework6.0基础上的重大版...
- 「技术分享」猪八戒基于Quartz分布式调度平台实践
-
点击原文:【技术分享】猪八戒基于Quartz分布式调度平台实践点击关注“八戒技术团队”,阅读更多技术干货1.背景介绍1.1业务场景调度任务是我们日常开发中非常经典的一个场景,我们时常会需要用到一些不...
- 14. 常用框架与工具(使用的框架)
-
本章深入解析Go生态中的核心开发框架与工具链,结合性能调优与工程化实践,提供高效开发方案。14.1Web框架(Gin,Echo)14.1.1Gin高性能实践//中间件链优化router:=...
- SpringBoot整合MyBatis-Plus:从入门到精通
-
一、MyBatis-Plus基础介绍1.1MyBatis-Plus核心概念MyBatis-Plus(简称MP)是一个MyBatis的增强工具,在MyBatis的基础上只做增强不做改变,为简化开发、提...
- Seata源码—5.全局事务的创建与返回处理
-
大纲1.Seata开启分布式事务的流程总结2.Seata生成全局事务ID的雪花算法源码3.生成xid以及对全局事务会话进行持久化的源码4.全局事务会话数据持久化的实现源码5.SeataServer创...
- Java开发200+个学习知识路线-史上最全(框架篇)
-
1.Spring框架深入SpringIOC容器:BeanFactory与ApplicationContextBean生命周期:实例化、属性填充、初始化、销毁依赖注入方式:构造器注入、Setter注...
- OpenResty 入门指南:从基础到动态路由实战
-
一、引言1.1OpenResty简介OpenResty是一款基于Nginx的高性能Web平台,通过集成Lua脚本和丰富的模块,将Nginx从静态反向代理转变为可动态编程的应用平台...
- 你还在为 Spring Boot3 分布式锁实现发愁?一文教你轻松搞定!
-
作为互联网大厂后端开发人员,在项目开发过程中,你有没有遇到过这样的问题:多个服务实例同时访问共享资源,导致数据不一致、业务逻辑混乱?没错,这就是分布式环境下常见的并发问题,而分布式锁就是解决这类问题的...
- 近2万字详解JAVA NIO2文件操作,过瘾
-
原创:小姐姐味道(微信公众号ID:xjjdog),欢迎分享,转载请保留出处。从classpath中读取过文件的人,都知道需要写一些读取流的方法,很是繁琐。最近使用IDEA在打出.这个符号的时候,一行代...
- 学习MVC之租房网站(十二)-缓存和静态页面
-
在上一篇<学习MVC之租房网站(十一)-定时任务和云存储>学习了Quartz的使用、发邮件,并将通过UEditor上传的图片保存到云存储。在项目的最后,再学习优化网站性能的一些技术:缓存和...
- Linux系统下运行c++程序(linux怎么运行c++文件)
-
引言为什么要在Linux下写程序?需要更多关于Linux下c++开发的资料请后台私信【架构】获取分享资料包括:C/C++,Linux,Nginx,ZeroMQ,MySQL,Redis,fastdf...
- 2022正确的java学习顺序(文末送java福利)
-
对于刚学习java的人来说,可能最大的问题是不知道学习方向,每天学了什么第二天就忘了,而课堂的讲解也是很片面的。今天我结合我的学习路线为大家讲解下最基础的学习路线,真心希望能帮到迷茫的小伙伴。(有很多...
- 一个 3 年 Java 程序员 5 家大厂的面试总结(已拿Offer)
-
前言15年毕业到现在也近三年了,最近面试了阿里集团(菜鸟网络,蚂蚁金服),网易,滴滴,点我达,最终收到点我达,网易offer,蚂蚁金服二面挂掉,菜鸟网络一个月了还在流程中...最终有幸去了网易。但是要...
- 多商户商城系统开发全流程解析(多商户商城源码免费下载)
-
在数字化商业浪潮中,多商户商城系统成为众多企业拓展电商业务的关键选择。这类系统允许众多商家在同一平台销售商品,不仅丰富了商品种类,还为消费者带来更多样的购物体验。不过,开发一个多商户商城系统是个复杂的...
你 发表评论:
欢迎- 一周热门
-
-
Redis客户端 Jedis 与 Lettuce
-
高并发架构系列:Redis并发竞争key的解决方案详解
-
redis如何防止并发(redis如何防止高并发)
-
开源推荐:如何实现的一个高性能 Redis 服务器
-
redis安装与调优部署文档(WinServer)
-
Redis 入门 - 安装最全讲解(Windows、Linux、Docker)
-
一文带你了解 Redis 的发布与订阅的底层原理
-
Redis如何应对并发访问(redis控制并发量)
-
oracle数据库查询Sql语句是否使用索引及常见的索引失效的情况
-
Java SE Development Kit 8u441下载地址【windows版本】
-
- 最近发表
- 标签列表
-
- oracle位图索引 (63)
- oracle批量插入数据 (62)
- oracle事务隔离级别 (53)
- oracle 空为0 (50)
- oracle主从同步 (55)
- oracle 乐观锁 (51)
- redis 命令 (78)
- php redis (88)
- redis 存储 (66)
- redis 锁 (69)
- 启动 redis (66)
- redis 时间 (56)
- redis 删除 (67)
- redis内存 (57)
- redis并发 (52)
- redis 主从 (69)
- redis 订阅 (51)
- redis 登录 (54)
- redis 面试 (58)
- 阿里 redis (59)
- redis 搭建 (53)
- redis的缓存 (55)
- lua redis (58)
- redis 连接池 (61)
- redis 限流 (51)