尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

第45篇:Java中间件:RabbitMQ入门,实现消息队列通信

第45篇:Java中间件:RabbitMQ入门,实现消息队列通信 简介消息队列是Java后端开发中核心的中间件技术广泛应用于高并发、分布式系统架构中。本文将从零入门RabbitMQ完整讲解消息队列核心作用、RabbitMQ核心原理、环境搭建、Spring Boot集成、消息收发、消息可靠性确认机制以及死信队列实战全程附带可直接运行的代码案例适合零基础开发者快速上手。适用人群Java后端初学者、分布式架构学习者、需要掌握MQ实战的开发人员技术栈Spring Boot 2.x/3.x RabbitMQ一、消息队列核心作用为什么要用MQ在传统单体架构中业务代码同步执行接口链路长、响应慢、容错性差。引入消息队列MQ后可以异步处理业务彻底优化系统架构核心价值集中在三点解耦、削峰、异步通信。1.1 业务解耦传统业务流程高度耦合比如用户下单后需要同步执行库存扣减、短信通知、物流生成、积分发放等一系列操作任意一个环节出错都会导致下单失败。引入MQ后主流程只完成创建订单后续所有附属业务通过消息队列异步消费业务之间完全隔离互不影响大幅降低系统耦合度提升代码可维护性。1.2 流量削峰秒杀、限时活动等场景会出现瞬时海量请求直接冲击数据库和业务接口极易导致系统雪崩。MQ可以作为流量缓冲区瞬时请求全部存入队列消费者按照系统最大处理能力匀速消费避免瞬时高并发压垮后端服务实现流量削峰填谷。1.3 异步通信同步调用需要等待所有业务执行完毕才能返回结果接口响应耗时极长。MQ支持异步通信主线程发送消息后直接返回无需等待后续业务执行极大提升接口响应速度和系统吞吐量。二、RabbitMQ核心概念底层原理必懂RabbitMQ是一款基于AMQP协议的开源消息中间件可靠性高、稳定性强、社区活跃是企业主流MQ选型之一。其核心架构由生产者、交换机、队列、绑定、消费者五部分组成核心三要素交换机、队列、绑定。2.1 核心角色介绍生产者Producer消息的发送方负责创建消息并发送到RabbitMQ交换机消费者Consumer消息的接收方持续监听队列获取并处理消息队列Queue消息的存储载体消息最终落地在队列中等待消费者消费持久化存储不丢失消息交换机Exchange消息路由中转站接收生产者消息根据路由规则分发到对应队列绑定Binding建立交换机和队列之间的关联关系是消息路由的桥梁2.2 交换机四大类型重点交换机没有存储消息的能力只负责路由核心四种类型Direct直连交换机精准匹配根据路由键完全匹配分发消息一对一通信适用于单消息单消费场景Topic主题交换机模糊匹配支持通配符*和#多对多通信适用于复杂业务订阅场景Fanout扇形交换机广播模式无视路由键绑定该交换机的所有队列都会接收消息适用于群发通知场景Headers头交换机根据消息头属性匹配极少使用2.3 绑定Binding绑定是交换机与队列的映射关系只有完成绑定交换机才能将消息路由到指定队列每个绑定会关联对应的路由键RoutingKey作为消息分发的匹配规则。三、RabbitMQ安装与配置Windows/Linux通用RabbitMQ基于Erlang语言开发安装前需提前安装Erlang环境推荐Docker快速安装简单高效、无需配置环境变量。3.1 Docker一键安装推荐# 1. 拉取RabbitMQ镜像带管理控制台 docker pull rabbitmq:3-management # 2. 启动容器 docker run -d \ --name rabbitmq \ -p 5672:5672 \ -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASS123456 \ rabbitmq:3-management3.2 端口说明5672MQ服务通信端口程序连接使用15672Web管理控制台端口浏览器访问使用3.3 访问控制台浏览器访问http://localhost:15672账号admin密码123456登录后可查看交换机、队列、消息状态、连接信息等方便调试。四、Spring Boot集成RabbitMQ基础环境搭建4.1 引入Maven依赖Spring Boot整合RabbitMQ核心依赖spring-boot-starter-amqp自动封装连接、消息收发、确认机制等核心功能。dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency !-- lombok 简化代码 -- dependency groupIdorg.projectlombok/groupId artifactIdlombok/artifactId optionaltrue/optional /dependency4.2 配置文件application.yml配置MQ连接信息、消息确认模式、持久化等核心参数开启生产者确认、消费者手动ACK保障消息可靠性。spring: rabbitmq: # 服务连接配置 host: localhost port: 5672 username: admin password: 123456 virtual-host: / # 开启生产者确认机制 publisher-confirm-type: correlated # 开启消息投递失败返回 publisher-returns: true listener: simple: # 消费者手动确认消息 acknowledge-mode: manual # 开启重试机制 retry: enabled: true max-attempts: 34.3 RabbitMQ核心配置类配置交换机、普通业务队列、绑定关系同时注入消息转换器支持JSON消息传输。import org.springframework.amqp.core.*; import org.springframework.amqp.rabbit.connection.ConnectionFactory; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.amqp.support.converter.Jackson2JsonMessageConverter; import org.springframework.amqp.support.converter.MessageConverter; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; import java.util.HashMap; import java.util.Map; Configuration public class RabbitMQConfig { // 普通交换机、队列、路由键定义 public static final String NORMAL_EXCHANGE normal.exchange; public static final String NORMAL_QUEUE normal.queue; public static final String NORMAL_ROUTING_KEY normal.key; // 死信相关定义 public static final String DLX_EXCHANGE dlx.exchange; public static final String DLX_QUEUE dlx.queue; public static final String DLX_ROUTING_KEY dlx.key; /** * JSON消息转换器 */ Bean public MessageConverter jsonMessageConverter() { return new Jackson2JsonMessageConverter(); } /** * 自定义RabbitTemplate开启确认回调 */ Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate new RabbitTemplate(connectionFactory); rabbitTemplate.setMessageConverter(jsonMessageConverter()); // 开启强制消息返回 rabbitTemplate.setMandatory(true); return rabbitTemplate; } /** * 普通直连交换机 */ Bean public DirectExchange normalExchange() { return ExchangeBuilder.directExchange(NORMAL_EXCHANGE).durable(true).build(); } /** * 普通队列绑定死信交换机消息超时/异常则进入死信队列 */ Bean public Queue normalQueue() { MapString, Object args new HashMap(); // 绑定死信交换机 args.put(x-dead-letter-exchange, DLX_EXCHANGE); // 死信路由键 args.put(x-dead-letter-routing-key, DLX_ROUTING_KEY); // 消息TTL10秒超时未消费则进入死信队列 args.put(x-message-ttl, 10000); return QueueBuilder.durable(NORMAL_QUEUE).withArguments(args).build(); } /** * 普通队列与交换机绑定 */ Bean public Binding normalBinding() { return BindingBuilder.bind(normalQueue()).to(normalExchange()).with(NORMAL_ROUTING_KEY); } // 死信队列、交换机配置 Bean public DirectExchange dlxExchange() { return ExchangeBuilder.directExchange(DLX_EXCHANGE).durable(true).build(); } Bean public Queue dlxQueue() { return QueueBuilder.durable(DLX_QUEUE).build(); } Bean public Binding dlxBinding() { return BindingBuilder.bind(dlxQueue()).to(dlxExchange()).with(DLX_ROUTING_KEY); } }五、消息发送与接收实战5.1 生产者发送消息通过RabbitTemplate实现消息发送支持普通文本消息、JSON对象消息。import lombok.RequiredArgsConstructor; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.web.bind.annotation.GetMapping; import org.springframework.web.bind.annotation.RestController; RestController RequiredArgsConstructor public class MQProducerController { private final RabbitTemplate rabbitTemplate; GetMapping(/send/msg) public String sendMsg() { String msg Hello RabbitMQ 入门实战消息; // 发送消息交换机、路由键、消息内容 rabbitTemplate.convertAndSend(RabbitMQConfig.NORMAL_EXCHANGE, RabbitMQConfig.NORMAL_ROUTING_KEY, msg); return 消息发送成功; } }5.2 消费者监听接收消息使用RabbitListener注解监听指定队列实现消息消费。import com.rabbitmq.client.Channel; import org.springframework.amqp.rabbit.annotation.RabbitListener; import org.springframework.amqp.support.AmqpHeaders; import org.springframework.messaging.handler.annotation.Header; import org.springframework.stereotype.Component; Component public class MQConsumer { /** * 监听普通业务队列 */ RabbitListener(queues RabbitMQConfig.NORMAL_QUEUE) public void consumeMsg(String msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception { System.out.println(消费者接收消息 msg); // 后续手动ACK确认此处先注释 // channel.basicAck(tag, false); } /** * 监听死信队列 */ RabbitListener(queues RabbitMQConfig.DLX_QUEUE) public void consumeDlxMsg(String msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception { System.out.println(死信队列接收异常消息 msg); channel.basicAck(tag, false); } }六、消息确认机制保障消息不丢失MQ消息丢失是生产环境常见问题RabbitMQ通过生产者确认、消费者确认双重机制保障消息可靠性。6.1 生产者确认机制Confirm Return生产者确认分为两种场景消息成功投递到交换机、消息未成功路由到队列通过回调函数感知投递结果。import org.springframework.amqp.rabbit.connection.CorrelationData; import org.springframework.amqp.rabbit.core.RabbitTemplate; import org.springframework.stereotype.Component; import javax.annotation.PostConstruct; Component public class RabbitConfirmCallback implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnCallback { private final RabbitTemplate rabbitTemplate; public RabbitConfirmCallback(RabbitTemplate rabbitTemplate) { this.rabbitTemplate rabbitTemplate; } PostConstruct public void init() { // 注入确认回调 rabbitTemplate.setConfirmCallback(this); rabbitTemplate.setReturnCallback(this); } /** * 交换机投递确认 */ Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { if (ack) { System.out.println(消息成功投递到交换机); } else { System.err.println(消息投递交换机失败原因 cause); // 可自定义重试、日志记录、告警逻辑 } } /** * 队列路由失败回调 */ Override public void returnedMessage(org.springframework.amqp.core.Message message, int replyCode, String replyText, String exchange, String routingKey) { System.err.println(消息路由队列失败交换机 exchange 路由键 routingKey); } }6.2 消费者确认机制手动ACK默认自动ACK会导致消息被删除若消费者业务异常消息会丢失。手动ACK可以保证业务执行成功后再确认消息异常时拒绝消息。basicAck成功消费确认消息队列删除消息basicNack消费失败拒绝消息可选择重回队列或丢弃优化后的消费者代码RabbitListener(queues RabbitMQConfig.NORMAL_QUEUE) public void consumeMsg(String msg, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long tag) throws Exception { try { System.out.println(消费者处理消息 msg); // 模拟业务逻辑 // int i 1/0; // 手动确认消息消费成功 channel.basicAck(tag, false); } catch (Exception e) { System.err.println(消息消费失败进入重试逻辑); // 消费失败消息重回队列false不批量拒绝true重回队列 channel.basicNack(tag, false, true); } }七、死信队列DLX实战详解7.1 死信队列核心作用当消息出现以下三种情况时会变为死信消息自动路由到绑定的死信队列消息超时未被消费配置TTL过期时间消费者手动拒绝消息且不重回队列队列消息数量达到最大限制死信队列主要用于处理异常消息、实现延迟任务、消息兜底重试、故障排查。7.2 实战测试流程启动项目访问接口/send/msg发送消息注释消费者的basicAck确认代码让消息无法被正常消费等待10秒TTL超时消息自动转为死信死信交换机将消息路由到死信队列死信消费者监听并处理异常消息7.3 业务场景落地实际开发中可利用死信队列实现订单超时取消、支付超时回滚、异常消息兜底处理等经典场景是企业级RabbitMQ开发的必备方案。八、完整项目总结本文从零完成RabbitMQ全流程实战核心知识点回顾消息队列三大核心价值解耦、削峰、异步解决传统同步业务的性能与耦合问题RabbitMQ核心架构交换机、队列、绑定四大交换机适配不同业务场景Docker快速搭建RabbitMQ环境开箱即用无需复杂配置Spring Boot快速集成实现基础消息收发双重消息确认机制生产者ConfirmReturn、消费者手动ACK彻底解决消息丢失问题死信队列实现消息兜底处理适配超时、异常消费场景该项目代码可直接用于学习、二次开发适配中小型项目的消息队列基础架构。
返回列表