RabbitMQ从入门到实战:消息队列核心用法详解

1181 字
6 分钟
RabbitMQ从入门到实战:消息队列核心用法详解

本文基于作者实际开发经验,已在项目中使用 RabbitMQ 实现异步消息处理,涵盖交换机配置、队列监听及消费者实践。

最近处理访客记录所以,来学习下RabbitMQ。之前同事已经写好了,这里只需要进行消费,后续会逐渐完善。

0.介绍#

0.1 交换机(Exchanges)#

RabbitMQ中生产者发送的消息都是发送到交换机,再由交换机推入队列。所以生产者不知道队列去了哪里,就靠Exchange来控制,交换机总共有以下几种类型。

0.1.1 广播模式(fanout)#

扇出所有消息进入队列,类似广播。

0.1.2 直接交换(direct)#

绑定相关的routerKey分发到不同的队列,简单说就是direct交换机接收了消息后,根据关键词分发队列。

0.1.3 主题模式(topic)#

direct路由比较单一,所以提升了routerKey的能力,在关键词标记下加上了通配符。

  • *(星号)可以代替一个单词
  • #(井号)可以替代零个或多个单词

1.公共配置类#

spring:
rabbitmq:
host: 127.0.0.1
port: 5672
username: guest
password: guest
/**
* RabbitMQ公共配置类
*/
public class RabbitMQConfig {
/** RabbitMQ的队列主题名称 */
public static final String RABBITMQ_TOPIC = "rabbitmqTopic";
/** RabbitMQ的DIRECT交换机名称 */
public static final String RABBITMQ_DIRECT_EXCHANGE = "rabbitmqDirectExchange";
/** RabbitMQ的Direct交换机和队列绑定的匹配键 */
public static final String RABBITMQ_DIRECT_ROUTING = "rabbitmqDirectRouting";
}

2.消费消息的两种方式#

把记录塞进队列里的时候,只是完成了第一步,那你肯定要对他进行消费。分为两种推模式和拉模式:推模式就是生产者发布消息时,主动推送给消费者;拉模式则是消费者发送请求后才会发送。

3.监听队列的两种方式#

一种是@RabbitListener注解的方式,一种是实现SpringBoot的ChannelAwareMessageListener接口的方式。

3.1 @RabbitListener#

如果demoData想不转换成String直接推,得在这个数据流实现序列化。

innerRabbitTemplate.convertAndSend(InnerMQConfig.TOPIC_EXCHANGE, msgKey, JSONObject.toJSONString(demoData));
@Component
public class DemoRabbitMQListener {
@RabbitListener(queues = "demo_queue")
@RabbitHandler
public void demoQueue(Message message){
System.out.println("message:"+message.getBody());
}
}

3.2 实现ChannelAwareMessageListener接口#

听前辈说直接实现这个接口,就不用管底层是谁的消息队列了,因为是基于SpringBoot。这个实现起来有点麻烦,我总结了以下顺序:

3.2.1 创建连接工厂(ConnectionFactory)#

publisherConfirms:消息发送到exchange,返回成功或者失败。 publishReturns:消息从exchange到queue,发送成功或者失败。

@Bean(name = "DemoConnectionFactory")
@Primary
public ConnectionFactory connectionFactory() {
CachingConnectionFactory connectionFactory = new CachingConnectionFactory();
connectionFactory.setHost(host);
connectionFactory.setPort(port);
connectionFactory.setUsername(username);
connectionFactory.setPassword(password);
connectionFactory.setVirtualHost(virtualHost);
connectionFactory.setPublisherConfirms(true);
connectionFactory.setPublisherReturns(true);
return connectionFactory;
}

3.2.2 初始化组件(RabbitAdmin)#

@Bean(name = "DemoRabbitAdmin")
@Primary
public RabbitAdmin rabbitAdmin(@Qualifier("DemoConnectionFactory") ConnectionFactory connectionFactory) {
RabbitAdmin rabbitAdmin = new RabbitAdmin(connectionFactory);
rabbitAdmin.setAutoStartup(true);
return rabbitAdmin;
}

3.2.3 创建交换器(Exchange)#

durable:是否持久化,RabbitMQ关闭后,没有持久化的Exchange将被清除 autoDelete:是否自动删除,如果没有与之绑定的Queue,直接删除 internal:是否内置的,如果为true,只能通过Exchange到Exchange

@Bean(DEMO_EXCHANGE)
public TopicExchange exchange() {
return new TopicExchange(DEMO_EXCHANGE, true, false);
}

