【基于 Swoole+Hyperf 的微服务实战】消息可靠性:手动确认、死信队列与幂等消费
今天我们进入的主题是 消息可靠性:手动确认、死信队列与幂等消费。昨天我们已经搭建了基本的 RabbitMQ 生产者和消费者,但消息可能因为消费异常而丢失,或者因为重复消费导致数据不一致。今天我们将为订单消息加上可靠投递和消费保障:通过手动确认、死信队列、幂等机制,确保消息被正确处理且不丢失、不重复。

今日目标
一、环境准备(约 15 分钟)
继续使用昨天的 RabbitMQ 容器,确保 hyperf-app 项目已安装 hyperf/amqp,并且已有 OrderCreatedProducer 和 OrderCreatedConsumer 基础代码。
进入 PHP 容器:
docker-compose exec swoole bash
cd /var/www/hyperf-app
二、知识核心:消息确认、死信与幂等(约 1 小时)
1. 消息确认(ACK/NACK)
在 RabbitMQ 中,消费者处理消息后可以:
- ACK:确认消息处理成功,消息从队列中删除。
- NACK:拒绝消息,可选是否重新入队。如果 requeue=true,消息会回到队列头部,可能被再次投递(可能造成死循环)。若 requeue=false,消息会被丢弃或转入死信队列。
- 超时:若消费者在处理过程中崩溃或连接断开,未确认的消息会自动重新入队,实现可靠的重新投递。
hyperf/amqp 默认要求消费者必须返回 Result::ACK 或 Result::NACK,否则消息会一直挂起。
2. 死信队列(DLX)
死信队列(Dead Letter Exchange)是 RabbitMQ 的一个特性:当消息满足以下条件之一时,会被转发到指定的死信交换机,然后进入绑定的死信队列:
- 消费者主动拒绝且 requeue=false。
- 消息达到最大重试次数(通常通过添加重试计数头实现,或利用队列的 x-message-ttl 和 x-dead-letter-exchange)。
- 消息过期(TTL 超时)。
我们可以在创建队列时定义参数 x-dead-letter-exchange 和 x-dead-letter-routing-key,RabbitMQ 就会自动将死信发送过去。
3. 幂等消费
在分布式环境中,网络抖动、消费者故障都可能导致同一条消息被投递多次。幂等消费意味着无论消费多少次,业务结果与消费一次相同。常用方案:
- 唯一索引:在数据库表中使用 order_id 唯一约束,重复插入会失败。
- 状态检查:处理前检查订单状态,已处理则跳过。
- Redis SETNX:处理前在 Redis 中尝试设置 msg_id 键,已存在则跳过。
三、实战:改造消费者为可靠且幂等(约 2.5 小时)
步骤 1:设计死信队列结构
我们需要创建以下交换机和队列:
- 正常交换机:order.exchange (direct)
- 正常队列:order.queue,设置死信交换机为 order.dlx.exchange,死信路由键为 order.dead
- 死信交换机:order.dlx.exchange (direct)
- 死信队列:order.dead.queue
hyperf/amqp 可以通过注解的 arguments 参数在声明队列时注入死信设置。
步骤 2:修改消费者注解,增加死信配置
编辑 app/Amqp/Consumer/OrderCreatedConsumer.php:
<?php
namespace App\\Amqp\\Consumer;
use Hyperf\\Amqp\\Annotation\\Consumer;
use Hyperf\\Amqp\\Message\\ConsumerMessage;
use Hyperf\\Amqp\\Result;
#[Consumer(
exchange: 'order.exchange',
routingKey: 'order.created',
queue: 'order.queue',
name: 'OrderCreatedConsumer',
nums: 1,
deadLetterExchange: 'order.dlx.exchange',
deadLetterRoutingKey: 'order.dead'
)]
class OrderCreatedConsumer extends ConsumerMessage
{
public function consume($data): string
{
$orderId = $data['order_id'] ?? 'unknown';
echo "[x] 处理订单: {$orderId}\\n";
// 模拟幂等性检查:Redis 记录已处理订单
// 此处用静态变量模拟,实际应使用 Redis SETNX
static $processed = [];
if (isset($processed[$orderId])) {
echo "[x] 订单 {$orderId} 已处理,跳过\\n";
return Result::ACK;
}
// 模拟随机失败(50% 几率)
if (rand(0, 1) === 0) {
echo "[x] 订单 {$orderId} 处理失败,拒绝消息且不重试\\n";
// 返回 NACK 且 requeue=false,消息将进入死信
return Result::NACK;
}
// 处理成功,标记已处理
$processed[$orderId] = true;
// 实际业务:扣减库存…
echo "[x] 订单 {$orderId} 处理成功\\n";
return Result::ACK;
}
}
解释:
- 在 #[Consumer] 中增加了 deadLetterExchange 和 deadLetterRoutingKey 参数,Hyperf 会在声明队列时自动设置死信规则。
- 模拟随机失败,当处理失败时返回 Result::NACK(Hyperf 自动设置 requeue=false),消息将变为死信。
- 幂等检查使用静态变量演示,实际应替换为持久化存储。
步骤 3:创建死信消费者(用于监控或人工处理)
新建 app/Amqp/Consumer/OrderDeadConsumer.php:
<?php
namespace App\\Amqp\\Consumer;
use Hyperf\\Amqp\\Annotation\\Consumer;
use Hyperf\\Amqp\\Message\\ConsumerMessage;
use Hyperf\\Amqp\\Result;
#[Consumer(
exchange: 'order.dlx.exchange',
routingKey: 'order.dead',
queue: 'order.dead.queue',
name: 'OrderDeadConsumer',
nums: 1
)]
class OrderDeadConsumer extends ConsumerMessage
{
public function consume($data): string
{
$orderId = $data['order_id'] ?? 'unknown';
echo "[死信] 收到失败订单: {$orderId},记录告警并通知运维\\n";
// 记录到数据库或发送钉钉通知
return Result::ACK;
}
}
这样任何进入死信的订单都会被该消费者记录,而不至于默默丢失。
步骤 4:配置自动创建死信队列和交换机
为了让消费者启动时自动创建所需的交换机和队列,我们需要在消费者注解中声明足够的参数。Hyperf 的 @Consumer 注解会自动根据名称创建交换机和队列,但死信交换机和队列需要手动声明,或者通过 RabbitMQ 管理界面预先创建。更优雅的方式是编写一个 BootApplication 监听器,在应用启动时通过 AMQP 管理接口创建拓扑。
但为了简化,我们可以在消费者注解中使用 arguments 直接传递死信配置,并且手动创建死信队列。然而 Hyperf 的消费者若没有对应的生产者和声明,不会主动创建死信交换机和队列。我们可以用 RabbitMQ 管理界面手动创建,或者依赖消费者首次运行时的声明。
其实,当消费者带有 deadLetterExchange 参数时,Hyperf 会在声明 order.queue 时设置 x-dead-letter-exchange,但死信交换机本身需要被创建。我们可以通过 #[Producer] 注解定义一个死信生产者(即使不用它发送消息),它会在启动时创建交换机和队列。
新建 app/Amqp/Producer/OrderDeadProducer.php(仅用于声明拓扑):
<?php
namespace App\\Amqp\\Producer;
use Hyperf\\Amqp\\Annotation\\Producer;
use Hyperf\\Amqp\\Message\\ProducerMessage;
#[Producer(exchange: 'order.dlx.exchange', routingKey: 'order.dead')]
class OrderDeadProducer extends ProducerMessage
{
public function __construct(array $data = [])
{
$this->payload = $data;
}
}
同时,为了让队列 order.dead.queue 被创建,我们可以在死信消费者中确保队列声明存在(通过消费注解的 queue 参数,已包含)。启动消费者进程时,这些交换机和队列会被自动创建(如果 RabbitMQ 有创建权限)。重启服务后检查管理界面。
步骤 5:增强幂等性实现(使用 Redis)
在 OrderCreatedConsumer 中注入 Redis,使用 SETNX 保证唯一性。
use Hyperf\\Redis\\Redis;
#[Inject]
private Redis $redis;
public function consume($data): string
{
$orderId = $data['order_id'] ?? 'unknown';
$idempotentKey = 'order:processed:' . $orderId;
// 尝试设置幂等键,仅当不存在时设置成功
if (!$this->redis->set($idempotentKey, 1, ['NX', 'EX' => 3600])) {
echo "[x] 订单 {$orderId} 已处理,幂等跳过\\n";
return Result::ACK;
}
// … 业务处理
return Result::ACK;
}
这样就实现了基于 Redis 的幂等消费,即使消息重复投递,也只有第一次会被处理。
步骤 6:模拟异常和重试,验证死信
重启 hyperf-app 使配置生效。
使用 curl 创建多个订单:
for i in {1..10}; do
curl -X POST http://localhost:9501/orders/create -d "user_id=1&product_id=1&amount=99"
echo
done
观察控制台输出:
- 部分订单处理成功。
- 约 50% 的订单因随机失败返回 NACK,这些消息会进入死信队列。
- 死信消费者会打印 [死信] 收到失败订单: …。
登录 RabbitMQ 管理界面:
- 查看 order.queue,消息会被消费或进入死信。
- order.dead.queue 中会有死信消息。
- 死信消息被 OrderDeadConsumer 消费后消失。
步骤 7:实现延迟消息(订单超时取消)思路
RabbitMQ 没有内置延迟队列,但可以通过 TTL + 死信组合实现:
- 定义一个 order.delay.queue,设置 x-message-ttl: 30000(30秒)和 x-dead-letter-exchange: order.exchange、x-dead-letter-routing-key: order.cancel。
- 当需要延迟时,发送消息到该队列(而不是直接到 order.created)。消息在队列中存活 30 秒后过期,成为死信,重新被投递到 order.exchange 并路由到 order.cancel,由取消消费者处理。
今天可以先创建这样的拓扑结构,并在订单创建时同时发送一个延迟消息到 order.delay.queue,实现 30 秒后检查订单状态并自动取消。这一部分可以作为挑战任务。
四、成果测试与验证(约 1 小时)
测试清单
| 手动确认 ACK | 消费者成功处理后返回 Result::ACK,队列消息移除 | 控制台打印成功,管理界面队列消息数减少 |
| 手动拒绝 NACK | 消费者处理失败返回 Result::NACK(requeue=false) | 消息进入死信队列,不在原队列 |
| 死信队列 | 查看 RabbitMQ 管理界面 order.dead.queue | 存在死信消息,且被死信消费者处理 |
| 幂等性 (Redis) | 发送相同 order_id 的消息两次,观察日志 | 第二次显示“已处理,跳过” |
| 消费者崩溃消息不丢 | 处理过程中杀掉消费者进程(Ctrl+C),消息重新入队 | 消息不会被确认,重新投递 |
| 死信消费者通知 | 死信消费者可写入日志或发送告警 | 死信被记录,不再丢失 |
常见问题
- 死信队列未创建:检查消费者注解中的 deadLetterExchange 和 deadLetterRoutingKey 是否拼写正确;确认 RabbitMQ 中有对应死信交换机(可通过管理界面或生产者注解创建)。
- NACK 后消息直接消失:如果没有配置死信,且 requeue=false,消息会被丢弃。务必配置死信。
- 幂等键未生效:确认 Redis 连接正常,键的过期时间合理。
五、今日作业与学习产出
- 创建 order.delay.queue 和相关交换绑定,实现订单创建 30 秒后自动检查支付状态,未支付则取消。
- 在消费者中实现取消逻辑(更新订单状态为“已取消”,恢复库存)。
- 画出消息的生命周期流程图:发布 → 正常消费 → 重试 → 死信 → 告警。
- 解释为什么幂等性是异步消息系统的“必备”特性,并列出三种实现方式。
- 实现消息重试:消费者不直接 NACK,而是捕获异常后将消息重新发布到带 TTL 的重试队列,达到重试次数上限后再进入死信。这种方式更灵活。
- 编写一个消息补偿控制器,通过管理接口手动重新投递死信队列中的消息,实现人工修复。
通过今天的学习,你的消息系统具备了企业级的可靠性:消息不丢失、不重复、异常可追溯。明天我们将引入 Kafka,体验高吞吐场景下的另一种消息引擎,并与现有 RabbitMQ 体系形成互补。
网硕互联帮助中心

评论前必须登录!
注册