You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在RabbitMQ中手动确认消息并配置递增延迟重试策略?

手动确认RabbitMQ消息+阶梯式重试策略实现(Symfony 6.1)

核心实现思路

要实现手动确认消息和阶梯间隔重试,不能依赖RabbitMQ默认的立即重试逻辑,需采用「死信交换机(DLX)+ 延迟队列(TTL)」方案:

  • 主队列绑定死信交换机,失败消息先转发到对应延迟重试队列
  • 每个重试队列设置递增TTL(5/10/15/20分钟),到期后消息自动回到主队列重新消费
  • 记录重试次数,超过5次后将消息转入最终死信队列,终止重试流程

第一步:修改RabbitMQ配置(old_sound_rabbit_mq.yaml)

新增死信交换机、4个阶梯延迟重试队列、1个最终死信队列,同时配置主队列的死信规则:

old_sound_rabbit_mq:
    connections:
        default:
            host: '%rabbitmqHost%'
            port: '%rabbitmqPort%'
            user: '%rabbitmqUser%'
            password: '%rabbitmqPassword%'
            vhost: '%rabbitmqVhost%'
    producers:
        # 新增生产者,用于转发失败消息到重试队列
        retry_producer:
            connection: default
            exchange_options: { name: 'upload_file_retry_exchange', type: direct, durable: true }
    consumers:
        upload_file:
            connection: default
            exchange_options: { name: 'upload_file_exchange', type: direct, durable: true, auto_delete: false }
            queue_options:
                name: 'upload_file_queue'
                durable: true
                auto_delete: false
                arguments:
                    'x-max-priority': [ 'I', 20 ]
                    # 主队列绑定死信交换机,失败消息先进入此交换机
                    'x-dead-letter-exchange': 'upload_file_retry_exchange'
                    'x-dead-letter-routing-key': 'retry_1'
            callback: App\Consumer\UploadFileConsumer
            qos_options: { prefetch_size: 0, prefetch_count: 1, global: false }
    queues:
        # 第2次重试队列(延迟5分钟)
        upload_file_retry_1:
            connection: default
            name: 'upload_file_retry_1'
            durable: true
            arguments:
                'x-message-ttl': 300000 # 5分钟=300000毫秒
                'x-dead-letter-exchange': 'upload_file_exchange'
                'x-dead-letter-routing-key': 'upload_file_queue'
                'x-max-priority': [ 'I', 20 ]
            bindings:
                - { exchange: 'upload_file_retry_exchange', routing_key: 'retry_1' }
        # 第3次重试队列(延迟10分钟)
        upload_file_retry_2:
            connection: default
            name: 'upload_file_retry_2'
            durable: true
            arguments:
                'x-message-ttl': 600000 # 10分钟=600000毫秒
                'x-dead-letter-exchange': 'upload_file_exchange'
                'x-dead-letter-routing-key': 'upload_file_queue'
                'x-max-priority': [ 'I', 20 ]
            bindings:
                - { exchange: 'upload_file_retry_exchange', routing_key: 'retry_2' }
        # 第4次重试队列(延迟15分钟)
        upload_file_retry_3:
            connection: default
            name: 'upload_file_retry_3'
            durable: true
            arguments:
                'x-message-ttl': 900000 # 15分钟=900000毫秒
                'x-dead-letter-exchange': 'upload_file_exchange'
                'x-dead-letter-routing-key': 'upload_file_queue'
                'x-max-priority': [ 'I', 20 ]
            bindings:
                - { exchange: 'upload_file_retry_exchange', routing_key: 'retry_3' }
        # 第5次重试队列(延迟20分钟)
        upload_file_retry_4:
            connection: default
            name: 'upload_file_retry_4'
            durable: true
            arguments:
                'x-message-ttl': 1200000 # 20分钟=1200000毫秒
                'x-dead-letter-exchange': 'upload_file_exchange'
                'x-dead-letter-routing-key': 'upload_file_queue'
                'x-max-priority': [ 'I', 20 ]
            bindings:
                - { exchange: 'upload_file_retry_exchange', routing_key: 'retry_4' }
        # 最终死信队列(超过5次重试后进入,不再重试)
        upload_file_dead_letter:
            connection: default
            name: 'upload_file_dead_letter'
            durable: true
            bindings:
                - { exchange: 'upload_file_retry_exchange', routing_key: 'dead_letter' }

第二步:修改Consumer代码实现手动确认与重试逻辑

在UploadFileConsumer中完成:手动确认成功消息、捕获异常后判断重试次数、转发消息到对应队列或死信队列:

<?php

declare(strict_types=1);

namespace App\Consumer;

use OldSound\RabbitMqBundle\RabbitMq\ConsumerInterface;
use OldSound\RabbitMqBundle\RabbitMq\ProducerInterface;
use PhpAmqpLib\Message\AMQPMessage;
use Symfony\Component\DependencyInjection\Attribute\Autowire;

class UploadFileConsumer implements ConsumerInterface
{
    private const MAX_RETRIES = 5;

    public function __construct(
        #[Autowire(service: 'old_sound_rabbit_mq.retry_producer')]
        private ProducerInterface $retryProducer
    ) {
    }

    public function execute(AMQPMessage $msg): void
    {
        // 读取当前重试次数,首次消费默认0
        $retryCount = (int)($msg->get('application_headers')->get('retry_count') ?? 0);

        try {
            // 业务逻辑:处理上传文件
            $payload = json_decode($msg->getBody(), true);
            // 替换为你的实际业务代码...

            // 处理成功,手动确认消息,从主队列移除
            $msg->ack();
        } catch (\Exception $e) {
            $retryCount++;

            if ($retryCount > self::MAX_RETRIES) {
                // 超过最大重试次数,转入最终死信队列
                $this->retryProducer->publish(
                    $msg->getBody(),
                    'dead_letter',
                    [],
                    ['application_headers' => ['retry_count' => ['I', $retryCount]]]
                );
                $msg->ack();
                return;
            }

            // 根据重试次数选择对应路由键,转发到指定延迟队列
            $routingKey = sprintf('retry_%d', $retryCount);
            $this->retryProducer->publish(
                $msg->getBody(),
                $routingKey,
                [],
                ['application_headers' => ['retry_count' => ['I', $retryCount]]]
            );

            // 确认原消息,从主队列移除
            $msg->ack();
        }
    }
}

关键细节说明

  1. 手动确认规则:无论业务成功还是转发重试,都必须调用$msg->ack(),否则消息会一直滞留在主队列。
  2. 重试次数存储:通过消息的application_headers记录重试次数,类型指定为I(整数)符合AMQP协议要求。
  3. 延迟生效逻辑:重试队列的x-message-ttl设置延迟时长,到期后消息自动通过死信交换机回到主队列,触发重新消费。
  4. 死信队列用途:超过5次重试的消息进入此队列,可后续人工排查问题原因。

内容的提问来源于stack exchange,提问作者guilherme.souza

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.16 07:50:25