3.2.4 创建队列(Queue)#

创建队列主要掌握这几个参数:

  • name:队列名称
  • durable:队列是否持久化
  • exclusive:是否排他的队列
  • autoDelete:是否自动删除
  • arguments:队列中的消息什么时候会自动被删除
@Bean(QUEUE_NAME)
public Queue QUEUE_DEMO() {
return new Queue(QUEUE_NAME, true, false, false);
}

3.2.5 绑定队列到交换机(Binding)#

@Bean
public Binding BINGING_EXCHANGE_QUEUE(@Qualifier(QUEUE_NAME) Queue queue,
@Qualifier(DEMO_EXCHANGE) Exchange exchange) {
return BindingBuilder.bind(queue).to(exchange).with(ROUTING_KEY).noargs();
}

3.2.6 创建监听容器(SimpleMessageListenerContainer)#

@Bean
public SimpleMessageListenerContainer simpleMessageListenerContainer(
@Qualifier("DemoConnectionFactory") ConnectionFactory connectionFactory,
DemoRabbitMQListener demoRabbitMQListener,
@Qualifier(QUEUE_NAME) Queue queue
) throws AmqpException {
SimpleMessageListenerContainer listenerContainer = new SimpleMessageListenerContainer(connectionFactory);
listenerContainer.setConcurrentConsumers(listenerSize);
listenerContainer.setQueues(queue);
listenerContainer.setExposeListenerChannel(true);
listenerContainer.setAcknowledgeMode(AcknowledgeMode.AUTO);
listenerContainer.setMessageListener(demoRabbitMQListener);
return listenerContainer;
}

3.2.7 创建操作类(RabbitTemplate)#

setConfirmCallback的消息回调是在生产者端要把参数丢进去的。

@Bean(name = "DemoRabbitTemplate")
@Primary
public RabbitTemplate rabbitTemplate(@Qualifier("DemoConnectionFactory") ConnectionFactory connectionFactory) {
RabbitTemplate rabbitTemplate = new RabbitTemplate(connectionFactory);
rabbitTemplate.setMandatory(true);
rabbitTemplate.setConnectionFactory(connectionFactory);
rabbitTemplate.setConfirmCallback(new RabbitTemplate.ConfirmCallback() {
@Override
public void confirm(CorrelationData correlationData, boolean ack, String cause) {
if (ack) {
log.info("消息发送到交换机成功, 消息 = {}", correlationData.getId());
} else {
log.error("消息发送到交换机失败! 消息: {}; 错误原因: cause: {}", correlationData.getId(), cause);
}
}
});
rabbitTemplate.setReturnCallback(new RabbitTemplate.ReturnCallback() {
@Override
public void returnedMessage(Message message, int replyCode, String replyText, String exchange, String routingKey) {
String messageId = message.getMessageProperties().getMessageId();
String result = null;
try {
result = new String(message.getBody(), "UTF-8");
} catch (Exception e) {
log.error("消息发送失败", e);
}
log.error("消息发送失败, 消息ID = {}; 消息内容 = {}", messageId, result);
}
});
return rabbitTemplate;
}

3.2.8 监听消费(RabbitMQListener)#

这个类要注意用@Service或者@Component注解让他交给IOC

@Service
@Slf4j
public class DemoRabbitMQListener implements ChannelAwareMessageListener {
@Override
public void onMessage(Message message, Channel channel) throws Exception {
log.info("message:{}", message.getBody());
//todo: 接下来就是各自的业务逻辑,就是消费环节
}
}

相关推荐#

文章分享

如果这篇文章对你有帮助,欢迎分享给更多人!

RabbitMQ从入门到实战:消息队列核心用法详解
https://blog.wmovie.site/posts/rabbitmq-guide/
作者
朵朵
发布于
2023-08-18
许可协议
CC BY-NC-SA 4.0
Profile Image of the Author
朵朵
全栈开发者,5 年 Java/SpringBoot + Vue/React 开发经验,曾维护 Delphi 遗留系统。业余深耕 NAS 与软路由,专注 Homelab 家庭网络搭建及 Docker 自托管应用实战。坐标福州,正迈向独立开发之路。
公告
欢迎来到我的博客!有什么想要的邮箱告诉我,最近沉迷nas和软路由。
分类
标签
站点统计
文章
25
分类
4
标签
18
总字数
32,908
运行时长
0
最后活动
0 天前

文章目录