欢迎光临
我们一直在努力

Spring Boot 4 整合 RabbitMQ:从消息发送到可靠投递、重试与死信队列

Spring Boot 4 整合 RabbitMQ:从消息发送到可靠投递、重试与死信队列

RabbitMQ 是一个基于 AMQP 协议的消息中间件,常用于系统解耦、异步处理、流量削峰、延迟执行和事件驱动架构。

Spring Boot 通过 spring-boot-starter-amqp 提供 RabbitMQ 自动配置,可以直接注入 RabbitTemplate 发送消息,并通过 @RabbitListener 消费消息。

本文将完成一个订单事件示例,并覆盖:

  • RabbitMQ 核心概念;
  • Spring Boot 项目搭建;
  • Exchange、Queue 和 Binding 声明;
  • JSON 消息发送与消费;
  • Publisher Confirm;
  • Publisher Return;
  • 消费者确认;
  • 消费失败重试;
  • 死信队列;
  • 消费幂等;
  • 常见问题;
  • 生产环境可靠性设计。

本文示例基于:

组件版本
Java 21
Spring Boot 4.1.0
Spring AMQP 由 Spring Boot 管理
RabbitMQ 4.x
Maven 3.9+

Spring Boot 3.x 用户也可以参考本文。需要注意的是,Spring AMQP 4 使用 Jackson 3 的 JacksonJsonMessageConverter;Spring AMQP 3 通常使用 Jackson2JsonMessageConverter。

官方文档:

  • Spring Boot AMQP
  • Spring AMQP Reference
  • RabbitMQ Documentation

一、为什么使用 RabbitMQ

假设一个订单创建接口需要依次完成:

保存订单
→ 扣减库存
→ 发送短信
→ 发送邮件
→ 增加积分
→ 记录审计日志

如果全部同步执行,任何一个下游服务变慢,都会影响订单接口响应时间。

引入 RabbitMQ 后,可以调整为:

保存订单
→ 发送“订单已创建”消息
→ 立即返回

RabbitMQ
├─ 库存服务消费
├─ 短信服务消费
├─ 邮件服务消费
├─ 积分服务消费
└─ 审计服务消费

这样可以实现:

  • 系统解耦;
  • 异步处理;
  • 流量削峰;
  • 失败重试;
  • 独立扩容;
  • 事件驱动。

但 RabbitMQ 并不会自动解决所有一致性问题。生产环境还需要考虑消息丢失、重复消费、消费失败、消息堆积和数据库与消息的双写一致性。

二、RabbitMQ 核心概念

一条消息通常经过以下链路:

#mermaid-svg-URDOEXClk8vk2ZVX{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-URDOEXClk8vk2ZVX .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-URDOEXClk8vk2ZVX .error-icon{fill:#552222;}#mermaid-svg-URDOEXClk8vk2ZVX .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-URDOEXClk8vk2ZVX .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-URDOEXClk8vk2ZVX .marker{fill:#333333;stroke:#333333;}#mermaid-svg-URDOEXClk8vk2ZVX .marker.cross{stroke:#333333;}#mermaid-svg-URDOEXClk8vk2ZVX svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-URDOEXClk8vk2ZVX p{margin:0;}#mermaid-svg-URDOEXClk8vk2ZVX .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-URDOEXClk8vk2ZVX .cluster-label text{fill:#333;}#mermaid-svg-URDOEXClk8vk2ZVX .cluster-label span{color:#333;}#mermaid-svg-URDOEXClk8vk2ZVX .cluster-label span p{background-color:transparent;}#mermaid-svg-URDOEXClk8vk2ZVX .label text,#mermaid-svg-URDOEXClk8vk2ZVX span{fill:#333;color:#333;}#mermaid-svg-URDOEXClk8vk2ZVX .node rect,#mermaid-svg-URDOEXClk8vk2ZVX .node circle,#mermaid-svg-URDOEXClk8vk2ZVX .node ellipse,#mermaid-svg-URDOEXClk8vk2ZVX .node polygon,#mermaid-svg-URDOEXClk8vk2ZVX .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-URDOEXClk8vk2ZVX .rough-node .label text,#mermaid-svg-URDOEXClk8vk2ZVX .node .label text,#mermaid-svg-URDOEXClk8vk2ZVX .image-shape .label,#mermaid-svg-URDOEXClk8vk2ZVX .icon-shape .label{text-anchor:middle;}#mermaid-svg-URDOEXClk8vk2ZVX .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-URDOEXClk8vk2ZVX .rough-node .label,#mermaid-svg-URDOEXClk8vk2ZVX .node .label,#mermaid-svg-URDOEXClk8vk2ZVX .image-shape .label,#mermaid-svg-URDOEXClk8vk2ZVX .icon-shape .label{text-align:center;}#mermaid-svg-URDOEXClk8vk2ZVX .node.clickable{cursor:pointer;}#mermaid-svg-URDOEXClk8vk2ZVX .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-URDOEXClk8vk2ZVX .arrowheadPath{fill:#333333;}#mermaid-svg-URDOEXClk8vk2ZVX .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-URDOEXClk8vk2ZVX .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-URDOEXClk8vk2ZVX .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-URDOEXClk8vk2ZVX .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-URDOEXClk8vk2ZVX .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-URDOEXClk8vk2ZVX .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-URDOEXClk8vk2ZVX .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-URDOEXClk8vk2ZVX .cluster text{fill:#333;}#mermaid-svg-URDOEXClk8vk2ZVX .cluster span{color:#333;}#mermaid-svg-URDOEXClk8vk2ZVX div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-URDOEXClk8vk2ZVX .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-URDOEXClk8vk2ZVX rect.text{fill:none;stroke-width:0;}#mermaid-svg-URDOEXClk8vk2ZVX .icon-shape,#mermaid-svg-URDOEXClk8vk2ZVX .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-URDOEXClk8vk2ZVX .icon-shape p,#mermaid-svg-URDOEXClk8vk2ZVX .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-URDOEXClk8vk2ZVX .icon-shape .label rect,#mermaid-svg-URDOEXClk8vk2ZVX .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-URDOEXClk8vk2ZVX .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-URDOEXClk8vk2ZVX .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-URDOEXClk8vk2ZVX :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

