0


Spring Boot整合RabbitMQ详细教程

org.springframework.boot

spring-boot-starter-amqp

1.5.2.RELEASE

第二步:在application.properties文件当中引入RabbitMQ基本的配置信息

#对于rabbitMQ的支持

spring.rabbitmq.host=127.0.0.1

spring.rabbitmq.port=5672

spring.rabbitmq.username=guest

spring.rabbitmq.password=guest

第三步:编写RabbitConfig类,类里面设置很多个EXCHANGE,QUEUE,ROUTINGKEY,是为了接下来的不同使用场景。

/**

Broker:它提供一种传输服务,它的角色就是维护一条从生产者到消费者的路线,保证数据能按照指定的方式进行传输,

Exchange:消息交换机,它指定消息按什么规则,路由到哪个队列。

Queue:消息的载体,每个消息都会被投到一个或多个队列。

Binding:绑定,它的作用就是把exchange和queue按照路由规则绑定起来.

Routing Key:路由关键字,exchange根据这个关键字进行消息投递。

vhost:虚拟主机,一个broker里可以有多个vhost,用作不同用户的权限分离。

Producer:消息生产者,就是投递消息的程序.

Consumer:消息消费者,就是接受消息的程序.

Channel:消息通道,在客户端的每个连接里,可建立多个channel.

*/

@Configuration

public class RabbitConfig {

private final Logger logger = LoggerFactory.getLogger(this.getClass());

@Value(“${spring.rabbitmq.host}”)

private String host;

@Value(“${spring.rabbitmq.port}”)

private int port;

@Value(“${spring.rabbitmq.username}”)

private String username;

@Value(“${spring.rabbitmq.password}”)

private String password;

public static final String EXCHANGE_A = “my-mq-exchange_A”;

public static final String EXCHANGE_B = “my-mq-exchange_B”;

public static final String EXCHANGE_C = “my-mq-exchange_C”;

public static final String QUEUE_A = “QUEUE_A”;

public static final String QUEUE_B = “QUEUE_B”;

public static final String QUEUE_C = “QUEUE_C”;

public static final String ROUTINGKEY_A = “spring-boot-routingKey_A”;

public static final String ROUTINGKEY_B = “spring-boot-routingKey_B”;

public static final String ROUTINGKEY_C = “spring-boot-routingKey_C”;

@Bean

public ConnectionFactory connectionFactory() {

CachingConnectionFactory connectionFactory = new CachingConnectionFactory(host,port);

connectionFactory.setUsername(username);

connectionFactory.setPassword(password);

connectionFactory.setVirtualHost(“/”);

connectionFactory.setPublisherConfirms(true);

return connectionFactory;

}

@Bean

@Scope(ConfigurableBeanFactory.SCOPE_PROTOTYPE)

//必须是prototype类型

public RabbitTemplate rabbitTemplate() {

RabbitTemplate template = new RabbitTemplate(connectionFactory());

return template;

}

}

第四步:编写消息的生产者

@Component

public class MsgProducer implements RabbitTemplate.ConfirmCallback {

private final Logger logger = LoggerFactory.getLogger(this.getClass());

//由于rabbitTemplate的scope属性设置为ConfigurableBeanFactory.SCOPE_PROTOTYPE,所以不能自动注入

private RabbitTemplate rabbitTemplate;

/**

  • 构造方法注入rabbitTemplate

*/

@Autowired

public MsgProducer(RabbitTemplate rabbitTemplate) {

this.rabbitTemplate = rabbitTemplate;

rabbitTemplate.setConfirmCallback(this); //rabbitTemplate如果为单例的话,那回调就是最后设置的内容

}

public void sendMsg(String content) {

CorrelationData correlationId = new CorrelationData(UUID.randomUUID().toString());

//把消息放入ROUTINGKEY_A对应的队列当中去,对应的是队列A

rabbitTemplate.convertAndSend(RabbitConfig.EXCHANGE_A, RabbitConfig.ROUTINGKEY_A, content, correlationId);

}

/**

  • 回调

*/

@Override

public void confirm(CorrelationData correlationData, boolean ack, String cause) {

logger.info(" 回调id:" + correlationData);

if (ack) {

logger.info(“消息成功消费”);

} else {

logger.info(“消息消费失败:” + cause);

}

}

}

