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

资讯详情

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

RabbitMQ实战:SpringAMQP整合与三种交换机详解

RabbitMQ实战:SpringAMQP整合与三种交换机详解 直接上一门课就给你把消息队列讲清楚了说实话SpringCloud系列教程走到第五天往往是最关键也最容易劝退的一步。前四天我们还在用Feign同步调用搞服务间通信今天开始引入RabbitMQ从同步走向异步这是理解微服务架构解耦思想的一个坎。这篇实战笔记主要围绕RabbitMQ基础概念、SpringAMQP框架封装以及fanout、direct、topic三种交换机的使用场景展开内容适合已经跑通Nacos注册中心和OpenFeign调用、想进一步改造为异步通信的同学也适合面试前突击消息中间件核心知识的人。1. 内容整体设计与思路拆解1.1 为什么微服务架构里一定要引入消息队列微服务拆分之后服务间通信方式无非两种同步调用和异步消息。前几天的课程内容基本都是基于Feign的同步调用也就是A服务调用B服务A必须等到B返回结果才算执行完成。这种方式实现简单但一旦B服务响应变慢或者挂掉A服务会被直接拖垮。随着业务链路拉长你还会发现同步调用带来的一系列连锁问题接口性能取决于链条上最慢的那个服务、代码耦合严重、突发流量会导致级联故障。消息队列解决的就是这些问题。它不是在两个服务之间直接建立调用而是通过一个独立的中间件做缓冲。A服务只需要把消息扔到队列里自己就完事了不用管B服务什么时候处理。B服务根据自己的消费能力从队列里取消息处理快慢自己掌控。这种模型天然具备削峰填谷、异步解耦、流量控制等能力。在电商秒杀、订单状态流转、日志收集这些场景里你基本见不到纯同步调用的方案用的全部是消息队列。说实话我之前也带过不少新人发现大多数人第一次接触消息队列时容易陷入一个误区上来就研究怎么保证消息不丢、怎么处理重复消费、怎么保证顺序性这些可靠性问题当然重要但前提是你得先把最基础的消息收发流程跑通理解交换机、队列、路由这几个核心概念之间的关系否则后面聊可靠性全是空中楼阁。所以这篇笔记的顺序是概念模型 → 环境部署 → SpringAMQP基础用法 → 三种交换机实战 → 可靠性机制 → 问题排查。1.2 这个系列为什么选用RabbitMQ和SpringAMQP消息队列的选型市面上常见的有RabbitMQ、Kafka、RocketMQ。这里选择RabbitMQ很大程度上是因为它是传统消息队列里对中小型团队最友好的一个。RabbitMQ基于Erlang语言编写实现了AMQPAdvanced Message Queuing Protocol高级消息队列协议支持多种消息路由模式社区活跃管理界面功能完善学习曲线比Kafka平缓很多。Kafka解决的是海量日志、高吞吐流式数据的场景RocketMQ是阿里开源、在电商领域性能非常强但部署和运维复杂度都高于RabbitMQ。对于SpringCloud微服务这种业务系统间通信的场景RabbitMQ的吞吐量完全够用而且它对消息路由的灵活性远胜于纯Topic模型的Kafka。SpringAMQP则是Spring官方针对AMQP协议做的封装底层基于RabbitMQ的Java客户端实现。它并不是简单的API包装而是把Spring框架的强项——依赖注入、模板模式、注解驱动全部融入消息队列的使用中。你要发送消息直接注入RabbitTemplate掉方法你要消费消息在方法上加RabbitListener注解监听队列即可。相比手写RabbitMQ原生客户端SpringAMQP省掉了大量的样板代码同时基于Spring的自动装配机制配置也极其简单。可以这么理解原生RabbitMQ客户端就像手动挡汽车什么都要自己操作SpringAMQP就是自动挡踩油门就能走把换挡逻辑都封装好了。我们日常开发绝大多数情况下用SpringAMQP就够了。2. RabbitMQ核心模型与部署准备2.1 五大核心概念从一条消息的旅行说起RabbitMQ的消息路由模型核心就五个概念生产者Producer、交换机Exchange、队列Queue、消费者Consumer、绑定Binding。我习惯用快递站来类比。假设你是一家电商公司要给客户发不同类型的通知订单支付成功发短信、发货后发邮件、活动促销发App推送。生产者就是各个业务系统它们产生消息相当于你把包裹交给了快递站交换机就是快递站的分拣中心它决定每个包裹该送到哪个区域队列则是具体的配送站点等着快递员消费者来取货配送绑定关系就是分拣中心墙上贴的配送区域标识。这里有一个容易混淆的点很多人以为生产者直接把消息发到队列实际上在RabbitMQ里消息是不是直接进队列的不消息先到达交换机再由交换机根据路由规则把消息投递到一个或多个队列。交换机本身不存储消息它只做路由判断如果没有任何队列绑定到交换机上或者路由规则匹配不到任何队列消息就会直接丢失。所以使用RabbitMQ时你首先要明确一个问题我该选哪种类型的交换机来满足这条消息的路由需求。除此之外还有几个概念需要理解。**路由键RoutingKey**是生产者发送消息时携带的一个字符串标记交换机根据这个标记和绑定的规则决定将消息路由到哪个队列。**虚拟主机Virtual Host**是RabbitMQ里的逻辑隔离空间默认有一个名为/的vhost不同应用或团队可以使用不同vhost实现数据隔离权限控制的最小粒度也是vhost级别。**通道Channel**是客户端与RabbitMQ建立的逻辑连接它在TCP连接之上复用一个TCP连接可以创建多个Channel避免频繁建立TCP连接带来的性能开销。2.2 环境部署Docker Compose方式最省心RabbitMQ的安装方式有很多种Windows直接下载exe安装包、Linux用yum/apt安装、生产环境用Docker容器化部署。我这里强烈推荐Docker Compose方式不仅安装步骤少还能把管理插件一并带上后续数据卷挂载、容器启停也方便管理。先创建docker-compose.yml文件version: 3.8 services: rabbitmq: image: rabbitmq:3.12-management container_name: rabbitmq hostname: rabbitmq ports: - 5672:5672 # AMQP协议端口Java客户端连接用 - 15672:15672 # Web管理界面端口 environment: RABBITMQ_DEFAULT_USER: admin RABBITMQ_DEFAULT_PASS: admin123 volumes: - ./data:/var/lib/rabbitmq然后在命令行执行docker-compose up -d等容器启动后访问http://localhost:15672输入用户名admin、密码admin123就能进入管理控制台。需要留意的是5672和15672两个端口缺一不可我之前见过有人只映射了15672结果管理界面能打开但项目里怎么都连不上RabbitMQ最后排查半天发现是5672端口没映射出来。如果你的环境不方便用DockerWindows安装其实也不复杂但要注意版本匹配问题。RabbitMQ是基于Erlang编写的它对Erlang版本有严格的兼容要求安装RabbitMQ之前必须先装对应版本的Erlang。比较稳妥的做法是去RabbitMQ官网查看版本兼容对照表然后下载对应的Erlang OTP版本。装好之后进入RabbitMQ安装目录的sbin文件夹执行rabbitmq-plugins enable rabbitmq_management开启Web管理插件再执行rabbitmq-server start启动服务。容器启动成功之后建议先在管理界面手动创建一个队列和交换机试试。我通常在正式写代码前会先在UI上手动发一条消息、消费再删掉确认消息收发能力正常。这个习惯帮我排除过很多次环境层面的问题如果UI都不能正常收发那代码层面的报错大概率是连接配置或者网络问题不是业务代码的问题。3. SpringAMQP基础用法与核心配置3.1 引入依赖与配置连接参数SpringBoot项目引入SpringAMQP非常简便只需要在pom.xml中加一个依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency然后在application.yml中配置连接参数spring: rabbitmq: host: 127.0.0.1 port: 5672 virtual-host: / username: admin password: admin123 publisher-confirm-type: correlated publisher-returns: true template: mandatory: true这里配置的注意点有两个。第一host不要一不小心填成localhost在一些Linux环境下localhost解析的是IPv6地址::1而RabbitMQ默认监听IPv4会导致连接超时建议直接写127.0.0.1。第二最后那三行publisher-confirm-type、publisher-returns、mandatory是为了开启生产者消息到达确认和失败回退功能这部分后面专门讲但建议一开始就写上否则要补充可靠性能力的时候还得改配置重启服务。SpringBoot的自动配置会帮我们创建ConnectionFactory、RabbitTemplate、RabbitAdmin等核心Bean。RabbitAdmin这个类很关键它负责自动声明交换机、队列和绑定关系。默认情况下你在代码里定义好的Bean会被RabbitAdmin自动声明到RabbitMQ服务器上不需要手动到管理界面创建这也是SpringAMQP很爽的一点。3.2 发送消息RabbitTemplate的正确打开方式RabbitTemplate是SpringAMQP提供的消息发送模板类主要方法有convertAndSend和send两个系列。send系列需要传入Message对象用起来比较繁琐convertAndSend系列支持直接传Object对象SpringAMQP内部用消息转换器MessageConverter将Java对象转成消息体开发时绝大多数场景用的都是convertAndSend。先写一个最简单的生产者示例Service public class OrderService { Autowired private RabbitTemplate rabbitTemplate; public void createOrder(Order order) { // 业务逻辑保存订单 orderMapper.insert(order); // 发送消息通知其他服务订单已创建 rabbitTemplate.convertAndSend(order.exchange, order.create, order); } }这里convertAndSend方法三个参数依次是交换机名称、路由键、消息体。消息体会被SimpleMessageConverter序列化为字节数组默认使用的是JDK的Java序列化机制生成的消息是二进制数据在管理界面查看消息内容时是一堆乱码。实际项目中我强烈建议把消息转换器换成Jackson2JsonMessageConverter配置也很简单Configuration public class RabbitMQConfig { Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); } }换了Jackson转换器之后消息体就是标准的JSON字符串前端都能看懂跨语言消费也没问题排查问题时的体验能提升一个档次。需要注意的是生产者和消费者两端的消息转换器要一致否则消费者拿到消息后无法反序列化。3.3 接收消息注解监听器最省事消费端用RabbitListener结合RabbitHandler就能实现消息监听。RabbitListener标注在类或方法上声明当前方法监听哪个队列RabbitHandler在同一个监听类中有多个方法时用于区分不同的消息类型实现同一队列消费不同类型消息的自动分发。最常见的写法是方法直接标注RabbitListenerComponent public class OrderMessageListener { RabbitListener(queues order.queue) public void handleOrderMessage(Order order) { log.info(接收到订单消息{}, order); // 处理业务逻辑 } }SpringAMQP会自动把消息体反序列化为Order对象你在方法里直接处理业务即可。这里有一个Java基础的点要注意如果方法参数用的是自定义DTO对象那么这个DTO必须有无参构造函数否则反序列化会报错。另外方法参数也可以直接声明为String这样拿到的是消息体的原始字符串适合做日志记录或者前置校验。关于监听容器SpringAMQP的SimpleRabbitListenerContainerFactory默认使用SimpleMessageListenerContainer它底层会创建多个消费者线程并发消费。并发线程数由concurrency属性控制默认情况下一个监听器只有一个线程如果消费者处理消息的速度跟不上生产者就会出现消息积压你可以通过RabbitListener(queues order.queue, concurrency 3-5)指定线程数范围来提升消费能力。这个参数需要根据实际业务吞吐量和数据库连接池大小来设定不是越大越好设置过大会把数据库连接数打满。4. 三种交换机实战从广播到通配符路由4.1 fanout交换机广播模式不关心路由键fanout交换机是三种交换机里最直接的一种它把收到的每条消息都复制一份发送到所有绑定到它上面的队列路由键对它来说完全没用。这种模式适合一对多广播的场景典型应用是用户下单后库存服务、积分服务、短信服务、日志服务等多个服务都需要知道这个事件各自己做各的事。代码定义如下Configuration public class FanoutExchangeConfig { // 声明fanout交换机 Bean public FanoutExchange fanoutExchange() { return new FanoutExchange(fanout.order.exchange); } // 声明队列 Bean public Queue fanoutOrderQueueA() { return new Queue(fanout.order.queue.a); } Bean public Queue fanoutOrderQueueB() { return new Queue(fanout.order.queue.b); } // 绑定队列到交换机fanout不需要指定路由键 Bean public Binding bindingA(FanoutExchange fanoutExchange, Queue fanoutOrderQueueA) { return BindingBuilder.bind(fanoutOrderQueueA).to(fanoutExchange); } Bean public Binding bindingB(FanoutExchange fanoutExchange, Queue fanoutOrderQueueB) { return BindingBuilder.bind(fanoutOrderQueueB).to(fanoutExchange); } }发送消息时路由键随便传或者传空字符串都能投递成功rabbitTemplate.convertAndSend(fanout.order.exchange, , message);如果两个消费者分别监听fanout.order.queue.a和fanout.order.queue.b那么每发送一条消息两个消费者都会收到。我在实际项目中常用fanout来发系统通知比如用户注册成功之后邮件服务、短信服务、推荐服务对同一事件各取所需互不影响。这种模式还有一个好处新增一个感兴趣的队列时只需要新写一个绑定关系不需要改动生产者的代码扩展性非常好。4.2 direct交换机点对点精准投递direct交换机是工作中最常用的类型它的路由规则是精确匹配消息的路由键和绑定关系中的路由键完全一致时消息才会被路由到对应队列。这个特别像邮局收信信封上写的地址必须和门牌号完全对上才能投递成功。还是以订单场景为例假设生产端会根据订单状态发不同类型的消息下单消息路由键为order.create支付成功消息路由键为order.pay.success发货消息路由键为order.shipped。消费者A只关心订单创建和支付事件消费者B只关心发货事件配置如下Configuration public class DirectExchangeConfig { Bean public DirectExchange directExchange() { return new DirectExchange(direct.order.exchange); } Bean public Queue directOrderCreateQueue() { return new Queue(direct.order.create.queue); } Bean public Queue directOrderShipQueue() { return new Queue(direct.order.ship.queue); } Bean public Binding bindCreateQueue(DirectExchange directExchange, Queue directOrderCreateQueue) { return BindingBuilder.bind(directOrderCreateQueue).to(directExchange).with(order.create); } Bean public Binding bindCreateQueueToPay(DirectExchange directExchange, Queue directOrderCreateQueue) { return BindingBuilder.bind(directOrderCreateQueue).to(directExchange).with(order.pay.success); } Bean public Binding bindShipQueue(DirectExchange directExchange, Queue directOrderShipQueue) { return BindingBuilder.bind(directOrderShipQueue).to(directExchange).with(order.shipped); } }注意看上面的配置directOrderCreateQueue这个队列和一个交换机绑定的时候可以绑定多个路由键也就是说多个不同路由键的消息都能路由到同一个队列。实际开发中这个用法很普遍比如一个订单服务需要同时处理创建和取消事件就可以把这两个路由键绑定到同一个队列。发送消息的代码前面已经写过只需要保证convertAndSend里的路由键字符串与绑定关系完全一致。这里最容易踩的坑就是拼写不一致比如多一个空格、大小写不同消息就悄悄丢了你如果不做生产者确认根本发现不了。4.3 topic交换机通配符匹配最灵活的规则引擎topic交换机在direct的基础上增加了通配符支持它的路由键需要由点号.分隔成多个单词绑定关系可以包含带通配符的路由键模式。*代表一个单词#代表零个或多个单词。简单理解*.order.*能匹配create.order.success和pay.order.failorder.#能匹配order.create、order.create.success、order本身。这种灵活性让topic成为路由规则的瑞士军刀。还是订单场景假设所有消息的交换机是topic.order.exchange路由键设计为订单号.状态码比如10001.created、10002.paid。现在有消费者A只关心创建类的消息消费者B关心所有订单的消息消费者C只关心具体某个订单号的状态变化Configuration public class TopicExchangeConfig { Bean public TopicExchange topicExchange() { return new TopicExchange(topic.order.exchange); } Bean public Queue orderCreatedQueue() { return new Queue(topic.order.created.queue); } Bean public Queue allOrderQueue() { return new Queue(topic.order.all.queue); } Bean public Queue specialOrderQueue() { return new Queue(topic.order.special.queue); } Bean public Binding bindCreatedQueue(TopicExchange topicExchange, Queue orderCreatedQueue) { return BindingBuilder.bind(orderCreatedQueue).to(topicExchange).with(*.created); } Bean public Binding bindAllOrderQueue(TopicExchange topicExchange, Queue allOrderQueue) { return BindingBuilder.bind(allOrderQueue).to(topicExchange).with(order.#); } Bean public Binding bindSpecialQueue(TopicExchange topicExchange, Queue specialOrderQueue) { return BindingBuilder.bind(specialOrderQueue).to(topicExchange).with(10001.*); } }发送消息时路由键按照设计好的规则拼装即可。topic模式在真实项目中的价值在于它能一套交换机覆盖多种订阅需求不会像fanout那样所有绑定队列都收到消息造成消息冗余也不会像direct那样必须精确指定路由键缺少灵活性。我这里给的例子比较简化实际项目里通常会把路由键设计成有层级结构的业务标识例如trade.order.status.success再用通配符按需匹配不同粒度。4.4 三种交换机的选型对比今天课程里有个环节是面试题常客怎么回答交换机的区别。我整理一个表格方便你理解记忆交换机类型路由规则适用场景路由键角色fanout广播忽略路由键事件通知、系统公告、配置刷新无意义可以不填direct路由键精确匹配点对点定向发送、按消息类型分发需要精确匹配topic路由键通配符匹配复杂的多条件过滤、按主题订阅使用*和#通配符选型的原则其实很简单能精确就用direct需要模糊匹配就用topic想让所有人都收到就用fanout。实际项目中往往三层交换机混合使用比如一个订单事件同时进入fanout交换机做系统通知又通过topic交换机把按订单类型过滤的消息分发给特定服务。这不是什么花哨做法有时候业务需求就是这么多维度想用一个交换机搞定所有订阅关系反而会把路由键设计得很复杂。5. 消息可靠性保障生产端确认与消费端重试5.1 生产端消息确认机制用消息队列最怕什么消息发出去就丢了业务还浑然不知。我见过好几起线上事故排查到最后都是消息丢失引发的数据不一致用户投诉技术背锅。所以这一节讲的是可靠性的第一道防线生产端确认。RabbitMQ的确认机制分为两种一种是publisher-confirm确认消息是否到达交换机另一种是publisher-return消息从交换机路由到队列失败时触发回调。配置项在前面已经写过了现在看代码怎么处理回调。Component Slf4j public class RabbitMQConfirmCallback implements RabbitTemplate.ConfirmCallback, RabbitTemplate.ReturnsCallback { PostConstruct public void init() { rabbitTemplate.setConfirmCallback(this); rabbitTemplate.setReturnsCallback(this); } Autowired private RabbitTemplate rabbitTemplate; Override public void confirm(CorrelationData correlationData, boolean ack, String cause) { if (ack) { log.info(消息已到达交换机消息ID{}, correlationData.getId()); } else { log.error(消息到达交换机失败原因{}, cause); // 这里可以做补偿处理比如重新发送 } } Override public void returnedMessage(ReturnedMessage returned) { log.error(消息路由到队列失败交换机{}路由键{}消息体{}, returned.getExchange(), returned.getRoutingKey(), new String(returned.getMessage().getBody())); // 处理路由失败的消息 } }发消息的时候给每条消息带一个全局唯一的消息ID用CorrelationData包装这样回调处理时就能知道是具体哪条消息出了问题CorrelationData correlationData new CorrelationData(UUID.randomUUID().toString()); rabbitTemplate.convertAndSend(exchange, routingKey, message, correlationData);这套机制的核心价值在于是异步确认不会阻塞消息发送性能损失很小但能极大提升消息投递的可观测性。你千万不要觉得这是多余的生产环境一旦出现消息丢失这就是你排查的第一道线索。5.2 消费端手动确认与重回队列消息从交换机到达队列后下一个可靠性风险点就是消费者。消费者默认采用自动确认模式也就是只要消息被消费方法接收RabbitMQ就认为消息处理成功并删除。如果消费方法内部抛异常呢消息已经删了业务没有执行成功消息就丢了。所以对关键业务我建议改成手动确认模式。配置如下spring: rabbitmq: listener: simple: acknowledge-mode: manual消费者代码同步调整Component Slf4j public class OrderMessageListener { RabbitListener(queues order.queue) public void handleOrderMessage(Order order, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) throws Exception { try { // 业务处理 processOrder(order); // 处理成功手动确认 channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error(订单消息处理失败{}, order, e); // 处理失败拒绝消息并让消息重新入队继续尝试 channel.basicNack(deliveryTag, false, true); } } }这里basicAck表示消息处理成功basicNack表示处理失败第三个参数requeue设为true时消息会重新回到队列头部或尾部。不过这里有一个实际业务里很容易踩的坑如果消费者业务代码一直抛异常、一直requeue消息就会陷入无限循环形成死循环消息。更合理的做法是把requeue设为false让消息进入死信队列DLQ由专门的逻辑去处理异常消息人工介入或者延迟重试。5.3 死信队列和延迟消息的思路说到死信队列它本质上不是一个特殊类型的交换机或队列而是一种对普通队列的附加策略。当你定义一个普通队列时可以指定一个死信交换机x-dead-letter-exchange和死信路由键x-dead-letter-routing-key。当这条队列中的消息满足以下条件之一时就会被自动转发到指定的死信交换机消息被消费者拒绝且requeuefalse消息在队列中存活时间超过设置的TTL队列消息数量达到上限最早的未消费消息被丢弃死信队列的配置用SpringAMQP的QueueBuilder写起来非常简洁Bean public Queue orderQueue() { return QueueBuilder.durable(order.queue) .deadLetterExchange(dlx.order.exchange) .deadLetterRoutingKey(order.dead) .build(); } Bean public Queue deadLetterQueue() { return new Queue(order.dead.queue); } Bean public DirectExchange dlxExchange() { return new DirectExchange(dlx.order.exchange); } Bean public Binding deadLetterBinding(DirectExchange dlxExchange, Queue deadLetterQueue) { return BindingBuilder.bind(deadLetterQueue).to(dlxExchange).with(order.dead); }死信队列再配合TTL还能实现延迟消息的效果。比如给队列设置x-message-ttl10000消息在队列里待够10秒没被消费后自动进入死信队列消费者监听死信队列就相当于收到了延迟10秒的消息。这个方案在支付超时关闭订单、订单自动确认收货这类场景里特别常用不用额外引入延迟消息插件虽然精度不高但足够应付大多数业务场景。RabbitMQ官方也有延迟消息插件rabbitmq_delayed_message_exchange可以精确实现毫秒级延迟但从简单可控的角度用死信队列TTL已经能解决90%的问题。6. 常见问题与踩坑记录新手最容易翻车的地方6.1 消息一直发不出去先查这四件事消息队列的问题排查往往卡在前置环境上。这里整理一份我排查RabbitMQ连接问题的标准动作希望能给你省点时间。现象排查点解决思路连接超时端口是否映射、防火墙是否放行检查5672端口连通性telnet认证失败用户名密码是否正确、vhost是否配置在管理界面创建一个测试用户交换机找不到代码里的命名和实际创建命名不一致看管理界面Exchanges列表用同名重新绑定消息有堆积但消费不了消费者线程被阻塞、死锁dump线程栈检查是否有阻塞的数据库操作消息丢失且无报错没有配置确认回调、消息被错误路由开启confirm/return先确认交换机绑定关系连接层的问题最笨也最有效的方法是先看管理界面的Connections和Channels页面。应用启动后如果界面上出现了连接记录说明网络通、认证通过问题大概率在代码逻辑如果界面上完全看不到连接那就是网络或配置的问题。6.2 消费者收不到消息别再怀疑代码了消费者收不到消息是一类高频问题原因多种多样而且大多数时候不是代码逻辑的问题。我自己遇到过一个典型的场景消息生产者投递成功了管理界面能看到队列里有消息但消费者就是没反应。排查过程先看日志发现消费者一直没启动连监听容器都没创建成功检查依赖才发现项目里的spring-boot-starter-amqp版本和SpringBoot父版本不兼容导致自动配置没有生效。这种问题在SpringBoot版本升级后尤其常见解决方案是把父POM版本统一对齐。另一个常见原因是消息被上一轮的消费者截胡了造成队列消息被其他消费者抢走但你不知道。如果你本地开着多个实例比如一个测试环境、一个本地环境连的都是同一个RabbitMQ服务器监听同一个队列消息就会被负载均衡地分发到多个消费者你本地调试时感觉像“消息丢了”实际是被别的实例消费了。排查方式是看管理界面的Queues页面观察消息的消费情况或者临时把其他实例停掉。6.3 消息重复消费分布式系统的通病最后说一个所有消息队列都绕不开的问题重复消费。无论你怎么配置消费者意外重启、网络抖动、消费超时都会导致消息被投递两次以上。保证不重复消费的核心手段只有一个消费幂等。代码里用业务唯一键在数据库层面做去重比如消费订单消息时先查一下订单表里是否已经存在这条订单编号存在就直接跳过。我个人的习惯是建一张message_consume_log表记录每次消费的消息ID消费前先查这个表没有记录才执行业务逻辑并插入日志用数据库唯一约束来兜底。这套方案虽然简单但是通用性极强不管消息队列怎么重试都不会造成重复业务操作。7. 总结我的体会这条异步改造之路值得走今天我们通过RabbitMQ和SpringAMQP把服务间的同步调用改造为异步消息理解了交换机在消息传递中的核心地位。这一课的内容量确实不少概念多、代码也多但异步通信本来就是微服务架构里的一座高山翻过这座山你才能感受到微服务架构真正的魅力——各服务独立演进、互不拖累。我个人经验是消息中间件的学习一定要把三种交换机的路由规则反复敲几遍因为面试的时候高频考、项目里高频用理解和熟练程度直接决定你的技术底气。Day5的内容先到这接下来可以继续往消息可靠性、分布式事务方向深入那些才是真正考验架构功底的地方。
返回列表