尧图精选

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

🕒 发布时间:2026/10/1 17:15:56 📁 来源:尧图网络
今天我们进入第八周周三主题是Saga 失败补偿实战。昨天我们编写了 Saga 协调器的基本框架让它可以驱动正向流程并在失败时发出补偿命令。今天我们将让这些命令真正操作数据库并重点验证支付失败后的全链路补偿取消订单、释放冻结的库存并确保无论发生什么异常数据都不会处于中间不一致状态。这将是你首次亲手处理分布式事务的一致性恢复。今日目标创建orders、products表模拟真实库存和订单数据。改造订单服务和库存服务消费者使其执行实际的数据库操作创建订单、冻结库存、解冻、取消。人为让支付服务返回失败触发 Saga 补偿链。验证补偿完成后订单状态为cancelled库存恢复到下单前数量无脏数据残留。模拟补偿过程中的异常如解冻失败观察协调器重试和死信机制确保最终成功或告警。一、环境准备约 20 分钟继续在hyperf-app容器中操作确保 MySQL、RabbitMQ 可用。docker-composeexecswoolebashcd/var/www/hyperf-app创建业务表我们新增两张表来模拟库存和订单服务各自的数据库实际微服务中它们在不同库今天在同一个库中模拟。执行 SQL可通过数据库工具或 PHP 脚本CREATETABLEIFNOTEXISTSproducts(idINTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,nameVARCHAR(255)NOTNULL,total_stockINTNOTNULLDEFAULT0COMMENT总库存,frozen_stockINTNOTNULLDEFAULT0COMMENT冻结库存,available_stockINTAS(total_stock-frozen_stock)STOREDCOMMENT可用库存,created_atTIMESTAMPDEFAULTCURRENT_TIMESTAMP,updated_atTIMESTAMPDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP)ENGINEInnoDB;CREATETABLEIFNOTEXISTSorders(idINTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,order_idVARCHAR(64)NOTNULLUNIQUE,user_idINTNOTNULL,product_idINTNOTNULL,amountDECIMAL(10,2)NOTNULL,statusENUM(pending,frozen,paid,cancelled)NOTNULLDEFAULTpending,created_atTIMESTAMPDEFAULTCURRENT_TIMESTAMP,updated_atTIMESTAMPDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP)ENGINEInnoDB;-- 初始化商品库存INSERTINTOproducts(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 03. 支付失败时的补偿链正向流程创建订单 → 冻结库存 → 执行扣款支付失败触发补偿退款若扣款成功才需今天场景支付未成功所以跳过退款解冻库存取消订单协调器按逆序发送补偿命令。三、实战改造服务消费者验证支付失败补偿约 2.5 小时我们基于昨天的代码进行增强。步骤 1创建操作日志表幂等辅助为了避免重复执行每个服务维护一个简单的操作日志表saga_step_logsCREATETABLEIFNOTEXISTSsaga_step_logs(idBIGINTUNSIGNEDAUTO_INCREMENTPRIMARYKEY,saga_idVARCHAR(64)NOTNULL,stepVARCHAR(50)NOTNULL,statusVARCHAR(20)NOTNULLDEFAULTsuccess,created_atTIMESTAMPDEFAULTCURRENT_TIMESTAMP,UNIQUEKEYuk_saga_step(saga_id,step))ENGINEInnoDB;每次处理命令前尝试插入一条记录若主键冲突则视为已处理。步骤 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($stepSagaConstants::STEP_ORDER_CREATE){$this-createOrder($payload);}elseif($stepSagaConstants::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],statuspending,]);echo[订单服务] 订单{$orderId}创建成功\n;}privatefunctioncancelOrder(array$payload):void{$orderId$payload[order_id];$affectedDb::table(orders)-where(order_id,$orderId)-update([statuscancelled]);if($affected0){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($stepSagaConstants::STEP_INVENTORY_FREEZE){$this-freezeStock($payload);}elseif($stepSagaConstants::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];$affectedDb::update(UPDATE products SET frozen_stock frozen_stock 1 WHERE id ? AND total_stock - frozen_stock 1,[$productId]);if($affected0){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启动一个 Sagacurl-XPOST http://localhost:9501/saga/start-duser_id1product_id1amount99观察控制台日志[订单服务] 订单创建成功 [协调器] 发送命令: inventory.freeze [库存服务] 冻结库存 [协调器] 发送命令: payment.debit [支付服务] 处理扣款失败 [协调器] 开始补偿... [协调器] 补偿命令已发送 [库存服务] 解冻库存 [订单服务] 订单取消 [协调器] Saga 状态 failed查询数据库orders表中订单状态为cancelled。products表中frozen_stock恢复为 0total_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_idxxx;-- 库存恢复SELECT*FROMproductsWHEREid1;-- 操作日志完整SELECT*FROMsaga_step_logsWHEREsaga_idxxxORDERBYcreated_at;所有正向和补偿步骤都应有记录且frozen_stock 0。2. 重复启动同一 Saga 的幂等性使用相同的saga_id再次调用 start 接口或手动向命令队列发送相同消息各服务应跳过所有步骤日志表无重复记录。3. 并发启动多个 Saga同时启动多个下单请求观察库存冻结量是否等于正在进行中的订单数取消后全部释放。4. 测试清单检验项方法通过标准支付失败触发补偿启动 Saga观察日志依次执行 create、freeze、debit(失败)、unfreeze、cancel订单状态查 orders 表status cancelled库存解冻查 products 表frozen_stock 0total_stock不变幂等性重复发送命令操作日志唯一无重复业务操作空回滚订单未创建时发送 cancel取消操作被忽略不报错补偿失败重试手动让解冻失败消费者重试需实现最终进入死信协调器状态查 saga_transactionsstatus 为 failed步骤记录正确五、今日作业与学习产出提交代码将改造后的订单、库存、支付消费者操作日志相关逻辑提交。增强协调器为补偿步骤增加超时与重试协调器等待补偿回复若超时则重新发送。增加补偿失败死信处理超过最大重试次数后将 Saga 标记为compensation_failed并发送钉钉通知。学习笔记总结补偿事务中“幂等性”、“空回滚”、“防悬挂”的实现方法。绘制支付失败后的补偿时序图包含服务、协调器和数据库。挑战任务实现TCC 模式的部分思路将冻结改为try扣款改为confirm补偿改为cancel并对比 Saga 的区别。为 Saga 流程编写自动化集成测试使用 Docker 环境的 RabbitMQ 和 MySQL断言最终数据一致性。通过今天的学习你已经掌握了 Saga 模式在失败场景下的核心能力——补偿恢复。现在你的分布式事务协调器已经能够在真实数据库环境中保证数据的最终一致。明天我们将为这个系统加上分布式锁防止并发下的超卖进一步完善系统的隔离性。
上一篇/下一篇内容由系统自动关联 返回资讯列表 →