拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

【基于 Swoole+Hyperf 的微服务实战】第八周·周三:Saga 失败补偿实战

【基于 Swoole+Hyperf 的微服务实战】第八周·周三:Saga 失败补偿实战

今天我们进入第八周周三,主题是Saga 失败补偿实战。昨天我们编写了 Saga 协调器的基本框架,让它可以驱动正向流程并在失败时发出补偿命令。今天我们将让这些命令真正操作数据库,并重点验证支付失败后的全链路补偿:取消订单、释放冻结的库存,并确保无论发生什么异常,数据都不会处于中间不一致状态。这将是你首次亲手处理分布式事务的一致性恢复。


今日目标

  1. 创建orders、products表,模拟真实库存和订单数据。
  2. 改造订单服务和库存服务消费者,使其执行实际的数据库操作(创建订单、冻结库存、解冻、取消)。
  3. 人为让支付服务返回失败,触发 Saga 补偿链。
  4. 验证补偿完成后,订单状态为cancelled,库存恢复到下单前数量,无脏数据残留。
  5. 模拟补偿过程中的异常(如解冻失败),观察协调器重试和死信机制,确保最终成功或告警。

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

继续在hyperf-app容器中操作,确保 MySQL、RabbitMQ 可用。

docker-composeexecswoolebashcd/var/www/hyperf-app
创建业务表

我们新增两张表来模拟库存和订单服务各自的数据库(实际微服务中它们在不同库,今天在同一个库中模拟)。

执行 SQL(可通过数据库工具或 PHP 脚本):

CREATETABLEIFNOTEXISTS`products`(`id`INTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,`name`VARCHAR(255)NOTNULL,`total_stock`INTNOTNULLDEFAULT0COMMENT'总库存',`frozen_stock`INTNOTNULLDEFAULT0COMMENT'冻结库存',`available_stock`INTAS(total_stock-frozen_stock)STOREDCOMMENT'可用库存',`created_at`TIMESTAMPDEFAULTCURRENT_TIMESTAMP,`updated_at`TIMESTAMPDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP)ENGINE=InnoDB;CREATETABLEIFNOTEXISTS`orders`(`id`INTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,`order_id`VARCHAR(64)NOTNULLUNIQUE,`user_id`INTNOTNULL,`product_id`INTNOTNULL,`amount`DECIMAL(10,2)NOTNULL,`status`ENUM('pending','frozen','paid','cancelled')NOTNULLDEFAULT'pending',`created_at`TIMESTAMPDEFAULTCURRENT_TIMESTAMP,`updated_at`TIMESTAMPDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP)ENGINE=InnoDB;-- 初始化商品库存INSERTINTO`products`(`id`,`name`,`total_stock`)VALUES(1,'Hyperf微服务实战',100);

二、知识核心:补偿的幂等、空回滚与防悬挂(约 1 小时)

1. 补偿操作的三大挑战
  • 幂等性:同一个补偿命令可能因网络重试而执行多次(如order.cancel被重复发送)。必须保证多次执行结果相同,例如取消订单时检查订单状态,若已取消则不再操作。
  • 空回滚:当创建订单本身失败时,协调器可能仍发出order.cancel补偿命令。此时订单不存在,补偿操作必须忽略(而非报错)。
  • 防悬挂:补偿命令可能比正向命令先到达(如网络乱序)。取消订单时,若订单尚未创建,则需等待或直接忽略(不能取消一个还没创建的订单)。

我们的实现策略:

  • 各服务在消费命令前,根据saga_id和step检查操作日志表(或 Redis)是否已执行,实现幂等。
  • 对于空回滚和防悬挂,业务逻辑应能处理“订单不存在”的情况,并正常返回成功。
2. 库存冻结与解冻的隔离性

采用冻结库存模式:下单时先冻结库存(frozen_stock增加),支付成功后再实际扣减(total_stock减少,同时frozen_stock减少)。取消时解冻库存(frozen_stock减少)。这样可以避免超卖,并且取消时无需复杂计算。

SQL 示例:

  • 冻结:UPDATE products SET frozen_stock = frozen_stock + 1 WHERE id = ? AND total_stock - frozen_stock >= 1
  • 解冻:UPDATE products SET frozen_stock = frozen_stock - 1 WHERE id = ? AND frozen_stock > 0
3. 支付失败时的补偿链