第五步:把交换机,队列,通过路由关键字进行绑定,写在RabbitConfig类当中

/**

  • 针对消费者配置

    1. 设置交换机类型
    1. 将队列绑定到交换机

FanoutExchange: 将消息分发到所有的绑定队列,无routingkey的概念

HeadersExchange :通过添加属性key-value匹配

DirectExchange:按照routingkey分发到指定队列

TopicExchange:多关键字匹配

*/

@Bean

public DirectExchange defaultExchange() {

return new DirectExchange(EXCHANGE_A);

}

/**

  • 获取队列A

  • @return

*/

@Bean

public Queue queueA() {

return new Queue(QUEUE_A, true); //队列持久

}

@Bean

public Binding binding() {

return BindingBuilder.bind(queueA()).to(defaultExchange()).with(RabbitConfig.ROUTINGKEY_A);

}

一个交换机可以绑定多个消息队列,也就是消息通过一个交换机,可以分发到不同的队列当中去。

@Bean

public Binding binding() {

return BindingBuilder.bind(queueA()).to(defaultExchange()).with(RabbitConfig.ROUTINGKEY_A);

}

@Bean

public Binding bindingB(){

return BindingBuilder.bind(queueB()).to(defaultExchange()).with(RabbitConfig.ROUTINGKEY_B);

}

第六步:编写消息的消费者,这一步也是最复杂的,因为可以编写出很多不同的需求出来,写法也有很多的不同。

比如一个生产者,一个消费者

@Component

@RabbitListener(queues = RabbitConfig.QUEUE_A)

public class MsgReceiver {

private final Logger logger = LoggerFactory.getLogger(this.getClass());

@RabbitHandler

public void process(String content) {

logger.info("接收处理队列A当中的消息: " + content);

}

}

比如一个生产者,多个消费者,可以写多个消费者,并且他们的分发是负载均衡的。

@Component

@RabbitListener(queues = RabbitConfig.QUEUE_A)

public class MsgReceiverC_one {

private final Logger logger = LoggerFactory.getLogger(this.getClass());

@RabbitHandler

public void process(String content) {

logger.info("处理器one接收处理队列A当中的消息: " + content);

}

}

@Component

@RabbitListener(queues = RabbitConfig.QUEUE_A)

public class MsgReceiverC_two {

private final Logger logger = LoggerFactory.getLogger(this.getClass());

@RabbitHandler

public void process(String content) {

logger.info("处理器two接收处理队列A当中的消息: " + content);

}

}

另外一种消息处理机制的写法如下,在RabbitMQConfig类里面增加bean:

@Bean

public SimpleMessageListenerContainer messageContainer() {

//加载处理消息A的队列

SimpleMessageListenerContainer container = new SimpleMessageListenerContainer(connectionFactory());

//设置接收多个队列里面的消息,这里设置接收队列A

//假如想一个消费者处理多个队列里面的信息可以如下设置:

//container.setQueues(queueA(),queueB(),queueC());

container.setQueues(queueA());

container.setExposeListenerChannel(true);

//设置最大的并发的消费者数量

最后

笔者已经把面试题和答案整理成了面试专题文档

image

image

image

image

image

image

//假如想一个消费者处理多个队列里面的信息可以如下设置:

//container.setQueues(queueA(),queueB(),queueC());

container.setQueues(queueA());

container.setExposeListenerChannel(true);

//设置最大的并发的消费者数量

最后

笔者已经把面试题和答案整理成了面试专题文档

[外链图片转存中…(img-YrVM3PMb-1714452690743)]

[外链图片转存中…(img-HzjoMwqG-1714452690743)]

[外链图片转存中…(img-NwQsbI6W-1714452690744)]

[外链图片转存中…(img-1vfHO6gi-1714452690744)]

[外链图片转存中…(img-xQ87MftD-1714452690744)]

[外链图片转存中…(img-wS6rBuSa-1714452690744)]

本文已被CODING开源项目:【一线大厂Java面试题解析+核心总结学习笔记+最新讲解视频+实战项目源码】收录


本文转载自: https://blog.csdn.net/2401_84103936/article/details/138342008
版权归原作者 2401_84103936 所有, 如有侵权,请联系我们删除。

“Spring Boot整合RabbitMQ详细教程”的评论:

还没有评论