欢迎光临
我们一直在努力

常用:SpringBoot项目引入RabbitMQ步骤

1.引入依赖:

pom.xml添加依赖

<!–核心starter,内置amqp‑client–>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

2.配置文件application.yml

spring:
rabbitmq:
addresses: amqp://登录RabbitMQ的用户名:登录的密码@部署RabbitMQ服务的公网IP地址:端口/虚拟机名
listener:
simple:
acknowledge-mode: auto #消息接收确认(自动),默认是auto,这里是为了直观观察
#以上为必写,以下为可选
# #acknowledge-mode: manual #消息接收确认(手动)
# #prefetch: 1 #消费者持有最大待消费消息数
# #retry:
# #enabled: true # 开启消费者失败重试
# #initial-interval: 5000ms # 初始失败等待时长为5秒
# #max-attempts: 5 # 最大重试次数
# #publisher-confirm-type: correlated #消息发送确认
# #publisher-returns: true #设置回退

3.创建一个常量类(存相关名称)和一个配置类(交换机,队列,绑定关系都写里面)

Constants类:

package com.txm.rabbitmq.constant;

public class Constants {
//工作模式
public static final String WORK_QUEUE="work.queue";

//发布订阅模式fanout
public static final String FANOUT_QUEUE1="fanout.queue1";
public static final String FANOUT_QUEUE2="fanout.queue2";
public static final String FANOUT_EXCHANGE="fanout.exchange";

//路由模式
public static final String DIRECT_QUEUE1="direct.queue1";
public static final String DIRECT_QUEUE2="direct.queue2";
public static final String DIRECT_EXCHANGE="direct.exchange";

//通配符模式
public static final String TOPIC_QUEUE1="topic.queue1";
public static final String TOPIC_QUEUE2="topic.queue2";
public static final String TOPIC_EXCHANGE="topic.exchange";

}

Config类:

package com.txm.rabbitmq.config;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.core.*;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMQConfig {
//工作模式
//声明队列
@Bean("workQueue")
public Queue workQueue(){
return QueueBuilder.durable(Constants.WORK_QUEUE).build();
}

//发布订阅模式
//1.声明队列
@Bean("fanoutQueue1")
public Queue fanoutQueue1(){
return QueueBuilder.durable(Constants.FANOUT_QUEUE1).build();
}
@Bean("fanoutQueue2")
public Queue fanoutQueue2(){
return QueueBuilder.durable(Constants.FANOUT_QUEUE2).build();
}
//2.声明交换机
@Bean("fanoutExchange")
public FanoutExchange fanoutExchange(){
return ExchangeBuilder.fanoutExchange(Constants.FANOUT_EXCHANGE).durable(true).build();
}
//3.声明交换机和队列的绑定关系
@Bean("fanoutQueueBinding1")
public Binding fanoutQueueBinding1(@Qualifier("fanoutExchange") FanoutExchange fanoutExchange,@Qualifier("fanoutQueue1") Queue queue){
return BindingBuilder.bind(queue).to(fanoutExchange);
}
@Bean("fanoutQueueBinding2")
public Binding fanoutQueueBinding2(@Qualifier("fanoutExchange") FanoutExchange fanoutExchange,@Qualifier("fanoutQueue2") Queue queue){
return BindingBuilder.bind(queue).to(fanoutExchange);
}

//路由模式
//1.声明队列
@Bean("directQueue1")
public Queue DIRECT_QUEUE1(){
return QueueBuilder.durable(Constants.DIRECT_QUEUE1).build();
}
@Bean("directQueue2")
public Queue DIRECT_QUEUE2(){
return QueueBuilder.durable(Constants.DIRECT_QUEUE2).build();
}
//2.声明交换机
@Bean("directExchange")
public DirectExchange DIRECT_EXCHANGE(){
return ExchangeBuilder.directExchange(Constants.DIRECT_EXCHANGE).durable(true).build();
}
//3.声明交换机和队列的绑定关系
@Bean("directQueueBinding1")
public Binding directQueueBinding1(@Qualifier("directExchange") DirectExchange directExchange,@Qualifier("directQueue1") Queue queue){
return BindingBuilder.bind(queue).to(directExchange).with("orange");
}
@Bean("directQueueBinding2")
public Binding directQueueBinding2(@Qualifier("directExchange") DirectExchange directExchange,@Qualifier("directQueue2") Queue queue){
return BindingBuilder.bind(queue).to(directExchange).with("black");
}
@Bean("directQueueBinding3")
public Binding directQueueBinding3(@Qualifier("directExchange") DirectExchange directExchange,@Qualifier("directQueue2") Queue queue){
return BindingBuilder.bind(queue).to(directExchange).with("orange");
}

//通配符模式
//1.声明队列
@Bean("topicQueue1")
public Queue topicQueue1(){
return QueueBuilder.durable(Constants.TOPIC_QUEUE1).build();
}
@Bean("topicQueue2")
public Queue topicQueue2(){
return QueueBuilder.durable(Constants.TOPIC_QUEUE2).build();
}
//2.声明交换机
@Bean("topicExchange")
public TopicExchange TOPIC_EXCHANGE(){
return ExchangeBuilder.topicExchange(Constants.TOPIC_EXCHANGE).durable(true).build();
}
//3.声明交换机和队列的绑定关系
@Bean("topicQueueBinding1")
public Binding topicQueueBinding1(@Qualifier("topicExchange") TopicExchange topicExchange,@Qualifier("topicQueue1") Queue queue){
return BindingBuilder.bind(queue).to(topicExchange).with("*.orange.*");
}
@Bean("topicQueueBinding2")
public Binding topicQueueBinding2(@Qualifier("topicExchange") TopicExchange topicExchange,@Qualifier("topicQueue2") Queue queue){
return BindingBuilder.bind(queue).to(topicExchange).with("*.*.rabbit");
}
@Bean("topicQueueBinding3")
public Binding topicQueueBinding3(@Qualifier("topicExchange") TopicExchange topicExchange,@Qualifier("topicQueue2") Queue queue){
return BindingBuilder.bind(queue).to(topicExchange).with("lazy.#");
}

}