正向流程:创建订单 → 冻结库存 → 执行扣款
支付失败触发补偿:

  1. 退款(若扣款成功才需,今天场景支付未成功,所以跳过退款)
  2. 解冻库存
  3. 取消订单

协调器按逆序发送补偿命令。


三、实战:改造服务消费者,验证支付失败补偿(约 2.5 小时)

我们基于昨天的代码进行增强。

步骤 1:创建操作日志表(幂等辅助)

为了避免重复执行,每个服务维护一个简单的操作日志表saga_step_logs:

CREATETABLEIFNOTEXISTS`saga_step_logs`(`id`BIGINTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,`saga_id`VARCHAR(64)NOTNULL,`step`VARCHAR(50)NOTNULL,`status`VARCHAR(20)NOTNULLDEFAULT'success',`created_at`TIMESTAMPDEFAULTCURRENT_TIMESTAMP,UNIQUEKEY`uk_saga_step`(`saga_id`,`step`))ENGINE=InnoDB;

每次处理命令前,尝试插入一条记录,若主键冲突则视为已处理。

步骤 2:改造订单服务消费者

编辑app/Amqp/Consumer/OrderCommandConsumer.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:'order.*',queue:'order.command.queue',name:'OrderCommandConsumer',nums:1,type:'topic')]classOrderCommandConsumerextendsConsumerMessage{#[Inject]privateProducer$producer;publicfunctionconsume($data):string{$sagaId=$data['saga_id'];$step=$data['step'];$payload=$data['payload']??[];// 幂等检查if($this->isProcessed($sagaId,$step)){echo"[订单服务] 步骤{$step}已处理,跳过\n";returnResult::ACK;}try{if($step===SagaConstants::STEP_ORDER_CREATE){$this->createOrder($payload);}elseif($step===SagaConstants::COMPENSATE_ORDER_CANCEL){$this->cancelOrder($payload);}$this->markProcessed($sagaId,$step);$this->reply($sagaId,$step,'success');}catch(\Throwable$e){echo"[订单服务] 处理失败: ".$e->getMessage()."\n";$this->reply($sagaId,$step,'failed');}returnResult::ACK;}privatefunctioncreateOrder(array$payload):void{$orderId=$payload['order_id'];Db::table('orders')->insert(['order_id'=>$orderId,'user_id'=>$payload['user_id'],'product_id'=>$payload['product_id'],'amount'=>$payload['amount'],'status'=>'pending',]);echo"[订单服务] 订单{$orderId}创建成功\n";}privatefunctioncancelOrder(array$payload):void{$orderId=$payload['order_id'];$affected=Db::table('orders')->where('order_id',$orderId)->update(['status'=>'cancelled']);if($affected===0){echo"[订单服务] 订单{$orderId}不存在,忽略取消\n";}else{echo"[订单服务] 订单{$orderId}已取消\n";}}privatefunctionisProcessed(string$sagaId,string$step):bool{returnDb::table('saga_step_logs')->where(['saga_id'=>$sagaId,'step'=>$step])->exists();}privatefunctionmarkProcessed(string$sagaId,string$step):void{Db::table('saga_step_logs')->insert(['saga_id'=>$sagaId,'step'=>$step]);}privatefunctionreply(string$sagaId,string$step,string$status):void{/* 同昨天 */}}
步骤 3:改造库存服务消费者

新建/编辑app/Amqp/Consumer/InventoryCommandConsumer.php,实现冻结与解冻:

#[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 products SET frozen_stock = frozen_stock + 1 WHERE id = ? AND total_stock - frozen_stock >= 1',[$productId]);if($affected===0){thrownew\Exception('库存不足');}echo"[库存服务] 冻结商品{$productId}库存\n";}privatefunctionunfreezeStock(array$payload):void{$productId=$payload['product_id'];Db::update('UPDATE products SET frozen_stock = frozen_stock - 1 WHERE id = ? AND frozen_stock > 0',[$productId]);echo"[库存服务] 解冻商品{$productId}库存\n";}// ... 幂等和回复方法}
步骤 4:设置支付服务强制失败

为了触发补偿,我们修改支付消费者app/Amqp/Consumer/PaymentCommandConsumer.php,使其始终返回失败(或随机)。

privatefunctionprocessDebit(array$payload):void{// 模拟支付失败,可以随机或固定thrownew\Exception('支付网关超时');}

并在回复中发送status: failed。

步骤 5:验证补偿流程

重启所有消费者,确保拓扑存在。使用saga/start启动一个 Saga:

curl-XPOST http://localhost:9501/saga/start-d"user_id=1&product_id=1&amount=99"

观察控制台日志:

[订单服务] 订单创建成功 [协调器] 发送命令: inventory.freeze [库存服务] 冻结库存 [协调器] 发送命令: payment.debit [支付服务] 处理扣款失败 [协调器] 开始补偿... [协调器] 补偿命令已发送 [库存服务] 解冻库存 [订单服务] 订单取消 [协调器] Saga 状态 failed

查询数据库:

  • orders表中订单状态为cancelled。
  • products表中frozen_stock恢复为 0,total_stock不变(仍为100)。
  • saga_transactions状态failed。
  • saga_step_logs包含正向的 create、freeze、debit(失败),以及补偿的 unfreeze、cancel。
步骤 6:模拟补偿失败与死信

临时注释掉解冻库存的 SQL,或者故意让解冻失败(如抛出异常)。此时补偿中的解冻步骤会失败,协调器应如何处理?目前我们的消费者只是回复失败,协调器会认为补偿失败,但昨天设计的协调器在补偿失败后没有进一步重试。我们可以增加一个补偿重试逻辑:在协调器的startCompensation中,如果某个补偿命令失败(需要监听补偿步骤的回复),则重新发送或转入死信。

一个简单实现:在InventoryCommandConsumer中,如果解冻失败,不回复成功,而是返回 NACK,让消息重入队列,或者触发重试机制(trait RetryableConsumer)。同时协调器等待补偿步骤的回复,若多次失败后仍未成功,则更新 Saga 状态为compensation_failed,并向死信队列发送告警。

由于时间关系,今天重点在于确保补偿逻辑正确执行并保证数据一致性。死信重试可作为挑战任务。


四、成果测试与数据一致性验证(约 1 小时)

1. 正常补偿验证

执行以下 SQL 检查数据:

-- 订单已取消SELECT*FROMordersWHEREorder_id='xxx';-- 库存恢复SELECT*FROMproductsWHEREid=1;-- 操作日志完整SELECT*FROMsaga_step_logsWHEREsaga_id='xxx'ORDERBYcreated_at;

所有正向和补偿步骤都应有记录,且frozen_stock = 0。

2. 重复启动同一 Saga 的幂等性

使用相同的saga_id再次调用 start 接口(或手动向命令队列发送相同消息),各服务应跳过所有步骤,日志表无重复记录。

3. 并发启动多个 Saga

同时启动多个下单请求,观察库存冻结量是否等于正在进行中的订单数,取消后全部释放。

4. 测试清单
检验项方法通过标准
支付失败触发补偿启动 Saga,观察日志依次执行 create、freeze、debit(失败)、unfreeze、cancel
订单状态查 orders 表status = 'cancelled'
库存解冻查 products 表frozen_stock = 0,total_stock不变
幂等性重复发送命令操作日志唯一,无重复业务操作
空回滚订单未创建时发送 cancel取消操作被忽略,不报错
补偿失败重试手动让解冻失败消费者重试(需实现),最终进入死信
协调器状态查 saga_transactionsstatus 为 failed,步骤记录正确

五、今日作业与学习产出

  1. 提交代码:将改造后的订单、库存、支付消费者,操作日志相关逻辑提交。
  2. 增强协调器:
    • 为补偿步骤增加超时与重试:协调器等待补偿回复,若超时则重新发送。
    • 增加补偿失败死信处理:超过最大重试次数后,将 Saga 标记为compensation_failed,并发送钉钉通知。
  3. 学习笔记:
    • 总结补偿事务中“幂等性”、“空回滚”、“防悬挂”的实现方法。
    • 绘制支付失败后的补偿时序图,包含服务、协调器和数据库。
  4. 挑战任务:
    • 实现TCC 模式的部分思路:将冻结改为try,扣款改为confirm,补偿改为cancel,并对比 Saga 的区别。
    • 为 Saga 流程编写自动化集成测试,使用 Docker 环境的 RabbitMQ 和 MySQL,断言最终数据一致性。

通过今天的学习,你已经掌握了 Saga 模式在失败场景下的核心能力——补偿恢复。现在你的分布式事务协调器已经能够在真实数据库环境中保证数据的最终一致。明天我们将为这个系统加上分布式锁,防止并发下的超卖,进一步完善系统的隔离性。

返回列表