首页/新闻资讯/正文详情

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

发布时间:2026/9/27 22:48:32 来源:云帆数科 栏目:资讯中心
【基于 Swoole+Hyperf 的微服务实战】第七周·周五:综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制
【基于 SwooleHyperf 的微服务实战】第七周·周五综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制今天我们进入第七周周五也是异步消息章节的收官之战。我们将综合运用 RabbitMQ 生产者、消费者、死信队列、TTL 延迟消息、幂等消费和事件机制构建一个完整的订单通知链。这个实战将模拟真实电商场景用户下单后系统异步处理、定时提醒、超时自动取消并恢复库存并且任何处理环节出现异常都会进行重试确保最终一致性。今日目标设计一个完整的订单生命周期消息链路包含创建通知、延迟支付提醒、超时取消、库存恢复。实现多级延迟消息使用 TTL 死信队列分别实现 30 分钟支付提醒和 10 分钟后最终取消。为消费者添加重试机制处理失败的消息自动重试指定次数超过上限进入死信队列便于人工介入。通过模拟订单生成验证消息流转的时序正确性和延迟精度。使用 RabbitMQ 管理界面观察整个链路的拓扑和消息轨迹。一、环境准备约 15 分钟确保 RabbitMQ 容器已运行hyperf-app项目已具备 AMQP 组件。docker-composeexecswoolebashcd/var/www/hyperf-app二、知识核心通知链设计与延迟消息精度约 1 小时1. 订单通知链业务逻辑我们定义如下的订单状态流转待支付 (pending)用户刚下单。已支付 (paid)用户完成支付。已取消 (cancelled)超时未支付自动取消。消息驱动流程order.created订单创建后立即触发库存预扣或实际扣减等同步逻辑同时发送一个延迟 30 分钟的消息到order.payment.reminder。order.payment.reminder30 分钟后消费者收到该消息检查订单状态若仍为pending则发送短信/邮件提醒用户支付并再发送一个延迟 10 分钟的消息到order.cancel。order.cancel消费者收到后再次检查订单状态若仍为pending则取消订单、恢复库存。此外每个消费者在处理消息时如果遇到异常如数据库临时不可用应进行重试重试次数用尽后转入死信队列由死信消费者记录日志并告警。2. 延迟消息实现方案选择我们继续采用TTL 死信队列的方式因为它与 RabbitMQ 原生功能配合不依赖额外插件。我们将创建两个延迟队列order.reminder.delay.queueTTL 30 分钟死信路由到order.reminder。order.cancel.delay.queueTTL 10 分钟死信路由到order.cancel。精度问题RabbitMQ 的 TTL 基于队列头消息的过期时间如果先入队的消息 TTL 很长后入队的短 TTL 消息会被阻塞直到前面消息过期。但因为我们两个延迟队列是独立的且队列内消息 TTL 相同30分钟/10分钟不会有阻塞问题。测试时可将 TTL 设为 30 秒和 15 秒快速验证。3. 重试机制设计hyperf/amqp原生在消费失败并返回NACK时可以通过设置requeuefalse将消息直接丢弃或进入死信。要实现带计数的重试我们需要在消息头中传递重试次数并配合多个队列或延迟重入。简化方案消费者捕获异常后检查消息头中的x-retry-count若未达最大次数如 3则将该消息重新发布到当前队列并在消息头中增加计数x-retry-count 1。可以利用 RabbitMQ 的priority或直接使用延迟队列进行退避。超过最大次数后返回NACK并指定requeuefalse消息进入配置的死信队列。今天我们将实现一种退避重试失败后不立即重试而是发送到延迟队列延迟一段时间如 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){$connnewAMQPStreamConnection(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分钟提醒延迟队列$reminderArgsnewAMQPTable([x-message-ttl1800000,// 30分钟测试可改为30000x-dead-letter-exchangeorder.exchange,x-dead-letter-routing-keyorder.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分钟取消延迟队列$cancelArgsnewAMQPTable([x-message-ttl600000,// 10分钟测试可改为15000x-dead-letter-exchangeorder.exchange,x-dead-letter-routing-keyorder.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);// 重试延迟队列退避用$retryArgsnewAMQPTable([x-message-ttl5000,// 5秒延迟x-dead-letter-exchangeorder.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分钟延迟提醒消息$delayMsgnew\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;$maxRetries3;if($retryCount$maxRetries){echo[订单创建] 重试耗尽转入死信\n;returnResult::NACK;// 因为配置了死信消息会进入死信}// 递增重试计数发送到延迟重试队列$data[_retry_count]$retryCount1;$retryMsgnew\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($statuspaid){echo[支付提醒] 订单已支付忽略\n;returnResult::ACK;}// 未支付发送提醒模拟echo[支付提醒] 发送短信/邮件提醒用户支付\n;// 发送10分钟延迟取消消息$cancelDelayMsgnewGenericProducer($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($statuspaid){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-duser_id1product_id1amount99控制台立即输出[订单创建] 开始处理订单 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 的消息不会重复处理业务五、今日作业与学习产出提交代码将拓扑监听器、所有消费者、通用生产者、订单控制器修改等提交到 Git。完善通知链集成真实的短信或邮件服务通过事件总线异步调用 API。使用RabbitMQ Delayed Message Plugin替代 TTLDLX简化延迟消息管理并对比两者的优劣。学习笔记画出完整的订单消息链路时序图标明各个队列、交换机和 TTL。总结基于消息队列实现最终一致性的核心要点幂等、重试、补偿。挑战任务实现动态 TTL通过消息属性expiration为每条消息单独设置 TTL而不依赖队列固定 TTL需使用x-dead-letter-exchange和目标路由键在消息属性中指定。使用Kafka实现类似的订单通知链对比两者的复杂度和性能。通过今天的综合实战你构建了一条企业级的消息驱动订单处理管线掌握了延迟消息、重试、死信和事件总线的组合拳。这标志着你已经能够运用异步消息解决复杂的分布式业务场景。下周我们将进入分布式事务的深水区探索 Saga 和事务消息。

相关推荐

【三个月 AI Agent 实战学习】Day 19:第一阶段测验准备 —— 手动模拟 ReAct 循环
【三个月 AI Agent 实战学习】Day 19:第一阶段测验准备 —— 手动模拟 ReAct 循环

Day 19:第一阶段测验准备 —— 手动模拟 ReAct 循环 欢迎来到第十九天!在前面的学习中,我们已经掌握了 Function Calling(Day 11),理解了模型如何请求调用工具。但 Function Calling 是 API 层面的机制&… · 2026/9/27 22:48:32

【八个月网安课程】第七周·周五:CSRF 防御——Token、Referer 校验、SameSite
【八个月网安课程】第七周·周五:CSRF 防御——Token、Referer 校验、SameSite

以下是第七周周五学习内容的详细展开。今天你将武装目标站点,从攻击者视角切换到防御者视角,系统学习并亲手配置 CSRF 的三大主流防御手段。你将看到昨天的攻击页面在加固后的环境中一一失效,并理解每种防御的精确边界与绕过风险。第七周周五… · 2026/9/27 22:48:32

烧烤店扫码点单选型:9 个维度对比加单、串数、分桌与后厨分单
烧烤店扫码点单选型:9 个维度对比加单、串数、分桌与后厨分单

烧烤店的单是餐饮里最不好扫的一类:一桌人边喝边加,既按串也按份,中途还有分桌、并台,凉菜、烤串、酒水要走出餐口。通用扫码点单放到烧烤场景,往往在“临时加单不重开单”和“按串/按份混合计价”这两步就卡住。本文把… · 2026/9/27 22:48:32

Xcelium混合仿真避坑指南:Verilog/VHDL/SystemC协同实战
Xcelium混合仿真避坑指南:Verilog/VHDL/SystemC协同实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:21

华为NE05E/NE08E物理层时钟同步配置与故障排查实战
华为NE05E/NE08E物理层时钟同步配置与故障排查实战

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:14

自动语音识别(ASR)技术全解析:从原理到工程实践
自动语音识别(ASR)技术全解析:从原理到工程实践

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:14

Y7000P Ubuntu 18.04 WiFi驱动修复:AIC8800编译安装教程
Y7000P Ubuntu 18.04 WiFi驱动修复:AIC8800编译安装教程

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:08

5分钟搞定ESP32/ESP8266开发环境:Arduino IDE 2.0国内镜像配置指南
5分钟搞定ESP32/ESP8266开发环境:Arduino IDE 2.0国内镜像配置指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:08

VSCODE加ESP-IDF配置指南:ESP32开发环境搭建与调试
VSCODE加ESP-IDF配置指南:ESP32开发环境搭建与调试

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views … · 2026/9/27 23:31:02

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码