4.写生产者(Controller类,发送消息):

直接注入RabbitTemplate发送消息:

package com.txm.rabbitmq.controller;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.beans.factory.annotation.Autowired;
import org.springframework.web.bind.annotation.PathVariable;
import org.springframework.web.bind.annotation.RequestMapping;
import org.springframework.web.bind.annotation.RestController;

@RequestMapping("/producer")
@RestController
public class ProducerController {

@Autowired
private RabbitTemplate rabbitTemplate;

//工作队列模式
@RequestMapping("/work")
public String work(){
for (int i = 0; i < 10; i++) {
//使用内置交换机,RoutingKey和队列名称一致
rabbitTemplate.convertAndSend("", Constants.WORK_QUEUE,"hello spring amqp:Lwork …"+i);

}
return "发送成功";
}

//发布订阅模式(fanout交换机)
@RequestMapping("/fanout")
public String fanout(){
rabbitTemplate.convertAndSend(Constants.FANOUT_EXCHANGE,"","hello spring amqp:fanout…");
return "发送成功";
}

//路由模式
@RequestMapping("/direct/{routingKey}")
public String direct(@PathVariable("routingKey") String routingKey){
rabbitTemplate.convertAndSend(Constants.DIRECT_EXCHANGE,routingKey,"hello spring amqp:direct,my routingKey is "+routingKey);
return "发送成功";
}

//通配符模式
@RequestMapping("/topic/{routingKey}")
public String topic(@PathVariable("routingKey") String routingKey){
rabbitTemplate.convertAndSend(Constants.TOPIC_EXCHANGE,routingKey,"hello spring amqp:topic,my routingKey is "+routingKey);
return "发送成功";
}
}

5.写消费者(listener类,监听消息)

<1>工作队列模式-消费者:

package com.txm.rabbitmq.listener;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

/**
* 工作队列模式-消费者
*/
@Component
public class WorkListener {

@RabbitListener(queues=Constants.WORK_QUEUE)
public void queueListener1(Message message){
System.out.println("listener 1 ["+ Constants.WORK_QUEUE+"] 接收到消息:"+message);
}

@RabbitListener(queues=Constants.WORK_QUEUE)
public void queueListener2(Message message){
System.out.println("listener 2 ["+ Constants.WORK_QUEUE+"] 接收到消息:"+message);
}
}

<2>发布订阅模式(广播模式)-消费者:

package com.txm.rabbitmq.listener;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

/**
* 发布订阅模式-消费者
*/
@Component
public class FanoutListener {

@RabbitListener(queues=Constants.FANOUT_QUEUE1)
public void queueListener1(String message){
System.out.println("队列["+ Constants.FANOUT_QUEUE1+"] 接收到消息:"+message);
}

@RabbitListener(queues=Constants.FANOUT_QUEUE2)
public void queueListener2(Message message){
System.out.println(" 队列["+ Constants.FANOUT_QUEUE2+"] 接收到消息:"+message);
}
}

<3>路由模式-消费者:

package com.txm.rabbitmq.listener;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.core.Message;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

/**
* 路由模式-消费者
*/
@Component
public class DirectListener {

@RabbitListener(queues=Constants.DIRECT_QUEUE1)
public void queueListener1(String message){
System.out.println("队列["+ Constants.DIRECT_QUEUE1+"] 接收到消息:"+message);
}

@RabbitListener(queues=Constants.DIRECT_QUEUE2)
public void queueListener2(String message){
System.out.println(" 队列["+ Constants.DIRECT_QUEUE2+"] 接收到消息:"+message);
}
}

<4>通配符模式-消费者:

package com.txm.rabbitmq.listener;

import com.txm.rabbitmq.constant.Constants;
import org.springframework.amqp.rabbit.annotation.RabbitListener;
import org.springframework.stereotype.Component;

/**
* 通配符模式-消费者
*/
@Component
public class TopicListener {

@RabbitListener(queues=Constants.TOPIC_QUEUE1)
public void queueListener1(String message){
System.out.println("队列["+ Constants.TOPIC_QUEUE1+"] 接收到消息:"+message);
}

@RabbitListener(queues=Constants.TOPIC_QUEUE2)
public void queueListener2(String message){
System.out.println(" 队列["+ Constants.TOPIC_QUEUE2+"] 接收到消息:"+message);
}
}

赞(0)
未经允许不得转载:171主机测评 » 常用:SpringBoot项目引入RabbitMQ步骤
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址