SpringBoot+RabbitMQ实现注册异步通知链路:幂等、重试与死信实践

发布时间:2026/9/11 20:38:34
SpringBoot+RabbitMQ实现注册异步通知链路:幂等、重试与死信实践 用户注册这个功能几乎所有后端开发都写过。但绝大多数人的版本只是“把数据塞进数据库、返回成功”就完事了后续的欢迎邮件、新手积分、行为埋点全部同步写在注册接口里接口越写越慢业务越堆越乱。这篇文章我想用一套最小但完整的工程化方案讲清楚在 SpringBoot 里怎么用消息队列把“注册主流程”和“注册后续动作”解开做成一条异步通知链路。整套代码可以直接跑通重点覆盖消息投递、消费幂等、手动 ACK、失败重试和死信兜底这几个工程上绕不开的点。适合已经会写 SpringBoot CRUD、但没正经在项目里用过消息队列的同学。1. 为什么注册链路要异步化先看懂同步调用的三个坑很多同学第一次接触异步思路时会觉得不就是开个线程池把邮件发送丢后台去嘛为什么非要引入一套消息队列这个问题如果不先想清楚后面写代码就只是在“照猫画虎”。所以我先拆一下同步链路到底哪里疼。1.1 同步调用的痛点响应时间、错误放大与峰值卡顿假设注册接口里要做四件事写入用户表、发送欢迎邮件、发放新手优惠券、同步用户信息到数据分析系统。如果是同步串行调用接口耗时就是四次操作时间的总和。举例来说写库正常 20ms邮件服务偶尔要 1 秒营销系统高峰期要 3 秒那你的注册接口平时不显山露水一旦某个下游服务抖动用户侧感受到的就是“点了注册按钮半天没反应”。更麻烦的是错误放大。邮件服务宕机时同步调用如果没做异常捕获整个注册事务都可能回滚用户填了半天的表单最后看到的是“系统错误”。你当然可以在代码里 catch 掉异常但这样写路由是湾的——你注册模块的代码里全是别人家服务的容错逻辑改一次别人接口的报错方式你就得跟着改一次。峰值卡顿就更直接了。逢年过节整点活动注册请求瞬时冲上来一万个下游每个服务都扛不住。同步链路下你的服务线程全被耗在下游等待上线程池一满新的注册请求直接排队或者被拒绝整个应用性能雪崩。1.2 消息队列的三大作用映射到注册场景消息队列在架构上常被概括为三个作用异步、解耦、削峰。落到注册这个场景异步注册接口只做写库和发消息两个必要动作邮件、积分、埋点这些业务自成消费者收到消息后再慢慢处理。主接口响应时间从“所有业务总耗时”降为“必要业务耗时”。解耦注册模块不需要知道邮件系统、营销系统、数据分析系统怎么实现。只要约定好消息体协议下游服务自己订阅队列即可。新增一个消费者不影响注册模块代码。削峰瞬时注册量再大消息队列作为缓冲区先把消息攒下来消费者按自己的处理能力匀速消费避免后端服务被瞬时流量打垮。1.3 为什么选 RabbitMQ对比线程池和 Kafka很多人会问异步我直接用线程池不也能实现是的单机场景下线程池确实能解决响应慢的问题但它有三个致命短板第一线程池的内存队列服务重启就丢消息没有持久化能力第二没有 ACK 机制任务执行失败不会自动重试你只能在 run() 里 try-catch 然后自己写重试逻辑第三线程池无法跨服务共享下游邮件服务如果是一个独立部署的服务线程池根本碰不到它。而 Kafka 又是另一个极端它是为海量日志、超高吞吐设计的功能上和 RabbitMQ 比缺少开箱即用的死信队列、延迟队列这些工程利器再加上运维和客户端相对重注册通知这种业务事件场景用 RabbitMQ 更顺手。小团队、单应用、微服务初阶RabbitMQ 是性价比极高的选择。2. 工程级最小闭环环境与骨架要先搭对写业务逻辑之前环境这块我建议一次性把版本踩实不然光版本兼容问题就能耗掉你半天。2.1 版本选型SpringBoot 2.7.x 是当前最稳的分水岭很多博客直接给 SpringBoot 3.x 的配置但对新手不太友好JDK 要求 17 起步很多老依赖和自定义配置类都要跟着改。个人建议工程实战直接用 SpringBoot 2.7.18配 JDK 8 或者 11这套组合在绝大多数公司里依然是绝对主力网上排坑经验也最全。RabbitMQ 我用的 3.12 版本直接通过 Docker 拉起不需要在本机装 Erlang——如果选择 Windows 原生安装最容易踩的坑就是 Erlang 版本和 RabbitMQ 版本不匹配服务起不来排查还特别费劲。2.2 Docker 一键拉起依赖服务真实开发环境里MySQL、Redis、RabbitMQ 全部用 Docker 管理启动快、清理方便。下面是本工程需要的最小依赖docker run -d --name rabbitmq \ -p 5672:5672 -p 15672:15672 \ -e RABBITMQ_DEFAULT_USERadmin \ -e RABBITMQ_DEFAULT_PASSadmin123 \ rabbitmq:3.12-management docker run -d --name redis \ -p 6379:6379 \ redis:7-alpine启动后访问http://localhost:15672用 admin/admin123 登录能看到管理控制台就算装好了。控制台非常有用查队列积压、看消息内容、人工重新投递消息都在这里操作。2.3 工程骨架只引入必要的依赖用 Spring Initializr 创建工程pom.xml 里我最终保留这五个依赖dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-web/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-amqp/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-data-redis/artifactId /dependency dependency groupIdorg.springframework.boot/groupId artifactIdspring-boot-starter-jdbc/artifactId /dependency dependency groupIdmysql/groupId artifactIdmysql-connector-j/artifactId scoperuntime/scope /dependency注意 starter 是spring-boot-starter-amqp不是spring-boot-starter-rabbitmq。我见过好几次有人写错依赖名然后项目启动直接报 ClassNotFound。Redis 在这里不是必须但做消息幂等去重时它非常好用所以一并加上。3. 核心链路实现注册、发消息、消费通知全流程下面进入正题。最小闭环的价值在于“完整跑通一条链路”所以不会堆太多业务逻辑但每个环节的工程处理都会给足。3.1 注册接口与用户服务写库成功后再发消息先定义一个最普通的用户注册接口省略掉复杂的参数校验重点是看消息发送的触发时机。RestController RequestMapping(/api/user) public class UserController { Resource private UserService userService; PostMapping(/register) public ResultLong register(RequestBody RegisterRequest request) { Long userId userService.register(request.getUsername(), request.getPassword(), request.getEmail()); return Result.success(userId); } }Service 层的核心逻辑是校验用户名是否重复、加密密码、写用户表、发送 MQ 消息。消息发送绝不能放在写库之前否则库没写成消息先发出去了消费者拿着不存在的数据做了一堆事情。Service Slf4j public class UserServiceImpl implements UserService { Resource private UserMapper userMapper; Resource private RegisterProducer registerProducer; Override Transactional(rollbackFor Exception.class) public Long register(String username, String password, String email) { // 1. 校验用户是否存在 User user userMapper.selectByUsername(username); if (user ! null) { throw new BizException(用户名已存在); } // 2. 密码加盐加密 String salt RandomUtil.randomString(8); String encryptedPwd DigestUtils.md5DigestAsHex((salt password).getBytes()); // 3. 写用户表 User newUser new User(); newUser.setUsername(username); newUser.setPassword(encryptedPwd); newUser.setSalt(salt); newUser.setEmail(email); userMapper.insert(newUser); // 4. 写库成功之后发送消息 registerProducer.sendRegisterMessage(newUser.getId(), username, email); return newUser.getId(); } }这里有个很关键的工程细节数据库事务和 MQ 消息发送存在“分布式事务”问题上面代码其实有极小的概率出现“数据库回滚了但消息已经发出去”。最小闭环里我用的是“先提交事务再发消息”的常规做法也就是让事务提交成功后再进入消息发送逻辑。如果你要更严格的可靠性需要引入本地消息表或者事务消息这个我在第 4 节展开说。3.2 RabbitMQ 配置交换机、队列、绑定一次说清先理清这三个概念。生产者不直接把消息扔给队列而是扔给交换机交换机再根据路由键把消息投递到绑定的队列。消费者只从队列取消息。这样设计的好处是生产者和队列完全解耦——生产者只管发队列怎么路由、消费者怎么处理它一概不关心。注册通知这个场景用 Direct 交换机即可路由键完全匹配队列绑定键。Fanout 是广播给所有绑定队列Topic 是通配符匹配当前场景用不上。Configuration public class RabbitConfig { public static final String REGISTER_EXCHANGE register.exchange; public static final String REGISTER_QUEUE register.queue; public static final String REGISTER_ROUTING_KEY register.notify; public static final String REGISTER_DLX_EXCHANGE register.dlx.exchange; public static final String REGISTER_DLX_QUEUE register.dlx.queue; public static final String REGISTER_DLX_ROUTING_KEY register.dlx.routing; // 业务队列绑定死信交换机 Bean public Queue registerQueue() { return QueueBuilder.durable(REGISTER_QUEUE) .deadLetterExchange(REGISTER_DLX_EXCHANGE) .deadLetterRoutingKey(REGISTER_DLX_ROUTING_KEY) .build(); } Bean public DirectExchange registerExchange() { return new DirectExchange(REGISTER_EXCHANGE, true, false); } Bean public Binding registerBinding() { return BindingBuilder.bind(registerQueue()).to(registerExchange()).with(REGISTER_ROUTING_KEY); } // 死信交换机与队列处理多次重投仍然失败的消息 Bean public Queue registerDlxQueue() { return QueueBuilder.durable(REGISTER_DLX_QUEUE).build(); } Bean public DirectExchange registerDlxExchange() { return new DirectExchange(REGISTER_DLX_EXCHANGE, true, false); } Bean public Binding registerDlxBinding() { return BindingBuilder.bind(registerDlxQueue()).to(registerDlxExchange()).with(REGISTER_DLX_ROUTING_KEY); } }所有队列和交换机我都设置了 durable持久化到磁盘RabbitMQ 重启后元数据不丢。消息本身能不能扛住重启取决于发送时是否带持久化属性SpringBoot 默认帮我们设置的MessageDeliveryMode.PERSISTENT这一点不用担心。3.3 生产者封装消息体要精简不要传整个对象生产者设计上注意两点消息体只用基础字段组装不要直接把数据库的 User 实体传过去发送时路由键对应配置类里的常量不要魔法字符串。Data Builder public class RegisterMessage { private Long userId; private String username; private String email; private Long timestamp; }Component Slf4j public class RegisterProducer { Resource private RabbitTemplate rabbitTemplate; public void sendRegisterMessage(Long userId, String username, String email) { RegisterMessage message RegisterMessage.builder() .userId(userId) .username(username) .email(email) .timestamp(System.currentTimeMillis()) .build(); rabbitTemplate.convertAndSend( RabbitConfig.REGISTER_EXCHANGE, RabbitConfig.REGISTER_ROUTING_KEY, message ); log.info(注册消息发送成功userId{}, userId); } }convertAndSend底层会把对象序列化成字节流默认序列化器是 JDK 序列化。为了后续排查方便我建议把默认序列化器换成 Jackson配一个Jackson2JsonMessageConverter这样在 RabbitMQ 控制台里直接能看到 JSON 格式的消息内容。Bean public MessageConverter messageConverter() { return new Jackson2JsonMessageConverter(); }换成这个消息转换器之后消费者的反序列化也要配套处理具体写法在 3.4 里。3.4 消费者重点幂等、手动 ACK、失败重试、死信兜底消费者是整个链路里坑最多的一块。我直接给一份完整可跑的代码每一段都在后面讲清为什么这么写。Component Slf4j public class RegisterConsumer { Resource private StringRedisTemplate stringRedisTemplate; RabbitListener(queues RabbitConfig.REGISTER_QUEUE) public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag message.getMessageProperties().getDeliveryTag(); RegisterMessage registerMessage null; try { // 1. 反序列化消息体 registerMessage (RegisterMessage) message.getPayload(); // 2. 幂等去重 String idempotentKey register:msg: registerMessage.getUserId(); Boolean first stringRedisTemplate.opsForValue() .setIfAbsent(idempotentKey, 1, 10, TimeUnit.MINUTES); if (Boolean.FALSE.equals(first)) { log.warn(重复消息直接确认并跳过userId{}, registerMessage.getUserId()); channel.basicAck(deliveryTag, false); return; } // 3. 业务处理发送欢迎邮件 sendWelcomeMail(registerMessage); // 4. 处理成功手动确认 channel.basicAck(deliveryTag, false); log.info(注册通知消费成功userId{}, registerMessage.getUserId()); } catch (Exception e) { log.error(注册通知消费失败deliveryTag{}, deliveryTag, e); // 失败不重回队列直接进入死信交换机兜底避免无限循环重投 channel.basicNack(deliveryTag, false, false); } } private void sendWelcomeMail(RegisterMessage message) { // 模拟发送邮件时偶发异常用于测试重试与死信链路 if (failtest.com.equals(message.getEmail())) { throw new RuntimeException(模拟邮件服务异常); } log.info(向 {} 发送欢迎邮件, message.getEmail()); } }这段代码有几个细节值得展开第一为什么用手动 ACKSpringBoot 默认是自动 ACK也就是消费者收到消息后立刻告诉 RabbitMQ 处理成功不管业务逻辑是否真正执行完。如果自动 ACK 后你的邮件服务抛异常消息已经找不回来了。手动 ACK 把确认权交给我们自己业务成功再确认失败可以决定是重回队列还是走死信。第二幂等性为什么必须做RabbitMQ 官方文档自己都承认消息可能被重复投递。典型场景消费者处理完业务、但还没来得及 ACK 时进程崩溃RabbitMQ 会把这条消息重新投递给别的消费者。所以“同一个 userId 的注册通知被消费两次”不是理论可能是现实必然。我用 Redis 的setIfAbsent做去重因为它是原子操作多个消费者同时处理同一条消息时只会有一个成功写入。第三为什么要basicNack(deliveryTag, false, false)而不是重回队列requeuetrue会把消息放回原队列头如果消费者每次都在同一处抛异常消息会被立刻重新投递然后再次抛异常形成死循环刷爆日志。更合理的方案是失败后进死信队列由专人或者定时任务去处理。死信队列本质是个普通队列绑定到死信交换机上代码在 3.2 已经写了。4. 链路可靠性消息不丢、不重、不乱序的三板斧一个消息队列工程功能跑通只是第一步。生产环境真正考验的是消息能不能可靠地送到消费者手里以及被正确处理。这块我把 RabbitMQ 在注册场景里最容易遇到的可靠性问题拆开讲。4.1 消息会丢在哪三个环节生产者发送消息到交换机、交换机把消息投递到队列、消费者从队列获取消息这三个环节任何一个出错都可能导致消息丢。生产者环节最容易出问题代码里convertAndSend执行完不代表 RabbitMQ 真的收到了因为网络传输可能出问题。交换机环节如果消息没有匹配到任何队列RabbitMQ 默认会直接静默丢弃。消费者环节自动 ACK 模式下业务代码抛异常消息就没了。针对这三个环节有三板斧。4.2 生产端确认ConfirmCallback 和 ReturnCallback在 yml 里开启两个配置spring: rabbitmq: publisher-confirm-type: correlated publisher-returns: truepublisher-confirm-type: correlated会返回消息投递到交换机的确认结果publisher-returns负责处理交换机投递到队列失败的消息。对应的回调要单独配置 RabbitTemplate。Configuration public class RabbitTemplateConfig { Bean public RabbitTemplate rabbitTemplate(ConnectionFactory connectionFactory) { RabbitTemplate rabbitTemplate new RabbitTemplate(connectionFactory); rabbitTemplate.setMandatory(true); rabbitTemplate.setMessageConverter(new Jackson2JsonMessageConverter()); // 消息到达交换机回调 rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (!ack) { log.error(消息投递交换机失败correlationId{}, cause{}, correlationData.getId(), cause); } }); // 消息从交换机到达队列失败回调 rabbitTemplate.setReturnsCallback(returned - { log.error(消息路由失败exchange{}, routingKey{}, message{}, returned.getExchange(), returned.getRoutingKey(), returned.getMessage()); }); return rabbitTemplate; } }注意setMandatory(true)必须开否则交换机路由不到队列时消息会被静默丢弃ReturnCallback 根本不会触发。这套配置能帮你发现“消息已发但没落队列”的隐性丢消息问题。生产端更可靠的做法是在发送时带上CorrelationData作为业务唯一 ID这样 ConfirmCallback 里能把失败的具体业务消息捞出来做补偿。4.3 消费端可靠性手动 ACK 与有限重试的组合消费端我在 3.4 里已经写了手动 ACK 的完整代码。这里补充一个注意点不要重复叠加 Spring 的Retryable到监听器上因为手动 ACK 模式下你本身就是自己控制确认再用重试模板会出现“消息被 Spring 重试 N 次后 ACK 状态混乱”的麻烦。如果你想用 Spring 重试就老老实实用自动 ACK 并提供RetryTemplate二选一不要混着来。我推荐手动 ACK 加死信队列的原因很朴素消费失败时你天然有一个“落盘”位置让消息积压着等人处理比无限重投要安全得多。生产上你只需要在死信队列后面再加一个消费者把死信消息推送到告警平台人工介入处理。4.4 数据库事务与消息发送的一致性问题3.1 里我提了一嘴“先写库再发消息”的写法其实存在一个极小概率的不一致窗口数据库事务提交成功但 MQ 消息发送时网络异常注册用户没有收到通知。反过来更危险如果消息发送先于事务提交且事务回滚消费者却已经处理了消息就出现了虚假用户。成熟项目的解法一般是“本地消息表”和“事务消息”两条路。本地消息表的核心思路在业务库建一张outbox_message表用户注册时在同一事务里写入用户数据和待发送消息记录。提交后由独立任务轮询消息表发送成功的标记为已发送。事务消息则是把消息的发送预提交到 MQ事务提交后再确认发送RabbitMQ 原生不支持 Kafka 那种完整事务消息需要用分布式事务中间件工程成本偏高。最小闭环里注册场景对通知即时性要求不高先写库再发消息够用但你要清楚它不是绝对强一致。5. 工程实战速查常见问题与排查技巧这一节把我在这个链路里真实踩过、以及帮别人排查过的典型问题整理成速查表每一条都是线上出现过的问题。现象可能原因排查方向解决方案消费者没报错但消息就是没处理消费者把队列名写错了或没绑定检查控制台 Queues 页面看队列消费者数量是否为 0核对消费者RabbitListener(queues...)和配置类队列名消息消费失败后疯狂刷日志失败后requeuetrue无限重回队列看日志时间间隔是否极短、循环重复改为basicNack(tag, false, false)进死信队列消息在队列里堆积消费者不消费消费者线程数配置太低或消费性能瓶颈看 RabbitMQ 控制台消费者数和消费速率调大spring.rabbitmq.listener.simple.concurrency或增加消费者实例重启消费者后部分消息丢失使用了自动 ACK检查配置文件 acknowledge-mode改成 manual业务成功后再确认消息消费重复消费者 ACK 前崩溃或超时MQ 重新投递看业务日志同一 userId 出现多次消费端做幂等推荐 RedissetIfAbsent控制台看不到消息内容只有二进制乱码默认 JDK 序列化器看消费者和生产者是否都配了 Jackson 转换器统一配置Jackson2JsonMessageConverter发送时消息丢失且无任何日志未开启 confirm 和 return 回调查看 yml 配置开启publisher-confirm-type: correlated和 returns5.1 消息重复消费是“必须解决”而不是“偶尔遇到”我发现不少初学者会把幂等当成性能优化觉得“重复就重复呗邮件多发一次又没事”。但注册场景里消息后续往往还要接积分增加、优惠券发放积分发两次就是事故。所以幂等不是附加能力是消费端的第一需求。用 RedissetIfAbsent做消息幂等时要注意过期时间不能太短否则长事务处理到一半 key 过期了第二个消费者就能再次处理也不能太长占用 Redis 内存。注册场景 10 分钟过期比较合理。5.2 手动 ACK 的一个隐蔽坑忘掉确认导致连接阻塞手动 ACK 模式下RabbitMQ 有一个“未确认消息数量上限”的概念具体值是prefetch。如果代码里漏掉了basicAck或basicNack消息会一直处于 Unacked 状态不断堆积。当未确认数量达到prefetch上限后消费者不会再收到新消息看起来就像“消费者死了”。所以无论业务执行成功还是失败finally块里一定要给消息一个明确的确认或拒绝结果。上面 3.4 的代码里我把 ACK 放在 try 里、NACK 放在 catch 里看着是覆盖了正常和异常两条路但更稳妥的写法是用finally 一个状态变量来统一处理。这里不贴冗余代码了记住这个原则每个RabbitListener都要确保所有路径都能到达 ACK/NACK。6. 从小闭环到工程级这套链路还能怎么扩展到这里一条从用户注册到异步通知的完整链路已经跑通了。它短小但该有的东西一样不少生产者、交换机、队列、消费者、幂等、ACK、重试、死信。接下来想把它往更深的地方练大概有三个方向。第一个方向是延迟队列。注册后 30 分钟未激活的用户发送提醒用 RabbitMQ 的延迟消息插件或者死信机制都能实现消息队列的用法就能从“异步”延伸到“定时”。第二个方向是多消费者的业务拆分注册消息进入一个 fanout 交换机邮件服务、积分服务、埋点服务各自绑定自己的队列彻底体验一下解耦的好处。第三个方向是链路追踪生产消息和消费日志打上统一的 traceId用它把一条消息从发起到落地的完整路径串起来这样线上排查问题就能定位到具体是生产端丢了、还是消费端处理慢了。当年我自己从 CRUD 进阶到工程链路最受益的不是看一堆高深理论而是把一个基于消息队列的异步场景自己完整敲了三遍。第一遍完全照抄第二遍默写第三遍不看任何参考从零搭项目。这个过程里遇到的每一个序列化报错、每一次消息丢失、每一轮死循环重投都在帮你建立真正的工程直觉。把文章里这个注册链路跑通再亲手造几个故障去排查你手里的牌就和“只会写同步接口”的阶段完全不同了。