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));@Componentpublic 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@Slf4jpublic class DemoRabbitMQListener implements ChannelAwareMessageListener {
@Override public void onMessage(Message message, Channel channel) throws Exception { log.info("message:{}", message.getBody()); //todo: 接下来就是各自的业务逻辑,就是消费环节 }}相关推荐
- CentOS7安装Jenkins:从零开始搭建CI/CD环境 - CI/CD流水线中集成消息队列
- SSE服务器推送:实时网络传输协议详解与实战 - 实时通信协议对比
- Metasploit渗透测试实战:白帽子从入门到精通 - 安全测试中的消息队列漏洞利用
文章分享
如果这篇文章对你有帮助,欢迎分享给更多人!