Producer 生产者

Exchange 交换机

Binding 路由规则

Queue 队列

Consumer 消费者

1. Producer

Producer 是消息生产者,负责把消息发送到 RabbitMQ。

在 Spring Boot 中通常使用:

rabbitTemplate.convertAndSend(exchange, routingKey, message);

2. Exchange

Exchange 接收生产者发送的消息,然后根据 Exchange 类型、Routing Key 和 Binding 将消息路由到队列。

常见 Exchange 类型:

类型路由规则典型场景
Direct Routing Key 完全匹配 订单、支付等精确事件
Topic 支持 *、# 通配符 业务事件订阅
Fanout 忽略 Routing Key,广播到所有绑定队列 广播通知
Headers 根据消息 Header 匹配 复杂属性路由

3. Queue

Queue 用于保存等待消费的消息。

生产环境中的重要消息通常需要:

  • 队列设置为 durable;
  • 消息设置为 persistent;
  • 开启 Publisher Confirm;
  • 消费成功后再确认;
  • 消费逻辑实现幂等。

RabbitMQ 官方指出,durable 队列只能保证队列定义可以在重启后恢复;要恢复队列中的消息,还必须将消息发布为 persistent。RabbitMQ Queues

4. Binding

Binding 用于建立 Exchange 和 Queue 之间的路由关系。

例如:

Exchange:demo.order.exchange
Routing Key:order.created
Queue:demo.order.created.queue

只有符合 Binding 规则的消息才会进入对应队列。

5. Consumer

Consumer 从 Queue 获取消息并执行业务逻辑。

Spring Boot 可以通过 @RabbitListener 创建消费者:

@RabbitListener(queues = "demo.order.created.queue")
public void consume(OrderCreatedEvent event) {
// 处理消息
}

三、使用 Docker 启动 RabbitMQ

创建 docker-compose.yml:

services:
rabbitmq:
image: rabbitmq:4management
container_name: rabbitmqdemo
hostname: rabbitmqdemo
restart: unlessstopped
ports:
"5672:5672"
"15672:15672"
environment:
RABBITMQ_DEFAULT_USER: app
RABBITMQ_DEFAULT_PASS: app123456
volumes:
rabbitmq_data:/var/lib/rabbitmq
healthcheck:
test: ["CMD", "rabbitmq-diagnostics", "-q", "ping"]
interval: 10s
timeout: 5s
retries: 5

volumes:
rabbitmq_data:

启动 RabbitMQ:

docker compose up -d

查看容器状态:

docker compose ps

管理后台地址:

http://localhost:15672

登录信息:

用户名:app
密码:app123456

RabbitMQ 端口:

端口用途
5672 AMQP 客户端连接
15672 RabbitMQ Management 管理页面

四、创建 Spring Boot 项目

项目结构如下:

rabbitmq-demo
├─ pom.xml
├─ docker-compose.yml
└─ src
└─ main
├─ java
│ └─ com.example.rabbitmqdemo
│ ├─ RabbitMqDemoApplication.java
│ ├─ config
│ │ └─ RabbitMqConfig.java
│ ├─ consumer
│ │ └─ OrderEventConsumer.java
│ ├─ event
│ │ └─ OrderCreatedEvent.java
│ ├─ producer
│ │ └─ OrderEventPublisher.java
│ └─ web
│ ├─ CreateOrderRequest.java
│ └─ OrderController.java
└─ resources
└─ application.yml

五、添加 Maven 依赖

pom.xml:

<?xml version="1.0" encoding="UTF-8"?>
<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
xsi:schemaLocation="
http://maven.apache.org/POM/4.0.0
https://maven.apache.org/xsd/maven-4.0.0.xsd"
>

<modelVersion>4.0.0</modelVersion>

<parent>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-parent</artifactId>
<version>4.1.0</version>
<relativePath/>
</parent>

<groupId>com.example</groupId>
<artifactId>rabbitmq-demo</artifactId>
<version>1.0.0</version>
<name>rabbitmq-demo</name>

<properties>
<java.version>21</java.version>
</properties>

<dependencies>
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-web</artifactId>
</dependency>

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-test</artifactId>
<scope>test</scope>
</dependency>

<dependency>
<groupId>org.springframework.amqp</groupId>
<artifactId>spring-rabbit-test</artifactId>
<scope>test</scope>
</dependency>
</dependencies>

<build>
<plugins>
<plugin>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-maven-plugin</artifactId>
</plugin>
</plugins>
</build>

</project>

核心依赖是:

<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>

引入后,Spring Boot 会自动配置:

  • RabbitMQ ConnectionFactory;
  • RabbitTemplate;
  • AmqpAdmin;
  • RabbitListenerContainerFactory;
  • @RabbitListener 相关基础设施。

六、配置 RabbitMQ 连接

创建 application.yml:

server:
port: 8080

spring:
application:
name: rabbitmqdemo

rabbitmq:
host: localhost
port: 5672
username: app
password: app123456
virtual-host: /

connection-timeout: 5s
requested-heartbeat: 30s

# 开启生产者关联确认
publisher-confirm-type: correlated

# 开启无法路由消息的退回通知
publisher-returns: true

template:
# 消息无法路由到任何队列时返回给生产者
mandatory: true

listener:
simple:
# AUTO 表示监听方法正常返回后由容器确认消息
acknowledge-mode: auto

# 消费者并发数
concurrency: 2
max-concurrency: 8

# 每个消费者允许同时持有的未确认消息数量
prefetch: 20

# 失败消息不无限重新入队
default-requeue-rejected: false

retry:
enabled: true
initial-interval: 1s
multiplier: 2
max-interval: 10s
max-retries: 3

logging:
level:
com.example.rabbitmqdemo: info

关键配置说明:

配置作用
publisher-confirm-type: correlated 开启 Publisher Confirm,并关联每条消息
publisher-returns: true 开启无法路由消息的返回通知
template.mandatory: true 消息无法进入任何队列时触发 Return
acknowledge-mode: auto 业务方法正常返回后确认消息
prefetch: 20 限制单个消费者未确认消息数量
default-requeue-rejected: false 防止失败消息无限重新入队
listener.simple.retry.enabled 开启消费重试

publisher-confirm-type 和 publisher-returns 是 Spring Boot 官方提供的 RabbitMQ 配置项。Spring Boot Application Properties

七、创建启动类

package com.example.rabbitmqdemo;

import org.springframework.boot.SpringApplication;
import org.springframework.boot.autoconfigure.SpringBootApplication;

@SpringBootApplication
public class RabbitMqDemoApplication {

public static void main(String[] args) {
SpringApplication.run(RabbitMqDemoApplication.class, args);
}
}

Spring Boot 自动配置已经启用 @RabbitListener,通常不需要额外添加 @EnableRabbit。

八、定义订单事件

创建 OrderCreatedEvent.java:

package com.example.rabbitmqdemo.event;

import java.math.BigDecimal;
import java.time.Instant;

public record OrderCreatedEvent(
String messageId,
String orderNo,
Long userId,
BigDecimal amount,
Instant createdAt
) {
}

建议业务事件至少包含:

  • 消息唯一 ID;
  • 业务唯一 ID;
  • 事件类型;
  • 事件发生时间;
  • 事件版本;
  • 必要业务字段。

生产环境可以进一步增加:

{
"messageId": "8580f765-f424-45e7-ae69-113f02681722",
"eventType": "ORDER_CREATED",
"eventVersion": 1,
"occurredAt": "2026-08-01T10:00:00Z",
"data": {
"orderNo": "SO202608010001"
}
}

事件版本有利于生产者和消费者独立升级。

九、声明 Exchange、Queue 和 Binding

创建 RabbitMqConfig.java:

package com.example.rabbitmqdemo.config;

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import org.springframework.amqp.core.Binding;
import org.springframework.amqp.core.BindingBuilder;
import org.springframework.amqp.core.DirectExchange;
import org.springframework.amqp.core.ExchangeBuilder;
import org.springframework.amqp.core.Queue;
import org.springframework.amqp.core.QueueBuilder;
import org.springframework.amqp.support.converter.JacksonJsonMessageConverter;
import org.springframework.amqp.support.converter.MessageConverter;
import org.springframework.beans.factory.annotation.Qualifier;
import org.springframework.boot.amqp.autoconfigure.RabbitTemplateCustomizer;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;

@Configuration
public class RabbitMqConfig {

private static final Logger log =
LoggerFactory.getLogger(RabbitMqConfig.class);

public static final String ORDER_EXCHANGE =
"demo.order.exchange";

public static final String ORDER_QUEUE =
"demo.order.created.queue";

public static final String ORDER_ROUTING_KEY =
"order.created";

public static final String DEAD_LETTER_EXCHANGE =
"demo.order.dead.exchange";

public static final String DEAD_LETTER_QUEUE =
"demo.order.dead.queue";

public static final String DEAD_LETTER_ROUTING_KEY =
"order.dead";

@Bean
public DirectExchange orderExchange() {
return ExchangeBuilder
.directExchange(ORDER_EXCHANGE)
.durable(true)
.build();
}

@Bean
public Queue orderQueue() {
return QueueBuilder
.durable(ORDER_QUEUE)
.withArgument(
"x-dead-letter-exchange",
DEAD_LETTER_EXCHANGE
)
.withArgument(
"x-dead-letter-routing-key",
DEAD_LETTER_ROUTING_KEY
)
.build();
}

@Bean
public Binding orderBinding(
@Qualifier("orderQueue") Queue queue,
@Qualifier("orderExchange") DirectExchange exchange) {

return BindingBuilder
.bind(queue)
.to(exchange)
.with(ORDER_ROUTING_KEY);
}

@Bean
public DirectExchange deadLetterExchange() {
return ExchangeBuilder
.directExchange(DEAD_LETTER_EXCHANGE)
.durable(true)
.build();
}

@Bean
public Queue deadLetterQueue() {
return QueueBuilder
.durable(DEAD_LETTER_QUEUE)
.build();
}

@Bean
public Binding deadLetterBinding(
@Qualifier("deadLetterQueue") Queue queue,
@Qualifier("deadLetterExchange")
DirectExchange exchange) {

return BindingBuilder
.bind(queue)
.to(exchange)
.with(DEAD_LETTER_ROUTING_KEY);
}

@Bean
public MessageConverter rabbitMessageConverter() {
return new JacksonJsonMessageConverter(
"com.example.rabbitmqdemo"
);
}

@Bean
public RabbitTemplateCustomizer rabbitTemplateCustomizer() {
return rabbitTemplate -> {
rabbitTemplate.setConfirmCallback(
(correlationData, ack, cause) -> {
String messageId =
correlationData == null
? "unknown"
: correlationData.getId();

if (ack) {
log.info(
"RabbitMQ 已确认消息,messageId={}",
messageId
);
} else {
log.error(
"RabbitMQ 拒绝消息,messageId={},cause={}",
messageId,
cause
);
}
}
);

rabbitTemplate.setReturnsCallback(returned -> {
log.error(
"消息无法路由,exchange={},routingKey={},replyCode={},replyText={}",
returned.getExchange(),
returned.getRoutingKey(),
returned.getReplyCode(),
returned.getReplyText()
);
});
};
}
}

应用启动时,Spring Boot 的 AmqpAdmin 会根据这些 Bean 自动声明:

demo.order.exchange
demo.order.created.queue
demo.order.dead.exchange
demo.order.dead.queue

主队列消费失败后,消息将按照下面的链路进入死信队列:

