文章目录
- 前言
- 一、WorkQueue
- 二、交换机
-
- 1.FanOut交换机
- 2.Direct交换机
- 3.Topic交换机
- 三、消息转换机
- 总结
前言
简单使用下RabbitMQ中特殊队列和交换机
一、WorkQueue
WorkQueue,任务模型,即多个消费者绑定在一个队列上,共同消费该队列的消息。核心价值是解决“耗时任务的分发与负载均衡”问题 —— 把需要耗时处理的任务封装成消息发送到队列,再由多个工作者(消费者)并行获取消息处理,避免单个进程积压任务导致阻塞,同时提高任务处理的吞吐量。
首先在consumer类中定义两个消费者
@RabbitListener(queues = "work.queue")
public void listenWorkQueueMessage1(String msg) throws InterruptedException {
System.out.println("spring 消费者1接收到消息:【" + msg + "】");
}
@RabbitListener(queues = "work.queue")
public void listenWorkQueueMessage2(String msg) throws InterruptedException {
System.err.println("spring 消费者2……..接收到消息:【" + msg + "】");
}
在publisher类中发送20条消息,
public void testSendWorkMessage() throws Exception {
String queue = "work.queue";
for (int i = 0; i < 20; i++) {
String message = "hello, spring amqp" + i;
rabbitTemplate.convertAndSend(queue, message);
}
}
运行结果如下:
不难发现,队列中的一条消息一次被一个消费者消费;多个消费者依次消费队列中的消息。
考虑到是否和两个线程的性能有关,让消费者2休眠20ms。
//消费者2
Thread.sleep(20);
运行结果如下:
虽然两个消费者处理的消息编号不规律,但是总数仍然为10条。且消费者1处理完后,并没有继续处理,而是交给消费者2处理。即正常情况下不会因为消费者处理快慢而改变消息分配。
要想实现能者多劳,就得给consumer类配置:
spring:
rabbitmq:
listener:
simple:
prefetch: 1 # 每次只能获取一条消息,处理完成才能获取下一个消息
运行结果如下:
成功实现能者多劳,完成了WorkQueue的搭建。
二、交换机
第一个介绍的是简单的队列模型,这一节介绍RabbitMQ中的三种特殊交换机,并看看各自的基础用法。
Java中有两种声明交换机的方式:一个是基于Bean;另一个是基于@RabbitMQListener注解。
基于Bean就是定义一个配置类,将交换机,队列和两者绑定关系注册到Bean中。
基于@RabbitMQListener注解,后续演示再详细介绍各个参数。
1.FanOut交换机
简单理解就是广播,将发送者的消息传递给所有绑定的队列。这里采用基于Bean的方式构建含Fanout交换机的异步结构。 启动配置类,即可在控制台自动创建队列,交换机,并且两者自动绑定。
@Configuration
public class FanoutConfig {
@Bean
public FanoutExchange fanoutExchange() {
return new FanoutExchange("fanout_exchange");
}
@Bean
public Queue fanoutQueue1() {
return new Queue("fanout_queue1");
}
@Bean
public Binding bindFanoutQueue1( Queue fanoutQueue1, FanoutExchange fanoutExchange) {
return BindingBuilder.bind(fanoutQueue1).to(fanoutExchange);
}
}
收发消息逻辑根据上述代码改一下变量参数,发送到指定队列即可,这里就不再粘贴代码。
2.Direct交换机
Direct交换机会将接收到的消息根据规则路由到指定的Queue,因此称为定向路由。每一个Queue都与交换机设置一个BindingKey,发布者发送消息时,指定消息的RoutingKey,交换机将消息路由到BindingKey与RoutingKey一致的队列。
这里就基于@RabbitMQListener注解创建异步结构。解释下对应参数:
- bingdings——队列和交换机的绑定关系,需要加上对应注解@QueueBinding,括号内为所要绑定的队列和交换机。
- 1
- value——绑定的队列,加上对应注解@Queue,括号内指定名字,还可以指定消息是否持久化,默认开启。
- 1
- exchange——绑定的交换机,加上对应注解@Exchange,括号内指定名字,类型和对应的BindingKey。
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "direct.queue1"),
exchange = @Exchange(name = "hmall.direct", type = ExchangeTypes.DIRECT),
key = {"red", "blue"}
))
public void listenDirectQueue1(String msg){
System.out.println("消费者1接收到direct.queue1的消息:【" + msg + "】");
}
@RabbitListener(bindings = @QueueBinding(
value = @Queue(name = "direct.queue2"),
exchange = @Exchange(name = "hmall.direct", type = ExchangeTypes.DIRECT),
key = {"red", "yellow"}
))
public void listenDirectQueue2(String msg){
System.out.println("消费者2接收到direct.queue2的消息:【" + msg + "】");
}
这里我们指定给队列1指定了red和blue;队列2指定了red和tyellow。接下来我们来收发信息:
//key为red
public void testSendDirectMessagered() throws Exception {
String exchangename= "test.direct";
String message = "hello, springred amqp";
rabbitTemplate.convertAndSend(exchangename, "red", message);
}
//key为blue
public void testSendDirectMessageblue() throws Exception {
String exchangename= "test.direct";
String message = "hello, springblue amqp";
rabbitTemplate.convertAndSend(exchangename, "blue", message);
}
//key为yellow
public void testSendDirectMessageyellow() throws Exception {
String exchangename= "test.direct";
String message = "hello, springyellow amqp";
rabbitTemplate.convertAndSend(exchangename, "yellow", message);
}
运行结果如下:
只有key相同的消费者接收到了信息。
3.Topic交换机
TopicExchange也是基于RoutingKey做消息路由,但是routingKey通常是多个单词的组合,并且以.分割。Queue与Exchange指定BindingKey时可以使用通配符:#:代指0个或多个单词;*:代指一个单词。相当于增强版Direct。其功能就不演示了,key可以设置为red.——以red开头的key的消息都可以接收。
三、消息转换机
这个和前面redis存储数据时类似,即保存的数据无法正常显示,所以需要进行序列化。 第一步,引入依赖
<dependency>
<groupId>com.fasterxml.jackson.core</groupId>
<artifactId>jackson–databind</artifactId>
</dependency>
第二步,配置MessageConvertor类
@Bean
public MessageConverter messageConverter(){
// 1.定义消息转换器
Jackson2JsonMessageConverter jackson2JsonMessageConverter = new Jackson2JsonMessageConverter();
// 2.配置自动创建消息id,用于识别不同消息,也可以在业务中基于ID判断是否是重复消息
jackson2JsonMessageConverter.setCreateMessageIds(true);
return jackson2JsonMessageConverter;
}
这里配置后,如果要能够在不同模块的不同包名下调用,需要依赖自动装配原理。
总结
今天就把整个RabbitMQ的基础理论实践了一番,操作简单,但如何运用到业务中还需深入学习。如果有一些概念不清楚的,可以回看昨天的日记后端学习日记2.2。





