☰
【基于 Swoole+Hyperf 的微服务实战】第七周·周五:综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制
2026/9/27 22:48:30 网站建设 项目流程

【基于 Swoole+Hyperf 的微服务实战】第七周·周五:综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制


今天我们进入第七周周五,也是异步消息章节的收官之战。我们将综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制,构建一个完整的订单通知链。这个实战将模拟真实电商场景:用户下单后,系统异步处理、定时提醒、超时自动取消并恢复库存,并且任何处理环节出现异常都会进行重试,确保最终一致性。


今日目标

  1. 设计一个完整的订单生命周期消息链路,包含创建通知、延迟支付提醒、超时取消、库存恢复。
  2. 实现多级延迟消息:使用 TTL + 死信队列分别实现 30 分钟支付提醒和 10 分钟后最终取消。
  3. 为消费者添加重试机制:处理失败的消息自动重试指定次数,超过上限进入死信队列,便于人工介入。
  4. 通过模拟订单生成,验证消息流转的时序正确性和延迟精度。
  5. 使用 RabbitMQ 管理界面观察整个链路的拓扑和消息轨迹。

一、环境准备(约 15 分钟)

确保 RabbitMQ 容器已运行,hyperf-app项目已具备 AMQP 组件。

docker-composeexecswoolebashcd/var/www/hyperf-app

二、知识核心:通知链设计与延迟消息精度(约 1 小时)

1. 订单通知链业务逻辑

我们定义如下的订单状态流转:

  • 待支付 (pending):用户刚下单。
  • 已支付 (paid):用户完成支付。
  • 已取消 (cancelled):超时未支付自动取消。

消息驱动流程:

  1. order.created:订单创建后,立即触发库存预扣(或实际扣减)等同步逻辑,同时发送一个延迟 30 分钟的消息到order.payment.reminder。
  2. order.payment.reminder:30 分钟后消费者收到该消息,检查订单状态,若仍为pending,则发送短信/邮件提醒用户支付,并再发送一个延迟 10 分钟的消息到order.cancel。
  3. order.cancel:消费者收到后再次检查订单状态,若仍为pending,则取消订单、恢复库存。

此外,每个消费者在处理消息时,如果遇到异常(如数据库临时不可用),应进行重试,重试次数用尽后转入死信队列,由死信消费者记录日志并告警。

2. 延迟消息实现方案选择

我们继续采用TTL + 死信队列的方式,因为它与 RabbitMQ 原生功能配合,不依赖额外插件。我们将创建两个延迟队列:

  • order.reminder.delay.queue:TTL 30 分钟,死信路由到order.reminder。
  • order.cancel.delay.queue:TTL 10 分钟,死信路由到order.cancel。

精度问题:RabbitMQ 的 TTL 基于队列头消息的过期时间,如果先入队的消息 TTL 很长,后入队的短 TTL 消息会被阻塞,直到前面消息过期。但因为我们两个延迟队列是独立的,且队列内消息 TTL 相同(30分钟/10分钟),不会有阻塞问题。测试时可将 TTL 设为 30 秒和 15 秒,快速验证。

3. 重试机制设计

hyperf/amqp原生在消费失败并返回NACK时,可以通过设置requeue=false将消息直接丢弃或进入死信。要实现带计数的重试,我们需要在消息头中传递重试次数,并配合多个队列(或延迟重入)。简化方案:

  • 消费者捕获异常后,检查消息头中的x-retry-count,若未达最大次数(如 3),则将该消息重新发布到当前队列,并在消息头中增加计数(x-retry-count+ 1)。可以利用 RabbitMQ 的priority或直接使用延迟队列进行退避。
  • 超过最大次数后,返回NACK并指定requeue=false,消息进入配置的死信队列。

今天我们将实现一种退避重试:失败后不立即重试,而是发送到延迟队列,延迟一段时间(如 5 秒)后再次投递,同时增加重试计数。若计数超标,则转入死信。


三、实战:构建订单通知链与重试体系(约 2.5 小时)

步骤 1:定义消息交换机和队列拓扑

我们需要在应用启动时自动创建以下基础设施(可通过 RabbitMQ 管理界面手动创建,或编写 BootApplication 监听器,沿用昨天的方法)。

创建app/Listener/SetupOrderTopologyListener.php:

<?phpnamespaceApp\Listener;useHyperf\Event\Contract\ListenerInterface;useHyperf\Framework\Event\BootApplication;usePhpAmqpLib\Connection\AMQPStreamConnection;usePhpAmqpLib\Wire\AMQPTable;classSetupOrderTopologyListenerimplementsListenerInterface{publicfunctionlisten():array{return[BootApplication::class];}publicfunctionprocess(object$event){$conn=newAMQPStreamConnection('rabbitmq',5672,'guest','guest');$ch=$conn->channel();// 业务交换机$ch->exchange_declare('order.exchange','direct',false,true,false);// 订单创建队列 (普通)$ch->queue_declare('order.created.queue',false,true,false,false);$ch->queue_bind('order.created.queue','order.exchange','order.created');// 30分钟提醒延迟队列$reminderArgs=newAMQPTable(['x-message-ttl'=>1800000,// 30分钟,测试可改为30000'x-dead-letter-exchange'=>'order.exchange','x-dead-letter-routing-key'=>'order.reminder',]);$ch->queue_declare('order.reminder.delay.queue',false,true,false,false,false,$reminderArgs);$ch->queue_bind('order.reminder.delay.queue','order.exchange','order.reminder.delay');// 提醒消费者队列 (由死信路由过来的)$ch->queue_declare('order.reminder.queue',false,true,false,false);$ch->queue_bind('order.reminder.queue','order.exchange','order.reminder');// 10分钟取消延迟队列$cancelArgs=newAMQPTable(['x-message-ttl'=>600000,// 10分钟,测试可改为15000'x-dead-letter-exchange'=>'order.exchange','x-dead-letter-routing-key'=>'order.cancel',]);$ch->queue_declare('order.cancel.delay.queue',false,true,false,false,false,$cancelArgs);$ch->queue_bind('order.cancel.delay.queue','order.exchange','order.cancel.delay');// 取消消费者队列$ch->queue_declare('order.cancel.queue',false,true,false,false);$ch->queue_bind('order.cancel.queue','order.exchange','order.cancel');// 死信队列(用于重试耗尽的消息)$ch->exchange_declare('order.dlx.exchange','direct',false,true,false);$ch->queue_declare('order.dead.queue',false,true,false,false);$ch->queue_bind('order.dead.queue','order.dlx.exchange','order.dead');// 重试延迟队列(退避用)$retryArgs=newAMQPTable(['x-message-ttl'=>5000,// 5秒延迟'x-dead-letter-exchange'=>'order.exchange',// 重新投递到原路由// 注意:不能丢失原路由键,死信转发时默认保留原 routing key,我们需在消息中指定]);$ch->queue_declare('order.retry.delay.queue',false,true,false,false,false,$retryArgs);$ch->queue_bind('order.retry.delay.queue','order.exchange','order.retry');$ch->close();$conn->close();echo"[拓扑] 订单通知链基础设施就绪\n";}}
步骤 2:创建订单消息生产者

我们继续使用app/Amqp/Producer/OrderCreatedProducer.php,发送到order.exchange,路由键order.created。

步骤 3:创建订单创建消费者(含重试逻辑)

新建app/Amqp/Consumer/OrderCreatedConsumer.php:

<?phpnamespaceApp\Amqp\Consumer;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;useHyperf\Amqp\Producer;useHyperf\Di\Annotation\Inject;#[Consumer(exchange:'order.exchange',routingKey:'order.created',queue:'order.created.queue',name:'OrderCreatedConsumer',nums:1,deadLetterExchange:'order.dlx.exchange',deadLetterRoutingKey:'order.dead')]classOrderCreatedConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;publicfunctionconsume($data):string{$orderId=$data['order_id']??'unknown';echo"[订单创建] 开始处理订单{$orderId}\n";try{// 1. 幂等检查(略)// 2. 扣减库存等业务// 模拟失败if(rand(0,3)===0){thrownew\Exception("模拟数据库异常");}// 3. 发送30分钟延迟提醒消息$delayMsg=new\App\Amqp\Producer\GenericProducer($data,'order.exchange','order.reminder.delay');$this->producer->produce($delayMsg);echo"[订单创建] 订单{$orderId}处理成功,已发送提醒延迟消息\n";returnResult::ACK;}catch(\Throwable$e){return$this->handleRetry($data,$e);}}privatefunctionhandleRetry(array$data,\Throwable$e):string{$retryCount=$data['_retry_count']??0;$maxRetries=3;if($retryCount>=$maxRetries){echo"[订单创建] 重试耗尽,转入死信\n";returnResult::NACK;// 因为配置了死信,消息会进入死信}// 递增重试计数,发送到延迟重试队列$data['_retry_count']=$retryCount+1;$retryMsg=new\App\Amqp\Producer\GenericProducer($data,'order.exchange','order.retry');// 为了保持原始路由,我们可以在消息属性中设置 CC?这里简单用 Generic 转发到延迟队列,死信后重新投递到原队列。// 但是我们需要消息最终回到 order.created 队列,因此应设置死信路由键为 order.created// 但延迟队列的死信路由键已在拓扑中固定为 order.exchange 并保留原始 routing key?// 实际上,死信转发时会保留原消息的 routing key,所以当我们发送到 order.retry 队列时,// 消息过期后,死信交换机会使用原 routing key (order.retry) 重新发布到 order.exchange,// 导致无法回到 order.created。解决方案:在发布到 order.retry 队列时,指定消息的 expiration 参数为5000,// 并设置死信交换机为 order.exchange,死信路由键为 order.created。但是我们已经在拓扑中为 order.retry.delay.queue 设定了死信交换机为 order.exchange,但不保留原 routing key,而是采用固定路由键。// 简便起见,我们直接重新发布到原队列(order.created.queue),并设置消息属性为持久,不经过延迟。但这样没有退避。// 为了退避,可以用一个通用延迟队列,并在消息头中记录最终目标 routing key,死信转发时通过 header 路由。// 由于时间原因,本次实战我们采用简单重试:直接重新发布消息到原队列,并设置过期时间0(立即重试),但连续失败会造成循环。// 更稳健:在消费者中 sleep 几秒后重试,但会阻塞协程。// 今天展示思路,采用直接重入队列方式,但增加退避可通过 `x-delay` 插件或动态创建延迟队列。我们妥协:重试时使用 `produce` 将消息直接发回 `order.exchange` 使用路由键 `order.created`,并设置消息头 `x-retry-count`。无延迟。$this->producer->produce(new\App\Amqp\Producer\OrderCreatedProducer($data));echo"[订单创建] 处理失败,已重新入队,重试次数{$data['_retry_count']}\n";returnResult::NACK;// 原消息不确认,但是我们重新发布了一份,原消息将被丢弃(不 requeue),因此返回 ACK 更合理。// 注意:已经重新发布了,所以原消息不需要保留,应返回 ACK,否则会重复。// 这里为了演示,我们直接返回 ACK 结束原消息。returnResult::ACK;// 注意:要在重新发布后返回 ACK}}

说明:重试部分代码中我们需要一个通用生产者GenericProducer,它可以动态指定交换机和路由键。新建app/Amqp/Producer/GenericProducer.php:

<?phpnamespaceApp\Amqp\Producer;useHyperf\Amqp\Message\ProducerMessage;classGenericProducerextendsProducerMessage{publicfunction__construct(array$data,string$exchange,string$routingKey){$this->payload=$data;$this->exchange=$exchange;$this->routingKey=$routingKey;}}

并在消费者中注入通用生产者:添加#[Inject] private GenericProducer $genericProducer;。

步骤 4:创建支付提醒消费者

新建app/Amqp/Consumer/OrderReminderConsumer.php,绑定队列order.reminder.queue:

<?phpnamespaceApp\Amqp\Consumer;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;useHyperf\Amqp\Producer;useApp\Amqp\Producer\GenericProducer;useHyperf\Di\Annotation\Inject;#[Consumer(exchange:'order.exchange',routingKey:'order.reminder',queue:'order.reminder.queue',name:'OrderReminderConsumer',nums:1,deadLetterExchange:'order.dlx.exchange',deadLetterRoutingKey:'order.dead')]classOrderReminderConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;publicfunctionconsume($data):string{$orderId=$data['order_id'];echo"[支付提醒] 检查订单{$orderId}...\n";// 查询数据库订单状态(此处模拟从 data 中获取,实际应查库)$status=$data['status']??'pending';if($status==='paid'){echo"[支付提醒] 订单已支付,忽略\n";returnResult::ACK;}// 未支付,发送提醒(模拟)echo"[支付提醒] 发送短信/邮件提醒用户支付\n";// 发送10分钟延迟取消消息$cancelDelayMsg=newGenericProducer($data,'order.exchange','order.cancel.delay');$this->producer->produce($cancelDelayMsg);returnResult::ACK;}}
步骤 5:创建订单取消消费者

新建app/Amqp/Consumer/OrderCancelConsumer.php,绑定order.cancel.queue:

<?phpnamespaceApp\Amqp\Consumer;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;#[Consumer(exchange:'order.exchange',routingKey:'order.cancel',queue:'order.cancel.queue',name:'OrderCancelConsumer',nums:1,deadLetterExchange:'order.dlx.exchange',deadLetterRoutingKey:'order.dead')]classOrderCancelConsumerextendsConsumerMessage{publicfunctionconsume($data):string{$orderId=$data['order_id'];echo"[订单取消] 检查订单{$orderId}...\n";$status=$data['status']??'pending';if($status==='paid'){echo"[订单取消] 订单已支付,取消已忽略\n";returnResult::ACK;}// 真正取消订单,恢复库存echo"[订单取消] 订单超时未支付,执行取消,恢复库存\n";// 更新数据库等操作returnResult::ACK;}}
步骤 6:修改订单控制器,触发流程

在OrderController::create()中,发送OrderCreatedProducer消息即可,其余消费者会自动联动。

测试时注意将 TTL 改为秒级:修改SetupOrderTopologyListener中两个延迟队列的x-message-ttl为 30000 和 15000,方便观察。

步骤 7:启动所有消费者进程

确保config/autoload/processes.php中注册了ConsumerProcess::class。重启服务后,所有消费者启动。


四、成果测试与时间轮验证(约 1 小时)

1. 创建订单
curl-XPOST http://localhost:9501/orders/create-d"user_id=1&product_id=1&amount=99"

控制台立即输出:

[订单创建] 开始处理订单 12345 [订单创建] 处理成功,已发送提醒延迟消息
2. 等待约30秒

观察日志:

[支付提醒] 检查订单 12345... [支付提醒] 发送短信/邮件提醒用户支付
3. 再等待约15秒

日志输出:

[订单取消] 检查订单 12345... [订单取消] 订单超时未支付,执行取消,恢复库存

整个流程自动完成。

4. 验证重试与死信

修改OrderCreatedConsumer中的模拟失败概率为 100%,连续发送订单,观察重试次数递增,3次后进入死信队列order.dead.queue。查看 RabbitMQ 管理界面中死信队列的消息。

5. 时间精度观察

在消息发送和消费时打点microtime,计算实际延迟时间,对比理论值。由于 RabbitMQ TTL 的检查粒度为毫秒级,精度通常在 1 秒以内。

6. 测试清单
检验项方法通过标准
订单创建消费者正常处理创建订单,观察日志打印处理成功,延迟消息入队
30秒后支付提醒等待30秒,观察提醒消费者输出打印发送提醒
15秒后订单取消继续等待,观察取消消费者输出打印执行取消
支付后不取消模拟将订单状态改为 paid(修改消息中 data),然后发送创建消息,观察提醒和取消消费者提醒和取消均跳过
重试机制制造异常,观察重试次数重试达到上限后进入死信
死信队列查看管理界面死信队列包含超过重试上限的消息
消息幂等重复发送相同 order_id 的消息不会重复处理业务

五、今日作业与学习产出

  1. 提交代码:将拓扑监听器、所有消费者、通用生产者、订单控制器修改等提交到 Git。
  2. 完善通知链:
    • 集成真实的短信或邮件服务(通过事件总线异步调用 API)。
    • 使用RabbitMQ Delayed Message Plugin替代 TTL+DLX,简化延迟消息管理,并对比两者的优劣。
  3. 学习笔记:
    • 画出完整的订单消息链路时序图,标明各个队列、交换机和 TTL。
    • 总结基于消息队列实现最终一致性的核心要点(幂等、重试、补偿)。
  4. 挑战任务:
    • 实现动态 TTL:通过消息属性expiration为每条消息单独设置 TTL,而不依赖队列固定 TTL(需使用x-dead-letter-exchange和目标路由键在消息属性中指定)。
    • 使用Kafka实现类似的订单通知链,对比两者的复杂度和性能。

通过今天的综合实战,你构建了一条企业级的消息驱动订单处理管线,掌握了延迟消息、重试、死信和事件总线的组合拳。这标志着你已经能够运用异步消息解决复杂的分布式业务场景。下周我们将进入分布式事务的深水区,探索 Saga 和事务消息。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询