#mermaid-svg-NEwXYC42EhRJxgeG{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;fill:#333;}@keyframes edge-animation-frame{from{stroke-dashoffset:0;}}@keyframes dash{to{stroke-dashoffset:0;}}#mermaid-svg-NEwXYC42EhRJxgeG .edge-animation-slow{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 50s linear infinite;stroke-linecap:round;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-animation-fast{stroke-dasharray:9,5!important;stroke-dashoffset:900;animation:dash 20s linear infinite;stroke-linecap:round;}#mermaid-svg-NEwXYC42EhRJxgeG .error-icon{fill:#552222;}#mermaid-svg-NEwXYC42EhRJxgeG .error-text{fill:#552222;stroke:#552222;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-thickness-normal{stroke-width:1px;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-thickness-thick{stroke-width:3.5px;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-pattern-solid{stroke-dasharray:0;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-thickness-invisible{stroke-width:0;fill:none;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-pattern-dashed{stroke-dasharray:3;}#mermaid-svg-NEwXYC42EhRJxgeG .edge-pattern-dotted{stroke-dasharray:2;}#mermaid-svg-NEwXYC42EhRJxgeG .marker{fill:#333333;stroke:#333333;}#mermaid-svg-NEwXYC42EhRJxgeG .marker.cross{stroke:#333333;}#mermaid-svg-NEwXYC42EhRJxgeG svg{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:16px;}#mermaid-svg-NEwXYC42EhRJxgeG p{margin:0;}#mermaid-svg-NEwXYC42EhRJxgeG .label{font-family:\”trebuchet ms\”,verdana,arial,sans-serif;color:#333;}#mermaid-svg-NEwXYC42EhRJxgeG .cluster-label text{fill:#333;}#mermaid-svg-NEwXYC42EhRJxgeG .cluster-label span{color:#333;}#mermaid-svg-NEwXYC42EhRJxgeG .cluster-label span p{background-color:transparent;}#mermaid-svg-NEwXYC42EhRJxgeG .label text,#mermaid-svg-NEwXYC42EhRJxgeG span{fill:#333;color:#333;}#mermaid-svg-NEwXYC42EhRJxgeG .node rect,#mermaid-svg-NEwXYC42EhRJxgeG .node circle,#mermaid-svg-NEwXYC42EhRJxgeG .node ellipse,#mermaid-svg-NEwXYC42EhRJxgeG .node polygon,#mermaid-svg-NEwXYC42EhRJxgeG .node path{fill:#ECECFF;stroke:#9370DB;stroke-width:1px;}#mermaid-svg-NEwXYC42EhRJxgeG .rough-node .label text,#mermaid-svg-NEwXYC42EhRJxgeG .node .label text,#mermaid-svg-NEwXYC42EhRJxgeG .image-shape .label,#mermaid-svg-NEwXYC42EhRJxgeG .icon-shape .label{text-anchor:middle;}#mermaid-svg-NEwXYC42EhRJxgeG .node .katex path{fill:#000;stroke:#000;stroke-width:1px;}#mermaid-svg-NEwXYC42EhRJxgeG .rough-node .label,#mermaid-svg-NEwXYC42EhRJxgeG .node .label,#mermaid-svg-NEwXYC42EhRJxgeG .image-shape .label,#mermaid-svg-NEwXYC42EhRJxgeG .icon-shape .label{text-align:center;}#mermaid-svg-NEwXYC42EhRJxgeG .node.clickable{cursor:pointer;}#mermaid-svg-NEwXYC42EhRJxgeG .root .anchor path{fill:#333333!important;stroke-width:0;stroke:#333333;}#mermaid-svg-NEwXYC42EhRJxgeG .arrowheadPath{fill:#333333;}#mermaid-svg-NEwXYC42EhRJxgeG .edgePath .path{stroke:#333333;stroke-width:2.0px;}#mermaid-svg-NEwXYC42EhRJxgeG .flowchart-link{stroke:#333333;fill:none;}#mermaid-svg-NEwXYC42EhRJxgeG .edgeLabel{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-NEwXYC42EhRJxgeG .edgeLabel p{background-color:rgba(232,232,232, 0.8);}#mermaid-svg-NEwXYC42EhRJxgeG .edgeLabel rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-NEwXYC42EhRJxgeG .labelBkg{background-color:rgba(232, 232, 232, 0.5);}#mermaid-svg-NEwXYC42EhRJxgeG .cluster rect{fill:#ffffde;stroke:#aaaa33;stroke-width:1px;}#mermaid-svg-NEwXYC42EhRJxgeG .cluster text{fill:#333;}#mermaid-svg-NEwXYC42EhRJxgeG .cluster span{color:#333;}#mermaid-svg-NEwXYC42EhRJxgeG div.mermaidTooltip{position:absolute;text-align:center;max-width:200px;padding:2px;font-family:\”trebuchet ms\”,verdana,arial,sans-serif;font-size:12px;background:hsl(80, 100%, 96.2745098039%);border:1px solid #aaaa33;border-radius:2px;pointer-events:none;z-index:100;}#mermaid-svg-NEwXYC42EhRJxgeG .flowchartTitleText{text-anchor:middle;font-size:18px;fill:#333;}#mermaid-svg-NEwXYC42EhRJxgeG rect.text{fill:none;stroke-width:0;}#mermaid-svg-NEwXYC42EhRJxgeG .icon-shape,#mermaid-svg-NEwXYC42EhRJxgeG .image-shape{background-color:rgba(232,232,232, 0.8);text-align:center;}#mermaid-svg-NEwXYC42EhRJxgeG .icon-shape p,#mermaid-svg-NEwXYC42EhRJxgeG .image-shape p{background-color:rgba(232,232,232, 0.8);padding:2px;}#mermaid-svg-NEwXYC42EhRJxgeG .icon-shape .label rect,#mermaid-svg-NEwXYC42EhRJxgeG .image-shape .label rect{opacity:0.5;background-color:rgba(232,232,232, 0.8);fill:rgba(232,232,232, 0.8);}#mermaid-svg-NEwXYC42EhRJxgeG .label-icon{display:inline-block;height:1em;overflow:visible;vertical-align:-0.125em;}#mermaid-svg-NEwXYC42EhRJxgeG .node .label-icon path{fill:currentColor;stroke:revert;stroke-width:revert;}#mermaid-svg-NEwXYC42EhRJxgeG :root{–mermaid-font-family:\”trebuchet ms\”,verdana,arial,sans-serif;}

