【基于 Swoole+Hyperf 的微服务实战】第九周·周五: 集成分布式事务 Saga 与延迟队列
今天我们进入第九周周五,核心任务是集成分布式事务 Saga 与延迟队列,实现下单→冻结库存→支付扣款的正向流程,以及支付失败时的补偿回滚,外加30 分钟未支付自动取消订单的延迟处理。这将补全电商系统最关键的一致性保障,让整个交易链路从“能跑通”进化为“能容错”。
今日目标
- 创建支付服务(轻量级),仅用于接收 Saga 命令并模拟支付结果。
- 在订单服务中实现Saga 协调器,驱动正向步骤:创建订单 → 冻结库存 → 支付扣款。
- 在库存服务中增加 AMQP 消费者,处理
inventory.freeze/unfreeze命令。 - 利用RabbitMQ 延迟队列(TTL+DLX)实现下单 30 分钟未支付自动取消,恢复库存。
- 修改网关下单接口为异步模式:接受订单返回
saga_id,不再同步等待结果。 - 测试正向成功流程、支付失败补偿流程以及超时自动取消,验证最终一致性。
一、环境准备与支付服务创建(约 45 分钟)
1. 启动依赖
docker-composeup-d确保 RabbitMQ、MySQL、Redis、Consul 均已运行。
2. 创建支付服务
docker-composeexecswoolebashcd/var/www/servicescomposercreate-project hyperf/hyperf-skeleton paymentcdpaymentcomposerrequire hyperf/amqp hyperf/redis hyperf/json-rpc hyperf/service-governance-consul支付服务只需要 AMQP 消费者,不需要 HTTP 服务端,所以我们只配置 AMQP 和 Consul 注册即可。端口设为9506(或可省略,但我们保留一个简单的 HTTP 服务用于健康检查)。
配置:
.env中APP_NAME=PaymentServiceconfig/autoload/amqp.php指向rabbitmq容器config/autoload/services.php配置 Consul 注册(可省略,因为支付服务不对外提供 RPC,但为了统一监控,可以注册)- 不需要数据库连接(模拟支付不存储数据)
二、知识核心:Saga 协调器与延迟队列设计(约 1 小时)
1. Saga 正向步骤(以订单创建为起点)
我们采用协调器式 Saga,消息驱动:
- 创建订单:订单服务接收命令
order.create,写入订单(状态pending),发送回复。 - 冻结库存:库存服务接收
inventory.freeze,冻结库存,回复成功/失败。 - 支付扣款:支付服务接收
payment.debit,模拟扣款(随机成功/失败),回复结果。
若任何一步失败,协调器启动补偿:
- 支付失败→ 无需退款(未扣款成功)→ 解冻库存 → 取消订单
- 冻结库存失败→ 取消订单
2. 延迟取消队列
订单创建成功后,发送一条延迟消息到延迟队列,TTL 为 30 分钟(测试时改为 30 秒)。过期后变为死信,路由到order.cancel队列,由订单服务的取消消费者处理。取消消费者检查订单状态,若仍为pending或frozen,则执行取消(调用库存解冻、更新订单为cancelled)。若已支付,则忽略。
三、实战:编码 Saga 与延迟队列(约 2.5 小时)
步骤 1:在库存服务中添加 AMQP 命令消费者
进入inventory-service,安装 AMQP(可能已安装?之前没有,需要添加):
cd/var/www/services/inventorycomposerrequire hyperf/amqp创建app/Amqp/Consumer/InventoryCommandConsumer.php:
<?phpnamespaceApp\Amqp\Consumer;useApp\Saga\SagaConstants;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;useHyperf\Amqp\Producer;useHyperf\Di\Annotation\Inject;useHyperf\DbConnection\Db;#[Consumer(exchange:SagaConstants::EXCHANGE_COMMANDS,routingKey:'inventory.*',queue:'inventory.command.queue',name:'InventoryCommandConsumer',nums:1,type:'topic')]classInventoryCommandConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;publicfunctionconsume($data):string{$sagaId=$data['saga_id'];$step=$data['step'];$payload=$data['payload']??[];if($this->isProcessed($sagaId,$step))returnResult::ACK;try{if($step===SagaConstants::STEP_INVENTORY_FREEZE){$this->freezeStock($payload);}elseif($step===SagaConstants::COMPENSATE_INVENTORY_UNFREEZE){$this->unfreezeStock($payload);}$this->markProcessed($sagaId,$step);$this->reply($sagaId,$step,'success');}catch(\Throwable$e){$this->reply($sagaId,$step,'failed');}returnResult::ACK;}privatefunctionfreezeStock(array$payload):void{$productId=$payload['product_id'];// 分布式锁已在内部实现,重试逻辑略$affected=Db::update('UPDATE stocks SET frozen = frozen + 1 WHERE product_id = ? AND total - frozen >= 1',[$productId]);if($affected===0)thrownew\Exception('库存不足');}privatefunctionunfreezeStock(array$payload):void{$productId=$payload['product_id'];Db::update('UPDATE stocks SET frozen = frozen - 1 WHERE product_id = ? AND frozen > 0',[$productId]);}privatefunctionisProcessed(string$sagaId,string$step):bool{/* 同前 */}privatefunctionmarkProcessed(string$sagaId,string$step):void{/* 同前 */}privatefunctionreply(string$sagaId,string$step,string$status):void{$msg=new\App\Amqp\Producer\GenericProducer(['saga_id'=>$sagaId,'step'=>$step,'status'=>$status,],SagaConstants::EXCHANGE_REPLIES,'saga.reply');$this->producer->produce($msg);}}注意:库存服务需要创建SagaConstants和GenericProducer,可以从订单服务复制一份通用代码。
步骤 2:在支付服务中创建支付命令消费者
进入payment-service,创建app/Amqp/Consumer/PaymentCommandConsumer.php:
<?phpnamespaceApp\Amqp\Consumer;useApp\Saga\SagaConstants;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;useHyperf\Amqp\Producer;useHyperf\Di\Annotation\Inject;#[Consumer(exchange:SagaConstants::EXCHANGE_COMMANDS,routingKey:'payment.*',queue:'payment.command.queue',name:'PaymentCommandConsumer',nums:1,type:'topic')]classPaymentCommandConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;publicfunctionconsume($data):string{$sagaId=$data['saga_id'];$step=$data['step'];$payload=$data['payload']??[];if($this->isProcessed($sagaId,$step))returnResult::ACK;// 模拟支付:80% 成功$success=rand(1,100)<=80;$status=$success?'success':'failed';$this->markProcessed($sagaId,$step);$this->reply($sagaId,$step,$status);returnResult::ACK;}// 幂等、回复方法同上}步骤 3:在订单服务中实现 Saga 协调器消费者
进入order-service,安装 AMQP(之前已安装?可能需要):
cd/var/www/services/ordercomposerrequire hyperf/amqp创建app/Amqp/Consumer/SagaCoordinatorConsumer.php(直接从第八周复制并调整步骤映射):
<?phpnamespaceApp\Amqp\Consumer;useApp\Saga\SagaConstants;useHyperf\Amqp\Annotation\Consumer;useHyperf\Amqp\Message\ConsumerMessage;useHyperf\Amqp\Result;useHyperf\Amqp\Producer;useHyperf\Di\Annotation\Inject;useHyperf\DbConnection\Db;#[Consumer(exchange:SagaConstants::EXCHANGE_REPLIES,routingKey:'saga.reply',queue:'saga.reply.queue',name:'SagaCoordinator',nums:1)]classSagaCoordinatorConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;privatearray$stepFlow=[SagaConstants::STEP_ORDER_CREATE=>['success'=>SagaConstants::STEP_INVENTORY_FREEZE,'compensate'=>SagaConstants::COMPENSATE_ORDER_CANCEL,],SagaConstants::STEP_INVENTORY_FREEZE=>['success'=>SagaConstants::STEP_PAYMENT_DEBIT,'compensate'=>SagaConstants::COMPENSATE_INVENTORY_UNFREEZE,],SagaConstants::STEP_PAYMENT_DEBIT=>['success'=>null,'compensate'=>SagaConstants::COMPENSATE_PAYMENT_REFUND,],];publicfunctionconsume($data):string{// ... 同第八周协调器逻辑,处理成功/失败并发送下一步或补偿命令}}需要创建SagaConstants类定义步骤常量,以及GenericProducer。
步骤 4:创建订单命令消费者
在订单服务中创建app/Amqp/Consumer/OrderCommandConsumer.php,处理order.create和order.cancel,并发送回复。
order.create:插入orders表(status=‘pending’),回复成功。order.cancel:更新订单状态为cancelled,若不存在则忽略,回复成功。
步骤 5:延迟取消消息发送
在OrderCommandConsumer的createOrder逻辑中,创建订单后,发送一条延迟消息到order.delay.exchange,路由键order.delay,消息体包含order_id和product_id。该消息进入 TTL 30 分钟(测试 30 秒)的延迟队列,过期后死信路由到order.cancel。
我们需要在拓扑中声明延迟队列(可通过启动监听器或 RabbitMQ 管理界面预先创建)。
步骤 6:创建订单取消消费者
订单服务中新建app/Amqp/Consumer/OrderCancelConsumer.php,监听order.cancel队列,收到消息后:
- 根据
order_id查询订单状态,若为pending或frozen,则更新为cancelled,并发送inventory.unfreeze命令(或直接调用库存服务的 AMQP 命令?为了解耦,通过 Saga 补偿机制?但超时取消独立于 Saga。可以直接调用库存 RPC 解冻,或发送一个补偿事件。简单起见,使用 RPC 调用库存解冻,然后更新订单状态。此处我们直接通过 RPC 调用库存服务的unfreeze接口(使用 RPC 客户端)实现。
因为取消消费者在订单服务中,可以注入InventoryServiceInterfaceRPC 客户端,调用解冻。确保幂等。
步骤 7:修改网关下单接口为异步
修改网关的OrderController::create(),不再同步调用订单服务 RPC 的createOrder,而是调用订单服务新暴露的submitOrderRPC,该 RPC 负责生成saga_id、写入saga_transactions、发送第一条命令order.create,并立即返回saga_id和order_id。客户凭此查询订单状态。
我们在订单服务中新增一个 RPC 方法submitOrder,完成上述操作。或者网关可以直接调用订单服务的submitOrder。
订单服务新增接口:submitOrder(int $userId, int $productId, float $amount): array
实现:
- 生成
order_id和saga_id - 插入
saga_transactions,状态running,当前步骤order.create - 发送
order.create命令到 Saga 命令交换机 - 返回
['saga_id' => $sagaId, 'order_id' => $orderId]
网关 OrderController改为调用submitOrder,并返回处理中信息。
步骤 8:配置所有拓扑
编写一个SetupSagaTopologyListener(或在启动脚本中)创建所需的交换机、队列、绑定,包括延迟队列。为快速,可以在 RabbitMQ 管理界面手动创建:
- 交换机:
saga.commands(topic)、saga.replies(direct)、order.delay.exchange(direct) - 队列:
order.command.queue、inventory.command.queue、payment.command.queue、saga.reply.queue、order.cancel.queue - 延迟队列:
order.delay.queue,绑定order.delay,TTL=30000ms,死信交换机order.dlx.exchange(这里简化,死信直接路由到order.cancel队列?我们需要将order.delay.queue的死信路由到saga.commands交换机并指定路由键order.cancel?或者直接让取消消费者绑定另一个交换机。为了简单,我们让延迟队列的死信交换机为saga.commands,路由键为order.cancel,这样取消命令进入命令总线,由订单服务的取消消费者处理。这需要OrderCancelConsumer监听路由键order.cancel。
调整:在InventoryCommandConsumer中我们监听inventory.*,同样订单命令消费者监听order.*,那么order.cancel也会被OrderCommandConsumer处理,我们可以在该消费者的consume中根据step调用cancelOrder。所以不需要单独的取消消费者,只需在OrderCommandConsumer中增加对order.cancel的处理。这样更统一。
于是OrderCommandConsumer支持:order.create、order.cancel。order.cancel的取消逻辑中,查询订单状态,若未支付则调用库存 RPC 解冻并更新订单为cancelled。注意幂等。
四、启动与集成测试(约 1 小时)
1. 启动所有服务
分别进入user-service,product-service,inventory-service,order-service,payment-service,gateway并启动。确认 RabbitMQ 中出现相关队列和消费者。
2. 正向流程测试
TOKEN=$(curl-s-XPOST http://localhost:9500/auth/login...)# 提交订单curl-XPOST http://localhost:9500/orders-H"Authorization: Bearer$TOKEN"-d'{"product_id":1,"amount":99.00}'# 返回 saga_id 和 order_id观察协调器日志,应依次处理order.create→inventory.freeze→payment.debit。若支付成功,最终订单状态应为paid(需在支付成功时更新状态,目前我们未实现,可在OrderCommandConsumer中处理payment.debit的回复?实际上回复由协调器接收,支付成功后协调器完成 Saga,不会通知订单服务变更状态。因此我们需要在 Saga 完成后向订单服务发送确认命令。可以在协调器的handleSuccess中,当payment.debit成功时,发送一个order.confirm命令到订单服务,订单服务将其状态改为paid。或者简化:在OrderCommandConsumer中,order.create时状态设为pending,当支付成功由协调器通知订单服务,订单服务更新状态。
完善协调器逻辑:步骤映射中,payment.debit成功后,发送order.confirm命令。订单服务处理order.confirm将状态改为paid。这样正向流程才完整。
3. 支付失败补偿测试
支付消费者返回failed,协调器启动补偿,依次发送inventory.unfreeze、order.cancel。库存解冻,订单取消,数据一致。
4. 超时取消测试
不下达支付成功指令,等待 30 秒(测试 TTL),观察取消消费者自动取消订单并恢复库存。
5. 幂等和重试测试
手动重复发送命令,观察操作日志表,确保不会重复处理。
五、今日作业与学习产出
- 提交代码:支付服务、各服务的 Saga 消费者、协调器、延迟队列配置等。
- 完善状态同步:
- 支付成功后通知订单服务更新状态为
paid。 - 补偿完成后更新 Saga 状态为
failed。
- 支付成功后通知订单服务更新状态为
- 学习笔记:
- 画出 Saga + 延迟取消的完整状态机图。
- 对比同步 RPC 与异步 Saga 的优缺点。
- 挑战任务:
- 使用Kafka替代 RabbitMQ 实现命令通道,并保证消息顺序。
- 为 Saga 协调器添加定时任务扫描长时间未完成的 Saga 主动查询各服务状态。
今天你成功地将分布式事务带入了电商系统,使订单流程具备了企业级的容错与补偿能力。下周我们将进行项目复盘、压力测试和调优,让整个系统更加健壮。