订单状态同步、库存扣减通知这类场景对消息丢失零容忍,而RabbitMQ的默认配置并不提供这个保证:生产者发完不确认、Broker内存存储、消费者自动ACK,三个环节任何一个出问题消息就丢了。可靠性投递需要针对发送、存储、消费三个环节分别设计,缺一不可。
发送端可靠性:Publisher Confirm与Return回调
Spring AMQP开启确认机制:
spring:
rabbitmq:
publisher-confirm-type: correlated
publisher-returns: true
template:
mandatory: true
rabbitTemplate.setConfirmCallback((data, ack, cause) -> {
if (!ack) {
// correlationData携带业务ID,记录后走补偿重发
log.error("消息未确认, id={}, cause={}", data.getId(), cause);
}
});
rabbitTemplate.setReturnsCallback(returned ->
log.warn("路由失败, exchange={}, routingKey={}",
returned.getExchange(), returned.getRoutingKey()));
Confirm回答消息是否到达Broker,Return回答消息到达Broker后是否成功路由到队列。两者任一失败,把消息落到本地消息表,由定时任务补偿重发,这是发送端不丢消息的兜底手段。
存储端可靠性:持久化与仲裁队列
交换机、队列、消息三者都要持久化,缺一个重启后仍会丢数据:
@Bean
public Queue orderQueue() {
return QueueBuilder.durable("order.notify") // 队列持久化
.quorum() // 仲裁队列,多数派写入
.build();
}
// 消息持久化
rabbitTemplate.convertAndSend("order.exchange", "order.notify",
message, msg -> {
msg.getMessageProperties().setDeliveryMode(MessageDeliveryMode.PERSISTENT);
return msg;
});
生产环境用仲裁队列(Quorum Queue)替代经典的镜像队列,基于Raft协议多数派写入,节点宕机不丢消息。
消费端可靠性:手动ACK与死信重试
spring:
rabbitmq:
listener:
simple:
acknowledge-mode: manual
default-requeue-rejected: false
@RabbitListener(queues = "order.notify")
public void onMessage(Message message, Channel channel) throws IOException {
long tag = message.getMessageProperties().getDeliveryTag();
try {
process(message);
channel.basicAck(tag, false);
} catch (Exception e) {
long delivered = getXDeliverCount(message);
if (delivered >= 2) {
channel.basicNack(tag, false, false); // 转入死信队列人工处理
} else {
channel.basicNack(tag, false, true); // requeue重试
}
}
}
自动ACK意味着消息一到消费者内存就标记完成,处理过程崩溃即丢消息,可靠性场景必须改手动ACK。
幂等消费:重复投递无法避免,只能消化
补偿重发、requeue重试都会造成重复投递,消费端必须幂等。常用方案是业务唯一键加状态表:
CREATE TABLE consumed_msg (
msg_id VARCHAR(64) PRIMARY KEY,
consumed_at DATETIME NOT NULL
);
public void process(Message message) {
String msgId = message.getMessageProperties().getMessageId();
// 唯一键插入,插入成功才执行业务,冲突说明已处理过
int rows = jdbc.update(
"INSERT IGNORE INTO consumed_msg(msg_id, consumed_at) VALUES(?, NOW())", msgId);
if (rows == 0) {
return; // 重复消息,直接跳过
}
doBusiness(message);
}
没有库表条件的场景用Redis SETNX加过期时间,但要评估Redis故障窗口期的漏判风险;资金类强一致场景还是数据库唯一键可靠。
可靠性成本的边界
确认回调、本地消息表、幂等校验,每一层都增加延迟和复杂度。判断标准是消息的业务代价:丢一条导致资金或库存错误的走全量可靠方案;丢一条只影响体验的通知类消息,允许少量丢失换取吞吐。可靠性设计跟着业务代价走,不是所有队列都需要这套完整配置。
原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/xiao-xi-dui-lie-ke-kao-xing-tou-di-shi-zhan-rabbitmq-que/