处理失败且不重新入队

订单生产者

demo.order.exchange

demo.order.created.queue

订单消费者

demo.order.dead.exchange

demo.order.dead.queue

本文为了让示例开箱即用,在代码中通过 x-dead-letter-exchange 声明死信交换机。

RabbitMQ 官方更推荐在生产环境中使用 Policy 配置 DLX,因为 Policy 可以动态调整;硬编码的队列参数通常需要删除并重新声明队列才能修改。RabbitMQ Dead Letter Exchanges

十、创建消息生产者

创建 OrderEventPublisher.java:

package com.example.rabbitmqdemo.producer;

import com.example.rabbitmqdemo.config.RabbitMqConfig;
import com.example.rabbitmqdemo.event.OrderCreatedEvent;
import org.springframework.amqp.core.MessageDeliveryMode;
import org.springframework.amqp.rabbit.connection.CorrelationData;
import org.springframework.amqp.rabbit.core.RabbitTemplate;
import org.springframework.stereotype.Service;

@Service
public class OrderEventPublisher {

private final RabbitTemplate rabbitTemplate;

public OrderEventPublisher(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}

public void publish(OrderCreatedEvent event) {
CorrelationData correlationData =
new CorrelationData(event.messageId());

rabbitTemplate.convertAndSend(
RabbitMqConfig.ORDER_EXCHANGE,
RabbitMqConfig.ORDER_ROUTING_KEY,
event,
message -> {
message.getMessageProperties()
.setMessageId(event.messageId());

message.getMessageProperties()
.setDeliveryMode(
MessageDeliveryMode.PERSISTENT
);

message.getMessageProperties()
.setHeader(
"businessOrderNo",
event.orderNo()
);

return message;
},
correlationData
);
}
}

