如何在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(); } } }
关键细节说明
- 手动确认规则:无论业务成功还是转发重试,都必须调用
$msg->ack(),否则消息会一直滞留在主队列。 - 重试次数存储:通过消息的
application_headers记录重试次数,类型指定为I(整数)符合AMQP协议要求。 - 延迟生效逻辑:重试队列的
x-message-ttl设置延迟时长,到期后消息自动通过死信交换机回到主队列,触发重新消费。 - 死信队列用途:超过5次重试的消息进入此队列,可后续人工排查问题原因。
内容的提问来源于stack exchange,提问作者guilherme.souza
相关产品推荐
相关产品推荐

