赞
踩
RabbitMQ是一个在AMQP基础上完成的,可复用的企业消息系统。他遵循Mozilla Public License开源协议。开发语言为Erlang。
linux系统中安装RabbitMQ比较繁琐,这里使用的是Docker安装。
docker search rabbitmq:management
docker pull macintoshplus/rabbitmq-management
docker images
docker run -d --hostname mzw-rabbitmq --name rabbitmq -e RABBITMQ_DEFAULT_USER=admin -e RABBITMQ_DEFAULT_PASS=admin -p 15672:15672 -p 5672:5672 c20
命令解读:
docker ps
访问地址:http://192.168.2.xx:15672/
输入启动容器时设置的用户密码登录
这就表示RabbitMQ安装成功了
创建SpringBoot项目并引入相关依赖
# RabbitMQ 配置
spring.rabbitmq.name=rabbitmq-demo01
spring.rabbitmq.host=192.168.2.22
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=admin
# 自定义一个属性,设置队列的名称
mq.queue.name=hello-queue
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.CustomExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
// 添加@Configuration 注解,表示一个注解类
@Configuration
public class QueueConfig {
@Value("${mq.queue.name}")
private String queueName;
/**
* 初始化短信队列
* @return
*/
@Bean
public Queue delayedSmsQueueInit() {
return new Queue(queueName);
}
}
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
* 创建一个rabbitmq消费者
*/
@Component
public class Receiver {
// 接受MQ消息 并 处理消息
@RabbitListener(queues = {"${mq.queue.name}"})
public void process(String msg){
// 处理消息
System.out.println("我是MQ消费者,我接收到的消息是:" + msg );
}
}
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.stereotype.Component;
/**
* 消息提供者
*/
@Component
public class Sender {
@Autowired
private AmqpTemplate template;
@Value("${mq.queue.name}")
private String queueName;
// 发送消息
public void send(String msg){
// 队列名,消息内容
template.convertAndSend(queueName,msg);
}
}
@Autowired
private Sender sender;
@Test
void contextLoads() {
sender.send("你好啊......");
}
RabbitMQ中有五种主要的交互器分别如下
交换器 | 说明 |
---|---|
direct | 发布与订阅 完全匹配 |
fanout | 广播 |
topic | 主体,规则匹配 |
fanout | 转发 |
custom | 自定义 |
上边已经演示,这里不重复演示。
死信队列就是在某种情况下,导致消息无法被正常消费(异常,过期,队列已满等),存放这些未被消费的消息的队列即为死信队列。
# RabbitMQ 配置
spring.rabbitmq.name=rabbitmq-demo01
spring.rabbitmq.host=192.168.2.22
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=admin
###死信队列
mq.dlx.exchange=mq_dlx_exchange
mq.dlx.queue=mq_dlx_queue
mq.dlx.routingKey=mq_dlx_key
###备胎交换机
mq.exchange=mq_exchange
mq.queue=mq_queue
mq.routingKey=routing_key
import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
@Configuration
public class MQConfig {
/**
* 普通交换机
*/
@Value("${mq.exchange}")
private String mqExchange;
/**
* 普通队列
*/
@Value("${mq.queue}")
private String mqQueue;
/**
* 普通路由key
*/
@Value("${mq.routingKey}")
private String mqRoutingKey;
/**
* 死信交换机
*/
@Value("${mq.dlx.exchange}")
private String dlxExchange;
/**
* 死信队列
*/
@Value("${mq.dlx.queue}")
private String dlxQueue;
/**
* 死信路由
*/
@Value("${mq.dlx.routingKey}")
private String dlxRoutingKey;
/**
* 声明死信交换机
* @return DirectExchange
*/
@Bean
public DirectExchange dlxExchange() {
return new DirectExchange(dlxExchange);
}
/**
* 声明死信队列
* @return Queue
*/
@Bean
public Queue dlxQueue() {
return new Queue(dlxQueue);
}
/**
* 声明普通业务交换机
* @return DirectExchange
*/
@Bean
public DirectExchange mqExchange() {
return new DirectExchange(mqExchange);
}
/**
* 声明普通队列
* @return Queue
*/
@Bean
public Queue mqQueue() {
// 普通队列绑定我们的死信交换机
Map<String, Object> arguments = new HashMap<>(2);
//死信交换机
arguments.put("x-dead-letter-exchange", dlxExchange);
//死信队列
arguments.put("x-dead-letter-routing-key", dlxRoutingKey);
return new Queue(mqQueue, true, false, false, arguments);
}
/**
* 绑定死信队列到死信交换机
* @return Binding
*/
@Bean
public Binding binding(Queue dlxQueue,DirectExchange dlxExchange) {
return BindingBuilder.bind(dlxQueue).to(dlxExchange).with(dlxRoutingKey);
}
/**
* 绑定普通队列到普通交换机
* @return Binding
*/
@Bean
public Binding mqBinding(Queue mqQueue,DirectExchange mqExchange) {
return BindingBuilder.bind(mqQueue).to(mqExchange).with(mqRoutingKey);
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.web.bind.annotation.RestController;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.web.bind.annotation.RequestMapping;
/**
* 生产者
*/
@RestController
@Slf4j
public class MQProducer {
@Autowired
private RabbitTemplate rabbitTemplate;
/**
* 普通交换机
*/
@Value("${mq.exchange}")
private String mqExchange;
/**
* 普通路由key
*/
@Value("${mq.routingKey}")
private String mqRoutingKey;
@RequestMapping("/sendMsg")
public String sendMsg() {
String msg = "Hello RabbitMQ ......";
//发送消息 参数一:交换机 参数二:路由键(用来指定发送到哪个队列)
rabbitTemplate.convertAndSend(mqExchange, mqRoutingKey, msg, message -> {
// 设置消息过期时间 10秒过期 如果过期时间内还没有被消费 就会发送给死信队列
message.getMessageProperties().setExpiration("10000");
return message;
});
log.info("生产者发送消息:{}", msg);
return "success";
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
* 消费者
*/
@Component
@Slf4j
public class MQConsumer {
/**
* 监听队列回调的方法
*
* @param msg
*/
@RabbitListener(queues = {"${mq.queue}"})
public void mqConsumer(String msg) {
log.info("正常普通消费者消息MSG:{}", msg);
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
* 死信消费者
*/
@Component
@Slf4j
public class MQDlxConsumer {
/**
* 死信队列监听队列回调的方法
*
* @param msg
*/
@RabbitListener(queues = {"${mq.dlx.queue}"})
public void mqConsumer(String msg) {
log.info("死信队列消费普通消息:msg{}", msg);
}
}
访问:http://127.0.0.1:9023/sendMsg 会被 消费者 消费掉
将 消费者 代码注释掉,在访问http://127.0.0.1:9023/sendMsg,等待10秒钟后会被死信队列接收到。
官网下载:https://www.rabbitmq.com/community-plugins.html
我的RabbitMQ是3.12 b版本的,下载此插件
wget https://github.com/rabbitmq/rabbitmq-delayed-message-exchange/releases/download/v3.12.0/rabbitmq_delayed_message_exchange-3.12.0.ez
docker cp ./rabbitmq_delayed_message_exchange-3.12.0.ez de24369edeb4:/plugins
docker exec -it de24369edeb4 /bin/bash
rabbitmq-plugins enable rabbitmq_delayed_message_exchange
rabbitmq-plugins list
E* 或 e* 代表 插件已启用
在RabbitMQ控制台可以看到
# RabbitMQ 配置
spring.rabbitmq.name=rabbitmq-demo01
spring.rabbitmq.host=192.168.2.22
spring.rabbitmq.port=5672
spring.rabbitmq.username=admin
spring.rabbitmq.password=admin
# 自定义一个属性,设置队列的名称
mq.queue.name=hello-queue
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.CustomExchange;
import org.springframework.amqp.core.Queue;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import java.util.HashMap;
import java.util.Map;
/**
* 使用x-delayed-message 延时队列插件
*/
@Configuration
public class QueueConfig {
@Value("${mq.queue.name}")
private String queueName;
/**
* 初始化短信队列
* @return
*/
@Bean
public Queue delayedSmsQueueInit() {
return new Queue(queueName);
}
/**
* 初始化延迟交换机
* @return
*/
@Bean
public CustomExchange delayedExchangeInit() {
Map<String, Object> args = new HashMap<>();
// 设置类型,可以为fanout、direct、topic
args.put("x-delayed-type", "direct");
// 第一个参数是延迟交换机名字,第二个是交换机类型,第三个设置持久化,第四个设置自动删除,第五个放参数
return new CustomExchange("delayed_exchange","x-delayed-message", true,false,args);
}
/**
* 短信队列绑定到交换机
* @param delayedSmsQueueInit
* @param customExchange
* @return
*/
@Bean
public Binding delayedBindingSmsQueue(Queue delayedSmsQueueInit, CustomExchange customExchange) {
// 延迟队列绑定延迟交换机并设置RoutingKey为sms
return BindingBuilder.bind(delayedSmsQueueInit).to(customExchange).with("sms").noargs();
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.core.AmqpTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.beans.factory.annotation.Value;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;
/**
* 生产者
*/
@RestController
@Slf4j
public class Sender {
@Autowired
private AmqpTemplate template;
@Value("${mq.queue.name}")
private String queueName;
// 发送消息
@RequestMapping("/sendMsg")
public void send(){
String msg = "Hello RabbitMQ ......";
// 队列名,消息内容
template.convertAndSend(queueName,msg);
log.info("生产者发送消息:{}", msg);
}
@RequestMapping("/sendDelayedMsg")
public void sendDelayedMsg(){
String msg = "Hello RabbitMQ Delayed ......";
// 第一个参数是延迟交换机名称,第二个是Routingkey,第三个是消息主题,第四个是X,并设置延迟时间,单位 是毫秒
template.convertAndSend("delayed_exchange","sms",msg,a -> {
a.getMessageProperties().setDelay(2000);
return a;
});
log.info("生产者发送延时消息:{}", msg);
}
}
import lombok.extern.slf4j.Slf4j;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;
/**
* 消费者
*/
@Component
@Slf4j
public class Receiver {
// 接受MQ消息 并 处理消息
@RabbitListener(queues = {"${mq.queue.name}"})
public void process(String msg){
// 处理消息
log.info("我是MQ消费者,我接收到的消息是:{}", msg);
}
}
访问:http://127.0.0.1:9022/sendMsg
访问:http://127.0.0.1:9022/sendDelayedMsg
Copyright © 2003-2013 www.wpsshop.cn 版权所有,并保留所有权利。