Java 基础体系 · 第 98/100 篇。示例统一以 Java 25 LTS 为语言和 JVM 基线;框架示例使用与其兼容的现代稳定版本。
Java RabbitMQ:Exchange、确认、重试、死信、顺序和幂等
RabbitMQ 生产架构中最容易被混淆的几个概念是:
- Exchange 接收消息,但通常不保存消息;
- Publisher Confirm 只能证明 RabbitMQ 接受了发布,不代表消费者已经处理;
- Consumer Ack 只能证明消费者确认了一次投递,不代表业务事务一定成功;
- 重试解决“暂时失败”,死信解决“无法继续正常处理”;
- RabbitMQ 可以提供局部顺序,但不能自动提供全局顺序;
- 只要存在重试、连接断开、消费者崩溃或发布端超时,就必须按“消息可能重复”设计幂等。
下面从消息流转开始,逐层建立这些机制之间的因果关系。
一、先明确 RabbitMQ 中的消息路径
RabbitMQ 中,一条消息通常经过以下路径:
flowchart LR
P[Producer] -->|publish| E[Exchange]
E -->|binding + routing key| Q[Queue]
Q -->|deliver| C[Consumer]
C -->|basic.ack| Q
Q -->|remove after ack| D[已确认消息]
E -->|unroutable + mandatory| R[Publisher Return]
Q -->|reject/nack requeue=false| X[Dead Letter Exchange]
Q -->|TTL expired / length limit| X
X --> Q2[Retry Queue or Dead Queue]
这里有三个不同的“成功”:
- 发布成功:Exchange 接受了发布请求,通常通过 Publisher Confirm 判断。
- 路由成功:Exchange 找到至少一个匹配的 Queue;找不到时可能触发 Return。
- 业务成功:Consumer 完成数据库、外部服务等业务操作,并发送 Ack。
它们不是同一个状态。一个消息可以:
- 发布成功,但因为没有绑定而未路由;
- 路由成功,但消费者处理失败;
- 消费者业务已提交,但 Ack 前进程崩溃,导致消息再次投递;
- 消费者发送了 Ack,但业务事务其实没有提交,造成消息丢失。
因此,不能用 Publisher Confirm 替代 Consumer Ack,也不能把 Consumer Ack 当作数据库事务提交的自动证明。
二、Exchange:消息为什么能到达某个 Queue
2.1 Exchange 的职责
Producer 不直接把消息写入 Queue,而是把消息发布到 Exchange。Exchange 根据:
- Exchange 类型;
- routing key;
- binding;
- 可选的 headers;
决定消息应该复制到哪些 Queue。
Queue 才负责保存消息和向消费者投递。Exchange 一般是路由组件,不应理解为消息存储区。
RabbitMQ 中常见的 Exchange 类型如下。
2.2 Direct Exchange
Direct Exchange 按精确 routing key 匹配。
例如存在以下绑定:
order.exchange --[routing key = order.created]--> order.queue
只有 routing key 为 order.created 的消息才会进入 order.queue。
rabbitTemplate.convertAndSend(
"order.exchange",
"order.created",
event
);
Direct Exchange 适合:
- 命令型消息;
- 明确的事件类型;
- 一个 routing key 对应一个或多个确定队列。
如果多个 Queue 使用同一个 routing key,它们都会收到消息。Exchange 的路由不是竞争消费,而是复制到每个匹配 Queue。
2.3 Topic Exchange
Topic Exchange 按点分隔的模式匹配 routing key。
例如:
order.created.eu
order.created.cn
order.cancelled.cn
绑定模式:
order.*.cn
order.#
其中:
*匹配一个单词;#匹配零个或多个单词。
因此:
order.*.cn
可以匹配:
order.created.cn
order.cancelled.cn
但不能匹配:
order.created.eu
Topic Exchange 适合事件订阅,但 routing key 应保持稳定。如果把用户输入直接拼成 routing key,可能造成绑定失控或路由规则难以审计。
2.4 Fanout Exchange
Fanout Exchange 忽略 routing key,把消息广播给所有绑定 Queue。
event.exchange
├── audit.queue
├── analytics.queue
└── notification.queue
每个 Queue 都会拥有一份独立消息。一个消费者 Ack 不会影响其他 Queue 中的副本。
2.5 Headers Exchange
Headers Exchange 根据消息 headers 路由,而不是主要依赖 routing key。它能表达更复杂的匹配条件,但使用和排查成本更高。除非路由条件确实不是简单的层级 key,否则 Topic 或 Direct 通常更容易维护。
2.6 Default Exchange
RabbitMQ 内置一个名称为空字符串的 Direct Exchange。Queue 声明后,RabbitMQ 会自动建立:
routing key = queue name
的隐式绑定。因此:
rabbitTemplate.convertAndSend("", "order.queue", event);
可以直接把消息发送到名为 order.queue 的 Queue。
Default Exchange 适合简单场景,但在生产架构中显式命名 Exchange 通常更清楚,因为业务路由、权限和监控都更容易表达。
三、声明 Exchange、Queue 和 Binding
下面使用 Spring Boot 和 Spring AMQP 声明一个订单消息拓扑:
@Configuration
class RabbitTopology {
static final String ORDER_EXCHANGE = "order.x";
static final String ORDER_QUEUE = "order.q";
static final String ORDER_ROUTING_KEY = "order.created";
static final String RETRY_EXCHANGE = "order.retry.x";
static final String RETRY_QUEUE = "order.retry.10s.q";
static final String DEAD_QUEUE = "order.dead.q";
@Bean
DirectExchange orderExchange() {
return new DirectExchange(ORDER_EXCHANGE, true, false);
}
@Bean
DirectExchange retryExchange() {
return new DirectExchange(RETRY_EXCHANGE, true, false);
}
@Bean
Queue orderQueue() {
return QueueBuilder.durable(ORDER_QUEUE)
.withArgument("x-dead-letter-exchange", RETRY_EXCHANGE)
.withArgument("x-dead-letter-routing-key", "order.retry")
.build();
}
@Bean
Binding orderBinding(Queue orderQueue, DirectExchange orderExchange) {
return BindingBuilder.bind(orderQueue)
.to(orderExchange)
.with(ORDER_ROUTING_KEY);
}
@Bean
Queue retryQueue() {
return QueueBuilder.durable(RETRY_QUEUE)
.withArgument("x-message-ttl", 10_000)
.withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
.withArgument("x-dead-letter-routing-key", ORDER_ROUTING_KEY)
.build();
}
@Bean
Binding retryBinding(Queue retryQueue, DirectExchange retryExchange) {
return BindingBuilder.bind(retryQueue)
.to(retryExchange)
.with("order.retry");
}
@Bean
Queue deadQueue() {
return QueueBuilder.durable(DEAD_QUEUE).build();
}
}
这段配置表达的路径是:
order.x
└── order.created
└── order.q
order.retry.x
└── order.retry
└── order.retry.10s.q
└── TTL 到期后回到 order.x
最终无法处理的消息
└── order.dead.q
需要注意几个事实:
- Queue 或 Exchange 已经存在时,重新声明必须使用兼容属性,否则 RabbitMQ 会返回
PRECONDITION_FAILED并关闭通道。 durable=true只表示拓扑在 Broker 重启后保留,不表示消息一定持久化。- 消息还需要设置持久化属性,通常使用
MessageDeliveryMode.PERSISTENT。 - 生产环境中,Queue 参数经常通过 RabbitMQ Policy 管理。硬编码
x-arguments后再修改,可能因为声明不兼容而无法滚动升级。 - DLX 指向的 Exchange、routing key 和目标 Queue 必须实际存在并有正确权限,否则死信可能无法按照预期路由。
四、Publisher Confirm:发布端确认了什么
4.1 Confirm 的准确含义
Publisher Confirm 是 RabbitMQ 对发布者的异步确认机制。
当 Producer 发布消息后,RabbitMQ 会在处理到相应阶段后发送:
ack:Broker 接受了该消息;nack:Broker 没有接受该消息。
Confirm 的确认对象是“发布操作”,不是业务处理结果。
可以把发布过程抽象为:
Producer --publish--> Broker
Producer <--confirm-- Broker
当 Broker 发送 ack 时,通常表示消息已经被 RabbitMQ 接受。对于持久化消息和持久化 Queue,确认时机还受到 RabbitMQ 存储和队列类型等因素影响,但它仍不等于“消费者已经处理”。
4.2 Return 和 Confirm 解决不同问题
如果消息发布到一个不存在的 Exchange,发布通道通常会出现协议错误,不能简单等待正常 Confirm。
如果 Exchange 存在,但没有任何 Queue 能匹配 routing key,则属于“不可路由”:
Producer -> Exchange -> 没有匹配 Queue
此时只有在发布设置 mandatory=true 时,RabbitMQ 才会把消息退回 Producer,Spring AMQP 中表现为 Return。
因此,可靠发布通常同时需要:
- Publisher Confirm:Exchange 是否接受了发布;
- Publisher Return:消息是否没有路由到任何 Queue;
- 对 Exchange 不存在、通道关闭等同步异常进行处理。
4.3 Spring Boot 配置
spring.rabbitmq.publisher-confirm-type=correlated
spring.rabbitmq.publisher-returns=true
spring.rabbitmq.template.mandatory=true
correlated 表示 Confirm 能关联到具体消息。mandatory 让不可路由消息返回给发布者。
示例发布代码:
@Service
class OrderPublisher {
private final RabbitTemplate rabbitTemplate;
OrderPublisher(RabbitTemplate rabbitTemplate) {
this.rabbitTemplate = rabbitTemplate;
}
CompletableFuture<CorrelationData.Confirm> publish(OrderCreated event) {
CorrelationData correlationData =
new CorrelationData(event.eventId());
rabbitTemplate.convertAndSend(
RabbitTopology.ORDER_EXCHANGE,
RabbitTopology.ORDER_ROUTING_KEY,
event,
message -> {
message.getMessageProperties()
.setDeliveryMode(MessageDeliveryMode.PERSISTENT);
message.getMessageProperties()
.setMessageId(event.eventId());
return message;
},
correlationData
);
return correlationData.getFuture();
}
}
调用方可以等待或异步观察 Confirm:
publisher.publish(event)
.whenComplete((confirm, error) -> {
if (error != null) {
// 发布调用本身失败,例如连接或通道异常
recordPublishFailure(event.eventId(), error);
return;
}
if (confirm == null || !confirm.isAck()) {
// Broker nack,或者实现返回了无效确认
recordPublishFailure(
event.eventId(),
new IllegalStateException(
confirm == null ? "no confirm" : confirm.getReason()
)
);
return;
}
recordPublishAccepted(event.eventId());
});
这里的 CompletableFuture 只表示发布确认结果。它不能告诉你:
- 是否有消费者在线;
- 消费者是否已经收到;
- 消费者是否提交了数据库事务;
- 业务是否调用外部系统成功。
4.4 Confirm 丢失时怎么办
假设 Producer 发布后连接断开:
1. Producer 发送消息
2. Broker 可能已经接受
3. Confirm 尚未返回
4. Producer 连接断开
5. Producer 不知道消息究竟是否存在
如果 Producer 直接重发,可能产生重复消息;如果不重发,可能丢失消息。
这不是通过网络协议简单消除的二选一问题。工程上通常选择:
- 给每个事件分配稳定的
eventId; - 发送端记录待确认状态;
- Confirm ack 后标记成功;
- 超时或 nack 后重发;
- 消费端使用幂等处理消除重复影响。
这就是“可靠发布通常是至少一次发布 + 消费端幂等”,而不是试图依靠一次网络调用实现绝对的一次且仅一次。
五、Consumer Ack:RabbitMQ 什么时候删除消息
5.1 自动 Ack 和手动 Ack
消费者获取消息时,可以使用自动确认或手动确认。
自动确认的逻辑大致是:
Broker 投递消息 -> 立即视为已确认
如果消费者刚收到消息就崩溃,消息可能已经从队列中移除,业务尚未执行完成。
手动确认的逻辑是:
Broker 投递消息
-> 消费者执行业务
-> 业务成功后 basic.ack
-> Broker 删除该投递对应的消息
Spring Boot 配置:
spring.rabbitmq.listener.simple.acknowledge-mode=manual
spring.rabbitmq.listener.simple.prefetch=1
消费者示例:
@Component
class OrderConsumer {
private final OrderService orderService;
OrderConsumer(OrderService orderService) {
this.orderService = orderService;
}
@RabbitListener(queues = RabbitTopology.ORDER_QUEUE)
public void consume(Message message, Channel channel) throws IOException {
long deliveryTag =
message.getMessageProperties().getDeliveryTag();
OrderCreated event = parse(message);
try {
orderService.handle(event);
// 业务成功后确认当前一条消息
channel.basicAck(deliveryTag, false);
} catch (TransientBusinessException e) {
// 让消息重新入队;如果没有退避,会快速重复投递
channel.basicNack(deliveryTag, false, true);
} catch (PermanentBusinessException e) {
// 不重新入队,进入 DLX 流程
channel.basicNack(deliveryTag, false, false);
} catch (Exception e) {
// 默认按暂时失败处理,实际项目应根据异常分类
channel.basicNack(deliveryTag, false, true);
}
}
private OrderCreated parse(Message message) {
// 这里可以使用 Jackson 等转换器
throw new UnsupportedOperationException("示例省略反序列化实现");
}
}
5.2 Ack 的作用域
RabbitMQ 的 delivery tag 是按 Channel 分配的。Ack 必须在产生该 delivery tag 的同一 Channel 上发送,否则会出现协议错误。
basicAck(deliveryTag, false) 只确认当前消息。
basicAck(deliveryTag, true) 会确认当前 tag 以及此前同一 Channel 上尚未确认的投递。批量 Ack 可以减少协议交互,但如果中间某条消息业务并未成功,就不能盲目使用 multiple=true。
5.3 Ack 与业务事务的顺序
正确的基本顺序通常是:
读取消息
-> 执行业务事务
-> 数据库提交成功
-> basic.ack
不能反过来:
读取消息
-> basic.ack
-> 数据库提交
后者在数据库提交前进程崩溃时会造成消息丢失。
但即使按正确顺序执行,也仍然存在另一个窗口:
1. 数据库事务提交
2. 进程在发送 Ack 前崩溃
3. RabbitMQ 重新投递
4. 业务再次执行
所以手动 Ack 解决的是“未成功处理时不要过早删除”,不能单独解决重复消费。
六、重试:暂时失败和永久失败必须分开
6.1 什么是暂时失败
暂时失败是指稍后重试可能成功,例如:
- 下游 HTTP 服务临时超时;
- 数据库连接池暂时耗尽;
- 外部服务返回限流;
- 某个依赖正在滚动发布。
永久失败包括:
- 消息格式不合法;
- 业务数据违反不可变约束;
- 必要字段缺失;
- 订单状态不允许当前操作。
永久失败不应无限重试,否则会形成消息热循环。
6.2 直接 requeue=true 的问题
如果消费者这样处理:
channel.basicNack(tag, false, true);
RabbitMQ 会把消息重新入队。没有任何延迟控制时,路径可能变成:
投递 -> 失败 -> 立即入队 -> 立即投递 -> 失败 -> ...
后果包括:
- CPU 持续消耗;
- 日志快速刷屏;
- 同一条坏消息阻塞有效消息;
- 下游服务被更高频率地冲击;
- 队列吞吐下降。
因此,“重新入队”不是完整的重试策略,只是重新安排一次投递。
6.3 TTL + DLX 实现延迟重试
RabbitMQ 没有把任意消息精确睡眠后再投递给消费者的通用语义。常见做法是使用一个带 TTL 的重试 Queue:
消费失败
-> 发布到 retry exchange
-> retry queue 暂存 10 秒
-> TTL 到期
-> retry queue 的 DLX 将消息送回主 exchange
-> 主 queue 再次消费
前面的拓扑中:
QueueBuilder.durable(RETRY_QUEUE)
.withArgument("x-message-ttl", 10_000)
.withArgument("x-dead-letter-exchange", ORDER_EXCHANGE)
.withArgument("x-dead-letter-routing-key", ORDER_ROUTING_KEY)
.build();
TTL 的变量是毫秒数。10_000 表示消息在该 Queue 中至少等待约 10 秒后才具备过期条件,但它不是实时定时器,实际重新投递还受队列调度和 Broker 状态影响。
一个重要边界是 Queue 中消息的过期处理可能受到队首消息影响。使用单一 Queue 存放不同 TTL 的消息,不能把它当作精确的时间轮。需要多个固定延迟 Queue 时,通常按 10 秒、1 分钟、10 分钟等档位分别建立拓扑,或使用 RabbitMQ 延迟消息插件;插件属于额外组件能力,不应当当作所有 RabbitMQ 部署的内置保证。
6.4 “发布到重试 Queue 后再 Ack”与失败窗口
消费者失败后,可以:
1. 将消息发布到 retry exchange
2. 等待 Publisher Confirm
3. Confirm 成功后 Ack 原消息
这比直接 nack(requeue=true) 更容易控制延迟和次数:
void retry(Message message, Channel channel, long tag, int attempt)
throws IOException {
RetryEvent retryEvent = new RetryEvent(
message.getBody(),
attempt + 1,
message.getMessageProperties().getMessageId()
);
CorrelationData correlationData =
new CorrelationData(retryEvent.eventId());
rabbitTemplate.convertAndSend(
RabbitTopology.RETRY_EXCHANGE,
"order.retry",
retryEvent,
correlationData
);
correlationData.getFuture().whenComplete((confirm, error) -> {
if (error == null && confirm != null && confirm.isAck()) {
try {
channel.basicAck(tag, false);
} catch (IOException ackError) {
// Ack 失败时,原消息可能再次投递;
// retryEvent 必须允许重复,不能依赖“恰好一次”
recordAckFailure(retryEvent.eventId(), ackError);
}
} else {
// 重试消息未被确认,不能确认原消息
recordRetryPublishFailure(retryEvent.eventId(), error);
}
});
}
这个流程仍然不是原子操作:
重试消息 Confirm 成功
原消息 Ack 失败
会导致重试消息存在两份。因此重试消息和业务处理都必须幂等。
6.5 重试次数
重试次数不能只依赖客户端内存计数,因为消费者重启后内存状态会丢失。常见方式是:
- 在消息 header 中保存
attempt; - 根据 RabbitMQ 的死信历史头部
x-death推断经历过的死信次数; - 使用单独的重试记录表;
- 把重试次数放入业务事件体。
Header 方案简单,但要注意消息在多次 DLX 转发过程中可能增加或更新死信历史;不同 RabbitMQ 版本和路由链路下,应在目标环境验证实际 header 结构,不应只凭日志猜测。
当达到最大次数后,消费者应将消息转入最终死信 Queue,而不是继续回到主 Queue:
attempt < 5 -> retry queue
attempt >= 5 -> dead queue
七、死信:消息为什么离开正常处理路径
7.1 死信的触发条件
消息成为死信,通常有以下原因:
- 消费者
basic.reject(requeue=false); - 消费者
basic.nack(requeue=false); - 消息在 Queue 中过期;
- Queue 达到长度限制后,旧消息被淘汰;
- 某些队列类型和 RabbitMQ 版本支持的投递次数限制触发。
消息成为死信后,RabbitMQ 会把它重新发布到该 Queue 配置的 Dead Letter Exchange。
因此:
DLX 不是一个神奇的“错误存储区”
它仍然是一次 Exchange 路由。必须为 DLX 配置目标 Queue,否则消息可能没有可见的最终落点。
7.2 消费失败进入 DLX
channel.basicReject(deliveryTag, false);
等价于拒绝当前投递且不重新入队。若原 Queue 设置了:
x-dead-letter-exchange = order.retry.x
消息将被发送到该 Exchange,再按 routing key 路由到 retry Queue 或 dead Queue。
7.3 DLX 的可靠性边界
必须区分:
- 正常 Producer 使用 Publisher Confirm 发布;
- RabbitMQ 内部把死信转发到 DLX。
死信转发并不自动等价于应用 Producer 的端到端 Confirm。RabbitMQ 对不同队列类型、版本和死信配置存在不同可靠性行为;例如某些版本和队列类型支持更强的 at-least-once dead-lettering,但需要显式配置并承担额外资源成本。不能仅凭“配置了 DLX”就断言死信绝不会丢失。
生产上应验证:
- 原 Queue 是否确实配置了 DLX;
- DLX Exchange 是否存在;
- routing key 是否匹配;
- DLX 目标 Queue 是否绑定;
- Broker 日志和权限是否正常;
- 目标队列是否有消息积压;
- RabbitMQ 版本和队列类型对死信确认的具体支持。
最终死信 Queue 不应只是“丢进去就结束”。通常还需要:
- 保存原始消息和 headers;
- 保存失败原因、重试次数和时间;
- 提供人工查看;
- 支持修复后重新投递;
- 防止人工重放绕过幂等检查。
八、完整的状态和故障路径
一个订单事件的正常路径可以表示为:
sequenceDiagram
participant P as Producer
participant E as order.x
participant Q as order.q
participant C as Consumer
participant DB as Database
P->>E: publish(eventId)
E-->>P: confirm ack
E->>Q: route by order.created
Q->>C: deliver(tag=17)
C->>DB: transaction
DB-->>C: commit
C->>Q: basic.ack(17)
Q-->>Q: remove message
消费者崩溃发生在数据库提交之后、Ack 之前:
sequenceDiagram
participant Q as order.q
participant C1 as Consumer 1
participant DB as Database
participant C2 as Consumer 2
Q->>C1: deliver(eventId=E1)
C1->>DB: commit business effect
C1--xQ: process crashes before ack
Q->>C2: redeliver(eventId=E1)
C2->>DB: detect duplicate eventId
C2->>Q: basic.ack
这条路径说明:
数据库已提交 != RabbitMQ 已 Ack
也说明幂等检查必须发生在业务操作的事务边界内,而不能只在内存中维护一个已处理集合。
九、幂等:为什么至少一次系统必须接受重复
9.1 重复消息的来源
重复可能来自多个位置:
Producer 重试
publish -> 网络中断 -> Producer 未收到 Confirm -> 重发
第一次发布可能已经成功,所以 Broker 中可能存在两条逻辑相同的消息。
Consumer Ack 丢失
业务提交 -> Ack 网络失败 -> Broker 重新投递
消费者主动重试
消息被发送到 retry Queue 后,原消息 Ack 失败,也会形成重复。
Broker 或客户端恢复
连接恢复、消费者重新注册和未确认消息重新投递,都可能让同一业务事件再次到达。
因此,所谓幂等不是“消息永远只投递一次”,而是:
同一个业务操作执行一次或多次,最终业务状态和执行一次相同。
9.2 使用唯一事件 ID
事件必须包含稳定且全局唯一的 ID:
public record OrderCreated(
String eventId,
String orderId,
long occurredAt,
BigDecimal amount
) {
}
重试和重发时必须保留原 eventId,不能每次重新生成,否则消费者无法判断它们是否属于同一个业务操作。
9.3 Inbox 表实现消费幂等
一种通用做法是建立消费记录表:
CREATE TABLE consumed_message (
consumer_name VARCHAR(100) NOT NULL,
event_id VARCHAR(100) NOT NULL,
consumed_at TIMESTAMP NOT NULL,
PRIMARY KEY (consumer_name, event_id)
);
处理流程:
开始数据库事务
-> INSERT consumer_name + event_id
-> 插入成功:继续执行业务
-> 唯一键冲突:说明已处理,跳过业务
-> 提交数据库事务
-> basic.ack
伪代码如下:
@Transactional
public void handle(OrderCreated event) {
int inserted = jdbc.update("""
INSERT INTO consumed_message
(consumer_name, event_id, consumed_at)
VALUES (?, ?, CURRENT_TIMESTAMP)
ON CONFLICT (consumer_name, event_id) DO NOTHING
""",
"order-consumer",
event.eventId()
);
if (inserted == 0) {
// 已经处理过
return;
}
jdbc.update("""
INSERT INTO order_projection(order_id, amount)
VALUES (?, ?)
ON CONFLICT (order_id) DO NOTHING
""",
event.orderId(),
event.amount()
);
}
不同数据库的冲突语法不同:
- PostgreSQL 使用
ON CONFLICT; - MySQL 可使用
INSERT IGNORE或ON DUPLICATE KEY UPDATE; - SQL Server 可使用唯一索引配合异常处理或其他写法。
关键不是具体 SQL,而是“唯一约束和业务写入在同一个事务中”。
如果先检查、后写入:
SELECT event_id 是否存在
业务写入
INSERT event_id
两个并发消费者可能同时查询到“不存在”,然后都执行业务。唯一约束必须成为最终并发仲裁机制。
9.4 非数据库操作的幂等
如果消费者调用外部支付接口、邮件服务或 HTTP API,数据库 Inbox 不能自动撤销外部副作用。
应优先使用:
- 外部接口支持的幂等键,例如
eventId; - 本地操作记录表;
- 可查询的状态机;
- 对账和补偿机制。
例如支付请求应把 eventId 作为下游幂等键:
POST /payments
Idempotency-Key: event-8b9...
如果下游不支持幂等,而操作本身又不可撤销,就无法仅靠 RabbitMQ Ack 保证 exactly-once。此时需要重新评估接口协议和业务补偿方案。
十、顺序:RabbitMQ 能保证哪一种顺序
10.1 Queue 内的局部顺序
对于一个普通 Queue,RabbitMQ 通常按照消息进入队列的顺序投递消息。但这不是“任意情况下的全局顺序保证”。
下列因素会改变可观察顺序:
- 一个 Queue 有多个消费者;
prefetch大于 1;- 消费者处理时间不同;
- 消息被 nack 后重新入队;
- 使用消息优先级;
- 消息经过重试 Queue 后重新返回;
- Producer 通过多个并发 Channel 发布。
例如:
Q: A, B
Consumer 1 收到 A,处理 10 秒
Consumer 2 收到 B,处理 1 秒
业务完成顺序是:
B -> A
即使投递顺序是:
A -> B
所以“Queue 中按序”不等于“业务完成按序”。
10.2 为什么 prefetch=1 仍不等于全局顺序
单消费者、prefetch=1 能减少并发未确认消息,使投递和处理更接近串行:
spring.rabbitmq.listener.simple.concurrency=1
spring.rabbitmq.listener.simple.max-concurrency=1
spring.rabbitmq.listener.simple.prefetch=1
但它仍不能解决:
- 失败消息重新入队后的顺序变化;
- 重试消息晚于后续消息返回;
- Producer 并发发布产生的实际到达顺序;
- 消费者进程崩溃后的重新投递;
- 多个 Queue 之间没有全局顺序。
因此它是降低并发的配置,不是分布式顺序协议。
10.3 按业务键分片
如果要求“同一订单内有序,不同订单可以并行”,应把顺序范围定义为业务聚合键,例如 orderId。
目标是:
同一个 orderId -> 固定进入同一个分片 Queue
不同 orderId -> 可以进入不同分片 Queue
可以采用:
partition = hash(orderId) mod N
Producer 根据分片结果选择 routing key 或 Queue。这样:
order-100 -> partition-2
order-101 -> partition-0
order-100 -> partition-2
同一个订单不会在多个分片之间漂移,多个分片可以并行消费。
但扩容时 N 改变会导致普通取模结果变化,使同一业务键可能进入不同分片。需要使用一致性哈希、固定分区数量,或维护稳定的分片映射。
10.4 用序列号检测乱序
仅靠 Queue 顺序不足以证明业务顺序。事件可以携带业务序列号:
public record OrderEvent(
String eventId,
String orderId,
long sequence,
String type
) {
}
消费者维护:
orderId = O1, lastSequence = 4
收到 sequence = 5 -> 应用
收到 sequence = 7 -> 暂存或拒绝,等待 6
收到 sequence = 3 -> 重复或过期,按幂等规则处理
这把“顺序”从传输层问题转化为业务状态机问题。若必须等待缺失事件,需要额外的暂存表、超时告警和人工恢复机制,否则一个丢失或永久失败的事件会阻塞整个聚合键。
10.5 重试会破坏简单顺序
考虑:
A: order.paid
B: order.shipped
如果 A 处理失败进入 10 秒重试 Queue,而 B 已经正常到达主 Queue,那么 B 可能先于 A 执行。
要保持同一订单的严格顺序,通常不能把 A 单独移出后让 B 继续执行。可选方案包括:
- 同一业务键串行消费,并在失败时暂停该键后续事件;
- 使用序列号和暂存机制;
- 将同一聚合键放入顺序队列,失败时阻塞后续;
- 让业务状态机拒绝非法状态转移并安排补偿。
“重试”和“严格顺序”存在天然张力。延迟重试越独立,越容易让后续消息越过失败消息。
十一、Outbox:解决数据库写入和消息发布之间的间隙
如果业务流程是:
数据库事务提交订单
发布 order.created
Producer 可能在两步之间崩溃:
订单已提交
消息尚未发布
此时 RabbitMQ Confirm 无法补救,因为消息甚至没有被发布。
Outbox 模式把业务数据和待发布事件写入同一个数据库事务:
CREATE TABLE outbox_event (
event_id VARCHAR(100) PRIMARY KEY,
event_type VARCHAR(100) NOT NULL,
aggregate_id VARCHAR(100) NOT NULL,
payload TEXT NOT NULL,
status VARCHAR(20) NOT NULL,
created_at TIMESTAMP NOT NULL
);
业务事务:
BEGIN
INSERT order
INSERT outbox_event(event_id, ...)
COMMIT
独立发布器扫描 status = 'NEW' 的事件:
读取 Outbox
-> 发布到 RabbitMQ
-> 等待 Publisher Confirm
-> Confirm ack 后标记 SENT
如果发布器在 Confirm 后、更新 SENT 前崩溃,事件会再次发布。这里仍然是至少一次发布,因此消费者仍需使用 eventId 幂等。
Outbox 解决的是:
数据库提交成功,但消息没有被发布
它没有自动解决:
RabbitMQ 发布成功,但消费者处理失败
后者仍然依赖 Consumer Ack、重试、死信和幂等。
十二、不要把 RabbitMQ 事务误认为端到端事务
RabbitMQ 提供事务发布和 Confirm 等机制,但它们只能覆盖 RabbitMQ 协议范围,不能把以下操作自动绑定为一个原子事务:
数据库提交
RabbitMQ 发布
外部 HTTP 调用
RabbitMQ Ack
例如:
数据库事务提交成功
RabbitMQ 发布失败
或者:
RabbitMQ 发布成功
数据库事务回滚
都可能发生。
通常的架构取舍是:
- 数据库写入 + 事件记录:使用 Outbox;
- RabbitMQ 发布:使用 Publisher Confirm;
- 消费业务 + 消费记录:使用 Inbox 和数据库事务;
- 消费成功后:发送 Consumer Ack;
- 外部副作用:使用下游幂等键和状态机。
这是一组组合协议,而不是 RabbitMQ 单独提供的 exactly-once 事务。
十三、一个可验证的 Spring Boot 处理骨架
依赖至少需要 Spring AMQP:
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-amqp</artifactId>
</dependency>
使用 Java 25 运行时即可,但 Spring Boot 和 Spring AMQP 的具体版本必须选择官方声明支持 Java 25 的版本组合。Java 25 是运行时和编译目标,不能据此推断某个 Spring 版本自动兼容。
消费者配置:
spring.rabbitmq.host=localhost
spring.rabbitmq.port=5672
spring.rabbitmq.username=app
spring.rabbitmq.password=app
spring.rabbitmq.listener.simple.acknowledge-mode=manual
spring.rabbitmq.listener.simple.prefetch=1
spring.rabbitmq.listener.simple.default-requeue-rejected=false
default-requeue-rejected=false 只能影响 Spring 容器对异常的默认处理,不能代替对异常分类。对于需要延迟重试的异常,应明确发布到 retry Queue 或配置 Spring Retry;对于永久失败,应明确进入死信路径。
一个更完整的消费决策可以写成:
try {
service.handle(event); // 数据库事务中完成 Inbox + 业务写入
channel.basicAck(tag, false);
} catch (TransientBusinessException e) {
publishToRetryAndConfirm(event);
channel.basicAck(tag, false);
} catch (PermanentBusinessException e) {
channel.basicReject(tag, false);
} catch (Exception e) {
// 不应无条件无限 requeue
channel.basicReject(tag, false);
}
其中 publishToRetryAndConfirm 必须等待重试消息的发布确认后才能确认原消息;否则可能出现原消息已删除、重试消息却未进入 Broker 的丢失窗口。
十四、失败表现和诊断方法
14.1 Producer 显示 Confirm ack,但 Queue 没有消息
优先检查:
- routing key 是否正确;
- Exchange 和 Queue 是否有 binding;
- 是否启用了
mandatory; - 是否收到 Return;
- 是否把消息发布到了错误的 vhost;
- Exchange 类型是否与预期一致。
Confirm ack 只说明 Exchange 接受发布,不说明一定路由成功。
14.2 消费者不断重复同一条消息
检查:
- 是否一直
nack(requeue=true); - 业务异常是否在 Ack 前抛出;
- 数据库事务是否实际提交;
- Ack 是否使用了错误 Channel;
- 消费者是否频繁重启;
prefetch是否过大;- 是否因为外部调用超时而重复执行。
查看 RabbitMQ 管理界面中的 messages_unacknowledged、消费者数和重新投递信息,并记录:
eventId
messageId
deliveryTag
redelivered
attempt
异常类型
redelivered 可以帮助识别 Broker 重新投递,但它不是可靠的业务重试次数,也不应替代持久化重试计数。
14.3 死信队列为空
检查:
- 原 Queue 是否真的设置了
x-dead-letter-exchange; - DLX Exchange 是否存在;
- DLX 的 routing key 是否正确;
- 目标 Queue 是否绑定;
- 消费者是否使用了
requeue=true,导致消息根本没有进入 DLX; - 是否因权限不足而转发失败;
- TTL 或长度限制是否真的触发。
14.4 顺序偶尔颠倒
检查:
- 是否有多个消费者;
prefetch是否大于 1;- 是否存在 retry Queue;
- 是否发生 nack 和重新入队;
- Producer 是否并发发布;
- 是否有多个 Queue 或多个分片;
- 业务是否观察的是“完成顺序”而非“投递顺序”。
如果业务确实要求顺序,应记录 aggregateId 和 sequence,只看日志时间通常无法准确判断消息因果关系。
十五、关键边界和取舍
Exchange
Exchange 解决路由,不解决持久化、不解决消费确认,也不保证同一业务键的全局顺序。
Publisher Confirm
Confirm 解决发布端与 Broker 之间的确认。它不确认路由到 Queue,也不确认消费者业务成功;不可路由消息还需要 Return。
Consumer Ack
Ack 解决 Broker 是否可以删除某次投递。Ack 必须晚于业务提交,但这会引入“业务已提交、Ack 未发送”的重复窗口。
重试
重试适合暂时失败。直接 requeue=true 没有退避,容易形成热循环;TTL + DLX 能提供固定延迟,但时间精度、队首阻塞和顺序都会受到影响。
死信
死信是离开正常处理路径的消息,不是自动可靠归档。DLX 仍然依赖 Exchange、binding、权限和目标 Queue;最终死信必须支持诊断和恢复。
顺序
RabbitMQ 更适合提供 Queue 内的局部顺序,而不是跨 Queue、跨消费者和跨重试链路的全局顺序。严格顺序需要业务分片、序列号和状态机共同实现。
幂等
在至少一次发布和至少一次消费模型下,幂等是业务正确性的基础。稳定事件 ID、唯一约束、Inbox、Outbox 和下游幂等键分别解决不同故障窗口,不能相互替代。
系列导航与关联阅读
- 系列入口:Java 完整学习路线:从 Java 25 语言与 JVM 到 Spring、微服务和生产交付
- 上一篇:Java Kafka:Producer、Consumer、分区、Offset、事务和再均衡
- 下一篇:Java 服务韧性:超时、重试、限流、熔断、隔离和降级
官方资料
本文依据 Java、Spring 与相关项目官方文档重新梳理;正文、示例与生产清单由 WR BLOG 编写。

评论
0 条讨论