消息队列可靠性投递实战:RabbitMQ确认机制、持久化与幂等消费

订单状态同步、库存扣减通知这类场景对消息丢失零容忍,而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/

(0)
小编小编
上一篇 21小时前
下一篇 21小时前

相关推荐