这里做了三件事:

  • 使用 CorrelationData 关联 Publisher Confirm;
  • 把 messageId 写入消息属性;
  • 将消息设置为 persistent。
  • 需要注意:

    Publisher Confirm 只说明 RabbitMQ 已经接收并承担了消息责任,不表示消费者已经处理成功。

    Publisher Confirm 和消费者 ACK 是两套相互独立的机制。RabbitMQ Consumer Acknowledgements and Publisher Confirms

    十一、创建消费者

    创建 OrderEventConsumer.java:

    package com.example.rabbitmqdemo.consumer;

    import com.example.rabbitmqdemo.config.RabbitMqConfig;
    import com.example.rabbitmqdemo.event.OrderCreatedEvent;
    import org.slf4j.Logger;
    import org.slf4j.LoggerFactory;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;

    import java.util.Set;
    import java.util.concurrent.ConcurrentHashMap;

    @Component
    public class OrderEventConsumer {

    private static final Logger log =
    LoggerFactory.getLogger(OrderEventConsumer.class);

    /*
    * 仅用于演示幂等。
    *
    * 生产环境不能使用内存 Set,
    * 应使用数据库唯一索引或 Redis 等持久化方案。
    */

    private final Set<String> processedMessageIds =
    ConcurrentHashMap.newKeySet();

    @RabbitListener(queues = RabbitMqConfig.ORDER_QUEUE)
    public void consume(OrderCreatedEvent event) {
    String messageId = event.messageId();

    if (!processedMessageIds.add(messageId)) {
    log.warn(
    "检测到重复消息,跳过处理,messageId={}",
    messageId
    );
    return;
    }

    try {
    log.info(
    "开始处理订单消息,messageId={},orderNo={},amount={}",
    messageId,
    event.orderNo(),
    event.amount()
    );

    /*
    * 用于测试消费失败、重试和死信队列。
    */

    if (event.orderNo().startsWith("FAIL")) {
    throw new IllegalStateException(
    "模拟订单处理失败"
    );
    }

    /*
    * 在这里执行业务操作,例如:
    *
    * 1. 扣减库存;
    * 2. 创建发货任务;
    * 3. 发送通知;
    * 4. 写入消费记录。
    */

    log.info(
    "订单消息处理成功,messageId={},orderNo={}",
    messageId,
    event.orderNo()
    );
    } catch (RuntimeException exception) {
    /*
    * 本次处理失败,需要允许后续重试。
    */

    processedMessageIds.remove(messageId);
    throw exception;
    }
    }
    }

    当前配置使用:

    acknowledge-mode: auto

    其行为可以理解为:

    监听方法正常返回
    → Spring AMQP 向 RabbitMQ 发送 ACK
    → RabbitMQ 删除消息

    监听方法抛出异常
    → Spring AMQP 执行有限次数重试
    → 重试仍然失败
    → 拒绝消息且不重新入队
    → RabbitMQ 将消息发送到死信交换机

    AUTO 并不等于 RabbitMQ 的 autoAck=true。

    Spring AMQP 的确认模式包括:

    模式说明
    AUTO 容器根据监听方法是否正常返回发送 ACK 或 NACK
    MANUAL 业务代码通过 Channel 手动确认
    NONE RabbitMQ autoAck=true,消息发送后立即视为成功

    对于大部分业务消费场景,建议优先使用 AUTO,由监听方法抛出异常控制失败流程。

    十二、创建测试接口

    创建 CreateOrderRequest.java:

    package com.example.rabbitmqdemo.web;

    import java.math.BigDecimal;

    public record CreateOrderRequest(
    String orderNo,
    Long userId,
    BigDecimal amount
    ) {
    }

    创建 OrderController.java:

    package com.example.rabbitmqdemo.web;

    import com.example.rabbitmqdemo.event.OrderCreatedEvent;
    import com.example.rabbitmqdemo.producer.OrderEventPublisher;
    import org.springframework.http.HttpStatus;
    import org.springframework.web.bind.annotation.PostMapping;
    import org.springframework.web.bind.annotation.RequestBody;
    import org.springframework.web.bind.annotation.RequestMapping;
    import org.springframework.web.bind.annotation.ResponseStatus;
    import org.springframework.web.bind.annotation.RestController;

    import java.time.Instant;
    import java.util.Map;
    import java.util.UUID;

    @RestController
    @RequestMapping("/api/orders")
    public class OrderController {

    private final OrderEventPublisher publisher;

    public OrderController(OrderEventPublisher publisher) {
    this.publisher = publisher;
    }

    @PostMapping
    @ResponseStatus(HttpStatus.ACCEPTED)
    public Map<String, Object> createOrder(
    @RequestBody CreateOrderRequest request) {

    String messageId = UUID.randomUUID().toString();

    OrderCreatedEvent event = new OrderCreatedEvent(
    messageId,
    request.orderNo(),
    request.userId(),
    request.amount(),
    Instant.now()
    );

    publisher.publish(event);

    return Map.of(
    "messageId", messageId,
    "orderNo", request.orderNo(),
    "status", "QUEUED"
    );
    }
    }

    十三、启动并测试项目

    启动 RabbitMQ:

    docker compose up -d

    启动 Spring Boot:

    mvn spring-boot:run

    1. 测试正常消息

    curl -X POST "http://localhost:8080/api/orders" \\
    -H "Content-Type: application/json" \\
    -d '{
    "orderNo": "SO202608010001",
    "userId": 1001,
    "amount": 199.99
    }'

    Windows PowerShell 可以使用:

    curl.exe X POST "http://localhost:8080/api/orders" `
    H "Content-Type: application/json" `
    d "{`"orderNo`":`"SO202608010001`",`"userId`":1001,`"amount`":199.99}"

    响应示例:

    {
    "messageId": "8580f765-f424-45e7-ae69-113f02681722",
    "orderNo": "SO202608010001",
    "status": "QUEUED"
    }

    日志示例:

    RabbitMQ 已确认消息,messageId=8580f765-f424-45e7-ae69-113f02681722
    开始处理订单消息,messageId=8580f765-f424-45e7-ae69-113f02681722,orderNo=SO202608010001,amount=199.99
    订单消息处理成功,messageId=8580f765-f424-45e7-ae69-113f02681722,orderNo=SO202608010001

    2. 测试消费失败

    发送一个以 FAIL 开头的订单号:

    curl -X POST "http://localhost:8080/api/orders" \\
    -H "Content-Type: application/json" \\
    -d '{
    "orderNo": "FAIL202608010001",
    "userId": 1001,
    "amount": 199.99
    }'

    消费者会抛出异常并执行有限次数重试。

    所有重试失败后,消息会进入:

    demo.order.dead.queue

    可以打开 RabbitMQ 管理页面查看:

    http://localhost:15672

    进入:

    Queues and Streams
    → demo.order.dead.queue

    即可看到死信消息。

    十四、Publisher Confirm 与 Publisher Return 的区别

    1. Publisher Confirm

    Publisher Confirm 用于确认 RabbitMQ 是否接收了消息。

    生产者
    → RabbitMQ
    ← ACK / NACK

    回调参数:

    (correlationData, ack, cause) -> {
    }

    其中:

    • ack=true:RabbitMQ 已接收消息;
    • ack=false:RabbitMQ 无法接收消息;
    • cause:失败原因;
    • correlationData:消息关联数据。

    2. Publisher Return

    Publisher Return 用于检测消息是否无法路由到任何队列。

    生产者
    → Exchange
    → 没有匹配的 Binding
    → 消息返回生产者

    要启用 Return,需要同时配置:

    spring:
    rabbitmq:
    publisher-returns: true
    template:
    mandatory: true

    一个容易误解的情况是:

    消息发送到了存在的 Exchange
    但 Routing Key 无法匹配任何 Queue

    RabbitMQ 仍可能对发布操作发送 Confirm ACK,因为 Exchange 已经接收并处理了发布请求;此时需要依靠 Return 判断消息没有进入队列。

    Spring AMQP 官方文档也明确区分了 Confirm 和 Return。Spring AMQP Publishing

    十五、消费者手动 ACK

    如果业务需要完全控制确认时机,可以使用 MANUAL 模式:

    package com.example.rabbitmqdemo.consumer;

    import com.example.rabbitmqdemo.config.RabbitMqConfig;
    import com.example.rabbitmqdemo.event.OrderCreatedEvent;
    import com.rabbitmq.client.Channel;
    import org.springframework.amqp.core.Message;
    import org.springframework.amqp.rabbit.annotation.RabbitListener;
    import org.springframework.stereotype.Component;

    import java.io.IOException;

    @Component
    public class ManualAckConsumer {

    @RabbitListener(
    queues = RabbitMqConfig.ORDER_QUEUE,
    ackMode = "MANUAL"
    )
    public void consume(
    OrderCreatedEvent event,
    Message message,
    Channel channel) throws IOException {

    long deliveryTag =
    message.getMessageProperties().getDeliveryTag();

    try {
    /*
    * 执行业务操作。
    */

    channel.basicAck(
    deliveryTag,
    false
    );
    } catch (NonRetryableBusinessException exception) {
    /*
    * 永久性错误:
    * 不重新入队,配置了 DLX 时进入死信队列。
    */

    channel.basicNack(
    deliveryTag,
    false,
    false
    );
    } catch (Exception exception) {
    /*
    * 临时性错误:
    * true 表示重新入队。
    *
    * 必须限制重试次数,否则可能形成无限循环。
    */

    channel.basicNack(
    deliveryTag,
    false,
    true
    );
    }
    }

    private static class NonRetryableBusinessException
    extends RuntimeException {
    }
    }

    手动 ACK 时需要注意:

    • ACK 必须在接收消息的同一个 Channel 上执行;
    • 不要对同一条消息重复 ACK;
    • requeue=true 可能造成无限消费;
    • 业务成功后才能 ACK;
    • 消费者进程在 ACK 前退出时,未确认消息会重新投递;
    • 手动 ACK 不会解决重复消费,消费者仍然必须幂等。

    RabbitMQ 的 delivery tag 只在当前 Channel 内有效,在其他 Channel 上 ACK 会导致 unknown delivery tag 错误。RabbitMQ Consumer Acknowledgements

    十六、为什么必须实现消费幂等

    RabbitMQ 的可靠消费通常是“至少一次”语义。

    下面的情况可能产生重复消息:

    消费者完成数据库操作
    → 准备发送 ACK
    → 网络断开
    → RabbitMQ 没有收到 ACK
    → 消息重新投递

    消费者实际上已经处理成功,但 RabbitMQ 不知道,因此会再次投递。

    所以业务逻辑必须保证:

    同一个 messageId 重复消费多次
    最终业务结果与消费一次相同

    数据库唯一索引方案

    创建消息消费记录表:

    CREATE TABLE mq_consumed_message (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    consumer_name VARCHAR(100) NOT NULL,
    message_id VARCHAR(100) NOT NULL,
    consumed_at DATETIME(3) NOT NULL,
    UNIQUE KEY uk_consumer_message (
    consumer_name,
    message_id
    )
    );

    消费时先插入记录:

    INSERT INTO mq_consumed_message (
    consumer_name,
    message_id,
    consumed_at
    ) VALUES (
    'order-created-consumer',
    '<messageId>',
    NOW(3)
    );

    如果唯一索引冲突,说明消息已经处理过。

    推荐将以下操作放在同一个本地事务中:

    插入消费记录
    + 执行业务更新
    + 提交事务

    伪代码:

    @Transactional
    public void handle(OrderCreatedEvent event) {
    boolean inserted = consumedMessageRepository.tryInsert(
    "order-created-consumer",
    event.messageId()
    );

    if (!inserted) {
    return;
    }

    orderService.process(event);
    }

    不能只依赖 RabbitMQ 的 redelivered 标志判断是否重复,因为它只能作为提示,不能替代业务幂等。

    RabbitMQ 官方也建议消费者按照可重复投递进行设计,并优先让业务处理天然幂等。RabbitMQ Reliability Guide

    十七、消费重试与死信队列

    消费失败大致可以分为两类。

    1. 临时性错误

    例如:

    • 数据库连接超时;
    • 下游接口暂时不可用;
    • 网络抖动;
    • 服务限流。

    这类错误可以重试。

    2. 永久性错误

    例如:

    • 消息格式错误;
    • 必填字段缺失;
    • 业务状态不允许;
    • 订单不存在;
    • 数据违反约束。

    这类错误继续重试通常没有意义,应进入死信队列并等待人工或补偿任务处理。

    推荐流程:

    消费消息
    → 处理成功
    → ACK

    → 临时失败
    → 有限次数重试
    → 仍然失败
    → DLQ

    → 永久失败
    → 直接拒绝
    → DLQ

    不要使用下面的无限循环:

    消费失败
    → requeue=true
    → 立即重新消费
    → 再次失败
    → 再次重新入队

    无限重新入队会导致:

    • CPU 和网络资源浪费;
    • 日志快速增长;
    • 正常消息被阻塞;
    • 队列吞吐下降;
    • 故障难以恢复。

    十八、数据库事务与消息发送的一致性

    下面的代码存在双写风险:

    @Transactional
    public void createOrder(Order order) {
    orderRepository.save(order);
    rabbitTemplate.convertAndSend(
    "order.exchange",
    "order.created",
    order
    );
    }

    可能出现:

    情况一

    数据库提交成功
    RabbitMQ 发送失败

    结果:订单存在,但下游没有收到消息。

    情况二

    RabbitMQ 发送成功
    数据库事务回滚

    结果:下游收到一条实际上不存在的订单消息。

    Publisher Confirm 只能确认 RabbitMQ 是否接收消息,不能让数据库事务和 RabbitMQ 事务自动形成一个原子事务。

    推荐使用 Outbox Pattern

    在同一个数据库事务中同时保存业务数据和待发送事件:

    本地事务:
    保存订单
    保存 outbox_event
    提交事务

    后台任务再发送 Outbox:

    扫描待发送事件
    → 发送 RabbitMQ
    → 等待 Publisher Confirm
    → 标记发送成功

    表结构示例:

    CREATE TABLE outbox_event (
    id BIGINT PRIMARY KEY AUTO_INCREMENT,
    event_id VARCHAR(100) NOT NULL,
    event_type VARCHAR(100) NOT NULL,
    aggregate_id VARCHAR(100) NOT NULL,
    payload JSON NOT NULL,
    status VARCHAR(20) NOT NULL,
    retry_count INT NOT NULL DEFAULT 0,
    next_retry_at DATETIME(3),
    created_at DATETIME(3) NOT NULL,
    sent_at DATETIME(3),
    UNIQUE KEY uk_event_id (event_id),
    KEY idx_status_retry (
    status,
    next_retry_at
    )
    );

    Outbox 可以避免数据库提交成功但消息完全没有记录的问题。

    仍然需要消费者幂等,因为下面的情况可能重复发送:

    RabbitMQ 已接收消息
    → Publisher Confirm 返回途中网络中断
    → Outbox 任务认为发送失败
    → 再次发送消息

    可靠消息系统通常采用:

    生产者 Outbox
    + Publisher Confirm
    + 持久化队列和消息
    + 消费者 ACK
    + 消费幂等
    + DLQ
    + 对账和补偿

    十九、生产环境建议

    1. 重要消息使用可靠性闭环

    Outbox
    → Publisher Confirm
    → Publisher Return
    → Durable Queue
    → Persistent Message
    → Consumer ACK
    → Idempotency
    → DLQ
    → Reconciliation

    2. 集群中的关键队列考虑 Quorum Queue

    Classic Queue 默认不会复制队列数据。

    对于订单、支付等不能轻易丢失的消息,可以在 RabbitMQ 集群中评估 Quorum Queue。

    Quorum Queue 更重视:

    • 数据复制;
    • 节点故障恢复;
    • 数据安全;
    • 与 Publisher Confirm 配合。

    它不一定适用于所有低延迟、高瞬时吞吐场景,需要根据业务进行压测。RabbitMQ Quorum Queues

    3. 控制 Prefetch

    Prefetch 过大可能导致:

    • 单个消费者持有大量未确认消息;
    • 内存占用增加;
    • 消费者故障后大量消息重新投递;
    • 多个消费者之间负载不均。

    Prefetch 过小则可能降低吞吐。

    建议根据单条消息处理耗时、并发数和服务资源进行压测,不要直接照搬固定值。

    4. 控制消息大小

    RabbitMQ 更适合传输业务事件,不适合直接传输大型文件。

    不推荐:

    {
    "fileContent": "几十 MB 的 Base64 文件"
    }

    推荐:

    {
    "fileId": "FILE202608010001",
    "objectKey": "orders/2026/08/invoice.pdf"
    }

    文件存储在对象存储中,消息只传递文件标识和地址。

    5. 监控消息堆积

    至少监控:

    • Queue Ready 数量;
    • Queue Unacked 数量;
    • 消息发布速率;
    • 消息消费速率;
    • 消费延迟;
    • Publisher NACK;
    • Publisher Return;
    • 消费失败次数;
    • DLQ 数量;
    • 连接和 Channel 数量;
    • RabbitMQ 内存和磁盘告警。

    6. 建立死信处理流程

    死信队列不是故障处理的终点。

    需要明确:

    谁处理死信
    → 如何查看失败原因
    → 是否允许重新投递
    → 如何避免重复副作用
    → 是否需要人工审批
    → 如何记录补偿结果

    不要在没有校验失败原因的情况下,把 DLQ 中的所有消息直接批量重放。

    7. 使用 TLS 和最小权限

    生产环境建议:

    • 使用 amqps;
    • 不使用默认 guest 用户;
    • 每个应用使用独立账号;
    • 按 Virtual Host 隔离;
    • 只授予需要的 configure、write、read 权限;
    • 密码放入 Secret 管理系统;
    • 限制 RabbitMQ 管理页面访问范围。

    8. 不依赖全局严格顺序

    当一个队列有多个并发消费者时,业务完成顺序不一定与发布顺序完全一致。

    如果业务必须按订单或用户顺序处理,可以考虑:

    • 按业务键分片到固定队列;
    • 同一分片使用单消费者;
    • 在业务层进行版本检查;
    • 使用状态机拒绝旧事件;
    • 评估 RabbitMQ Streams 或其他分区消息系统。

    二十、常见问题

    1. 消息发送成功,但队列没有数据

    检查:

    • Exchange 是否存在;
    • Routing Key 是否正确;
    • Exchange 和 Queue 是否绑定;
    • 是否启用 mandatory;
    • 是否注册 Returns Callback。

    2. 修改队列参数后应用启动失败

    如果队列已经存在,又修改了 durable、DLX 或其他声明参数,RabbitMQ 可能返回:

    PRECONDITION_FAILED – inequivalent arg

    RabbitMQ 不允许用不同参数重新声明同名队列。

    开发环境可以删除旧队列后重新启动。生产环境必须先设计迁移方案,不能直接删除包含业务消息的队列。

    3. 消费失败后一直循环

    通常是因为:

    default-requeue-rejected=true

    或者手动执行了:

    channel.basicNack(
    deliveryTag,
    false,
    true
    );

    需要增加有限重试和 DLQ,避免无限重新入队。

    4. 消息被重复消费

    这是至少一次投递下的正常可能性。

    需要使用:

    • 消息唯一 ID;
    • 数据库唯一索引;
    • 幂等业务更新;
    • 状态机;
    • 乐观锁。

    5. Publisher Confirm 成功是否代表消费成功

    不代表。

    Publisher Confirm
    = RabbitMQ 已接收消息

    Consumer ACK
    = 消费者已完成处理

    两者相互独立。

    6. Durable Queue 是否保证消息不丢失

    仅设置 durable 不够。

    还需要:

    • 消息设置为 persistent;
    • 开启 Publisher Confirm;
    • 使用可靠的队列类型和集群部署;
    • 消费成功后再 ACK;
    • 对未确认消息进行重发;
    • 消费者实现幂等。

    7. Spring Boot 3 找不到 JacksonJsonMessageConverter

    Spring AMQP 4 使用:

    JacksonJsonMessageConverter

    Spring AMQP 3 通常使用:

    Jackson2JsonMessageConverter

    Spring AMQP 4 中的 Jackson 2 转换器已经被标记为弃用,推荐迁移到基于 Jackson 3 的转换器。Spring AMQP Message Converters

    二十一、总结

    Spring Boot 使用 RabbitMQ 的基础代码并不复杂:

    配置连接
    → 声明 Exchange
    → 声明 Queue
    → 创建 Binding
    → RabbitTemplate 发送
    → @RabbitListener 消费

    真正复杂的是生产环境中的消息可靠性。

    一套相对完整的可靠消息方案通常包括:

    数据库与 Outbox 本地事务
    → 后台可靠投递
    → Publisher Confirm
    → Publisher Return
    → Durable 或 Quorum Queue
    → Persistent Message
    → Consumer ACK
    → 有限次数重试
    → 消费幂等
    → Dead Letter Queue
    → 监控告警
    → 对账与补偿

    需要特别记住:

  • Publisher Confirm 不代表消费者处理成功;
  • Durable Queue 不等于消息绝对不丢失;
  • 消费重试必须有上限;
  • 死信队列必须有后续处理流程;
  • 至少一次投递意味着消费者必须幂等;
  • RabbitMQ 不能自动解决数据库与消息的双写一致性;
  • 关键业务应采用 Outbox、对账和补偿形成闭环。
  • 参考资料

    • Spring Boot AMQP
    • Spring Boot RabbitMQ 配置项
    • Spring AMQP Reference
    • Spring AMQP RabbitTemplate
    • Spring AMQP Message Converters
    • RabbitMQ Consumer Acknowledgements and Publisher Confirms
    • RabbitMQ Reliability Guide
    • RabbitMQ Dead Letter Exchanges
    • RabbitMQ Queues
    • RabbitMQ Quorum Queues
    赞(0)
    未经允许不得转载:171主机测评 » Spring Boot 4 整合 RabbitMQ:从消息发送到可靠投递、重试与死信队列
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址