【基于 Swoole+Hyperf 的微服务实战】第七周·周四:事件总线与内部解耦 + 延迟消息实战
【基于 SwooleHyperf 的微服务实战】第七周·周四事件总线与内部解耦 延迟消息实战今天我们进入的主题是事件总线与内部解耦 延迟消息实战。前面我们已经掌握了 RabbitMQ 和 Kafka 的基本消息通信但在微服务内部业务逻辑往往需要在多个模块间传递事件如订单创建后需要发送短信、更新统计、触发风控。今天我们将利用 Hyperf 内置的事件机制与 RabbitMQ 结合构建一个松耦合、可扩展的事件总线并实现订单延迟取消这一典型的延迟消息场景让系统内部模块彻底解耦同时引入时间驱动的异步处理。今日目标深入理解 Hyperf 的 PSR-14 事件机制实现自定义事件与监听器的完全解耦。整合 RabbitMQ 作为事件总线的传输层让事件可以跨进程、跨服务传播。使用死信队列 TTL实现延迟消息完成下单 30 分钟未支付自动取消的完整流程。编写多个监听器发短信、记录日志、扣库存验证事件发布后自动异步执行。测试延迟消息的精准性观察消息在 TTL 后准时被消费。一、环境准备约 15 分钟继续使用现有的 Docker 环境确保 RabbitMQ 容器已启动。进入 PHP 容器docker-composeexecswoolebashcd/var/www/hyperf-app确认已安装hyperf/amqp和hyperf/eventHyperf 骨架自带。今天我们会复用昨天的 AMQP 配置。二、知识核心事件总线与延迟消息设计约 1 小时1. Hyperf 事件机制Hyperf 基于 PSR-14 提供了轻量级的事件组件事件Event任意 PHP 类封装需要传递的数据。监听器Listener实现Hyperf\Event\Contract\ListenerInterface在listen()中返回要监听的事件类。事件调度器EventDispatcher注入即可dispatch()事件所有注册的监听器会被自动调用。默认监听器在当前协程中同步执行但我们可以通过#[Listener]注解的async: true使其异步执行利用协程并发。如果想要更高可靠性的异步独立进程、持久化可以将事件通过消息队列发送再由独立消费者处理。2. 事件总线架构我们将事件分为本地事件和远程事件本地事件在同一个进程内通过 EventDispatcher 同步或异步处理适合实时性要求高的操作如缓存更新。远程事件通过 RabbitMQ 或 Kafka 发布由其他服务或独立消费者处理适合跨服务通信或耗时任务如发送邮件、风控审核。今天我们会为关键业务事件如OrderCreated同时发布本地事件异步处理缓存、日志和远程事件通知其他服务。3. 延迟消息的实现原理RabbitMQ 没有内置延迟队列但可以通过TTL生存时间 死信队列组合实现定义一个延迟队列设置x-message-ttl为延迟时间如 30 分钟并设置x-dead-letter-exchange和x-dead-letter-routing-key。消息发送到延迟队列后不会被任何消费者消费因为没有直接消费者绑定。消息在队列中驻留直到 TTL 过期自动变成“死信”被转发到指定的死信交换机进而路由到实际处理的队列。我们的取消订单消费者监听实际处理队列当收到消息时说明订单已到期。注意RabbitMQ 在 TTL 过期后才会检查并转发死信因此延迟时间的精度取决于队列的x-message-ttl和服务器时钟可能存在数秒误差对非实时场景足够。三、实战构建事件总线与订单延迟取消约 2.5 小时步骤 1定义订单创建事件本地新建app/Event/OrderCreatedEvent.php?phpnamespaceApp\Event;classOrderCreatedEvent{publicfunction__construct(publicarray$orderData){}}步骤 2创建本地监听器异步缓存、日志创建app/Listener/OrderCacheListener.php异步清除用户订单缓存?phpnamespaceApp\Listener;useApp\Event\OrderCreatedEvent;useHyperf\Event\Annotation\Listener;useHyperf\Event\Contract\ListenerInterface;useHyperf\Redis\Redis;useHyperf\Di\Annotation\Inject;#[Listener(async:true)]classOrderCacheListenerimplementsListenerInterface{#[Inject]privateRedis$redis;publicfunctionlisten():array{return[OrderCreatedEvent::class];}publicfunctionprocess(object$event){$userId$event-orderData[user_id];// 清除该用户的订单列表缓存$this-redis-del(user_orders:.$userId);echo[本地监听] 已清除用户{$userId}的订单缓存\n;}}创建app/Listener/OrderLogListener.php异步记录日志?phpnamespaceApp\Listener;useApp\Event\OrderCreatedEvent;useHyperf\Event\Annotation\Listener;useHyperf\Event\Contract\ListenerInterface;usePsr\Log\LoggerInterface;useHyperf\Di\Annotation\Inject;#[Listener(async:true)]classOrderLogListenerimplementsListenerInterface{#[Inject]privateLoggerInterface$logger;publicfunctionlisten():array{return[OrderCreatedEvent::class];}publicfunctionprocess(object$event){$this-logger-info(订单创建,$event-orderData);}}步骤 3创建远程事件生产者RabbitMQ我们希望订单创建事件也被外部服务感知因此创建一个专用的生产者。新建app/Amqp/Producer/OrderEventProducer.php?phpnamespaceApp\Amqp\Producer;useHyperf\Amqp\Annotation\Producer;useHyperf\Amqp\Message\ProducerMessage;#[Producer(exchange:order.event.exchange,routingKey:order.created)]classOrderEventProducerextendsProducerMessage{publicfunction__construct(array$orderData){$this-payload$orderData;}}步骤 4修改订单创建控制器集成事件总线在app/Controller/OrderController.php的create方法中注入事件调度器和远程生产者useApp\Event\OrderCreatedEvent;useApp\Amqp\Producer\OrderEventProducer;usePsr\EventDispatcher\EventDispatcherInterface;#[Inject]privateEventDispatcherInterface$eventDispatcher;#[Inject]privateOrderEventProducer$orderEventProducer;// 如果需要publicfunctioncreate(){// ... 订单数据生成$orderData[order_idrand(10000,99999),user_id$this-request-input(user_id,1),product_id$this-request-input(product_id,1),amount$this-request-input(amount,99.00),statuspending,created_atdate(Y-m-d H:i:s),];// 1. 发布本地事件异步缓存、日志$this-eventDispatcher-dispatch(newOrderCreatedEvent($orderData));// 2. 发送远程事件到 RabbitMQ 供其他服务消费$messagenewOrderEventProducer($orderData);$this-amqpProducer-produce($message);// 3. 发送延迟消息用于30分钟后自动取消重点$delayMessagenewOrderDelayProducer($orderData);$this-amqpProducer-produce($delayMessage);return[code201,message订单创建成功,data$orderData,];}步骤 5设计延迟队列拓扑与生产者我们需要延迟队列order.delay.queue设置 TTL 为 30 分钟死信交换机order.dlx.exchange死信路由键order.cancel。首先在 RabbitMQ 管理界面或通过代码声明拓扑。我们可以创建一个专用的生产者来声明延迟队列即使不通过它发送消息也可以用来声明交换机、队列和绑定。新建app/Amqp/Producer/OrderDelayProducer.php?phpnamespaceApp\Amqp\Producer;useHyperf\Amqp\Annotation\Producer;useHyperf\Amqp\Message\ProducerMessage;#[Producer(exchange:order.delay.exchange,routingKey:order.delay)]classOrderDelayProducerextendsProducerMessage{publicfunction__construct(array$orderData){$this-payload$orderData;}}但我们需要确保order.delay.queue被创建时带有 TTL 和死信参数。可以在消费者中声明然而延迟队列没有直接的消费者。一种方式是在配置中通过arguments在声明队列时加入参数。我们可以创建一个“假”消费者或者利用 Hyperf 的#[Producer]注解配合自定义队列声明实际上#[Producer]只会创建交换机和路由不会创建队列。我们需要在 RabbitMQ 中手动创建这个队列或者编写启动脚本。更简单的在docker-compose中启动一个一次性脚本来声明队列或者通过 RabbitMQ 管理界面 HTTP API。为了教学我们在项目启动时通过一个BootApplication监听器来创建。创建app/Listener/SetupDelayQueueListener.php?phpnamespaceApp\Listener;useHyperf\Event\Contract\ListenerInterface;useHyperf\Framework\Event\BootApplication;usePhpAmqpLib\Connection\AMQPStreamConnection;usePhpAmqpLib\Wire\AMQPTable;classSetupDelayQueueListenerimplementsListenerInterface{publicfunctionlisten():array{return[BootApplication::class];}publicfunctionprocess(object$event){$connectionnewAMQPStreamConnection(rabbitmq,5672,guest,guest);$channel$connection-channel();// 声明延迟交换机$channel-exchange_declare(order.delay.exchange,direct,false,true,false);// 声明延迟队列参数 TTL 和死信交换机$argsnewAMQPTable([x-message-ttl1800000,// 30分钟测试时可改为 30000 (30秒)x-dead-letter-exchangeorder.dlx.exchange,x-dead-letter-routing-keyorder.cancel,]);$channel-queue_declare(order.delay.queue,false,true,false,false,false,$args);$channel-queue_bind(order.delay.queue,order.delay.exchange,order.delay);// 死信交换机及取消队列已有但确保创建$channel-exchange_declare(order.dlx.exchange,direct,false,true,false);$channel-queue_declare(order.cancel.queue,false,true,false,false);$channel-queue_bind(order.cancel.queue,order.dlx.exchange,order.cancel);$channel-close();$connection-close();echo[初始化] 延迟队列拓扑创建完成\n;}}注意我们这里直接使用了php-amqplib它是hyperf/amqp的依赖已存在。如果不想用底层库可以调用 Hyperf 的 AMQP 管理方法。此监听器在框架启动时执行确保队列存在。测试时为快速看到效果将 TTL 设为 30 秒30000生产恢复 30 分钟。步骤 6创建订单取消消费者新建app/Amqp/Consumer/OrderCancelConsumer.php?phpnamespaceApp\Amqp\Consumer;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;#[Consumer(exchange:order.dlx.exchange,routingKey:order.cancel,queue:order.cancel.queue,name:OrderCancelConsumer,nums:1)]classOrderCancelConsumerextendsConsumerMessage{publicfunctionconsume($data):string{$orderId$data[order_id]??unknown;echo[订单取消] 检查订单{$orderId}...\n;// 模拟检查数据库若订单仍为 pending 状态则取消并恢复库存$status$data[status]??pending;// 实际情况应从数据库查询最新状态if($statuspending){echo[订单取消] 订单{$orderId}超时未支付已取消\n;// 恢复库存、更新状态等操作}else{echo[订单取消] 订单{$orderId}已支付无需取消\n;}returnResult::ACK;}}说明消费者绑定到死信交换机order.dlx.exchange和路由键order.cancel当延迟消息过期后会被投递到order.cancel.queue此消费者便会收到。步骤 7调整生产环境 TTL测试时我们希望 30 秒就能看到取消效果修改SetupDelayQueueListener中的 TTL 为30000(30秒)。实际生产应使用环境变量注入。四、成果测试与验证约 1 小时1. 测试本地事件解耦创建订单curl-XPOST http://localhost:9501/orders/create-duser_id1product_id1amount99立即观察控制台输出[本地监听] 已清除用户 1 的订单缓存订单创建日志出现在日志文件JSON 格式证明异步监听器执行且没有阻塞订单创建接口。2. 测试远程事件在 RabbitMQ 管理界面查看order.event.exchange应该能看到消息流入。如果另一个服务如通知服务绑定了相同队列就能收到事件。3. 测试延迟取消关键发送一个订单创建请求后在 30 秒内不要进行支付。30 秒后观察控制台出现[订单取消] 检查订单 12345... [订单取消] 订单 12345 超时未支付已取消如果订单在 30 秒内被支付模拟修改状态取消消费者应检查状态并跳过取消。验证时间精度记录消息发送时间和消费者收到时间差值大约为 30 秒允许数秒偏差。4. 测试清单检验项方法通过标准本地监听异步执行创建订单观察接口响应时间响应不等待监听器完成监听器日志稍后出现远程事件发布查看 RabbitMQ 管理界面order.event.exchange有消息发布延迟消息声明检查 RabbitMQ 中order.delay.queue属性TTL 和死信配置正确延迟消费发送订单 30 秒后取消消费者收到消息日志打印取消信息且业务逻辑执行幂等性重复消费同一取消消息如重启消费者不会重复取消事件总线解耦注释掉某个监听器创建订单仍成功失败不影响主流程五、今日作业与学习产出提交代码将事件类、监听器、延迟生产者、消费者、启动监听器等提交到 Git。完善事件总线添加一个短信通知监听器本地或远程当订单创建时异步发送短信可 Mock。将延迟消息的 TTL 迁移到 Nacos 配置中心实现动态调整延迟时间。学习笔记画出事件总线的架构图控制器 → 事件调度器 → 本地监听器 远程消息队列 → 外部消费者。总结延迟消息的两种实现方式TTLDLX 与 插件rabbitmq_delayed_message_exchange并比较优劣。挑战任务实现分级延迟例如订单创建后 10 分钟发送提醒短信30 分钟未支付则取消。需要创建不同 TTL 的队列或使用延迟插件。使用Kafka 的延迟消息如通过时间轮算法实现类似功能对比两者的实现复杂度。通过今天的学习你不仅实现了模块间的彻底解耦还掌握了时间驱动的异步处理利器——延迟消息。明天我们将整合 RabbitMQ 和 Kafka 的消息能力完成一个综合实战订单全流程的事件驱动架构。