NestJS中Saga模式实现:补偿事务最佳实践咨询
NestJS CQRS Saga 事务与补偿机制最佳实践示例
核心思路
Saga模式的核心是通过多个本地事务+补偿操作实现最终一致性。当Saga流程中任意步骤失败时,需要对之前已成功执行的步骤执行反向补偿操作,回滚已变更的数据。在NestJS CQRS中,try/catch是实现补偿逻辑的基础手段,但需要结合事务边界、幂等性等要点来保证可靠性。
完整代码示例
1. 定义业务命令与补偿命令
首先定义流程中的正向业务命令和对应的反向补偿命令:
// src/order/commands/order.commands.ts import { ICommand } from '@nestjs/cqrs'; // 正向命令:创建订单 export class CreateOrderCommand implements ICommand { constructor( public readonly orderId: string, public readonly userId: string, public readonly productId: string, public readonly quantity: number, ) {} } // 正向命令:扣减库存 export class DeductStockCommand implements ICommand { constructor( public readonly productId: string, public readonly quantity: number, ) {} } // 补偿命令:取消订单 export class CancelOrderCommand implements ICommand { constructor(public readonly orderId: string) {} } // 补偿命令:恢复库存 export class RestoreStockCommand implements ICommand { constructor( public readonly productId: string, public readonly quantity: number, ) {} }
2. Saga 流程实现(含补偿逻辑)
在Saga中,通过try/catch包裹业务步骤,捕获异常后按顺序触发补偿操作:
// src/order/sagas/order.saga.ts import { Injectable } from '@nestjs/common'; import { Saga, ofType, ICommand, CommandBus } from '@nestjs/cqrs'; import { Observable, concatMap, of } from 'rxjs'; import { OrderInitiatedEvent } from '../events/order-initiated.event'; import { CreateOrderCommand, DeductStockCommand, CancelOrderCommand, RestoreStockCommand, } from '../commands/order.commands'; @Injectable() export class OrderSaga { constructor(private readonly commandBus: CommandBus) {} @Saga() handleOrderInitiated = (events$: Observable<any>): Observable<ICommand> => { return events$.pipe( ofType(OrderInitiatedEvent), concatMap(async (event) => { const { orderId, userId, productId, quantity } = event; // 标记已执行的步骤,用于确定需要补偿的操作 const executedSteps: string[] = []; try { // 步骤1:创建订单(本地事务) await this.commandBus.execute( new CreateOrderCommand(orderId, userId, productId, quantity), ); executedSteps.push('order_created'); // 步骤2:扣减库存(本地事务) await this.commandBus.execute( new DeductStockCommand(productId, quantity), ); executedSteps.push('stock_deducted'); // 更多业务步骤... } catch (error) { // 按逆序执行补偿操作 const compensationCommands: ICommand[] = []; if (executedSteps.includes('stock_deducted')) { compensationCommands.push(new RestoreStockCommand(productId, quantity)); } if (executedSteps.includes('order_created')) { compensationCommands.push(new CancelOrderCommand(orderId)); } // 批量执行补偿 await Promise.all( compensationCommands.map((cmd) => this.commandBus.execute(cmd)), ); // 重新抛出异常,让上层感知流程失败 throw error; } return of(); }), ); }; }
3. 命令处理器(保证本地事务原子性)
每个正向/补偿命令的处理器内部,要通过ORM事务保证单个操作的原子性:
// src/order/handlers/deduct-stock.handler.ts import { CommandHandler, ICommandHandler } from '@nestjs/cqrs'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { Product } from '../../product/entities/product.entity'; import { DeductStockCommand } from '../commands/order.commands'; @CommandHandler(DeductStockCommand) export class DeductStockHandler implements ICommandHandler<DeductStockCommand> { constructor( @InjectRepository(Product) private readonly productRepo: Repository<Product>, ) {} async execute(command: DeductStockCommand): Promise<void> { // 开启TypeORM本地事务 await this.productRepo.manager.transaction(async (em) => { const product = await em.findOne(Product, { where: { id: command.productId }, }); if (!product || product.stock < command.quantity) { throw new Error('库存不足,无法扣减'); } product.stock -= command.quantity; await em.save(product); }); } }
4. 补偿命令处理器(实现幂等性)
补偿操作必须保证幂等,避免重复执行导致数据错误:
// src/order/handlers/restore-stock.handler.ts import { CommandHandler, ICommandHandler } from '@nestjs/cqrs'; import { InjectRepository } from '@nestjs/typeorm'; import { Repository } from 'typeorm'; import { Product } from '../../product/entities/product.entity'; import { RestoreStockCommand } from '../commands/order.commands'; @CommandHandler(RestoreStockCommand) export class RestoreStockHandler implements ICommandHandler<RestoreStockCommand> { constructor( @InjectRepository(Product) private readonly productRepo: Repository<Product>, ) {} async execute(command: RestoreStockCommand): Promise<void> { await this.productRepo.manager.transaction(async (em) => { const product = await em.findOne(Product, { where: { id: command.productId }, }); if (!product) { throw new Error('商品不存在'); } // 幂等处理:可通过事件日志或状态标记避免重复恢复 // 示例:假设我们有库存操作日志,先检查是否已执行过该补偿 // const exists = await em.findOne(StockOperationLog, { ... }); // if (exists) return; product.stock += command.quantity; await em.save(product); }); } }
关键最佳实践
- 事务隔离:每个命令处理器只负责单个本地事务,Saga仅做流程协调,不跨服务持有事务。
- 幂等优先:所有补偿操作必须实现幂等,可通过唯一标识(如补偿请求ID)或状态标记实现。
- 异常区分:区分可重试异常(如网络抖动)和不可重试异常(如业务规则失败),可重试异常建议先重试再触发补偿。
- 状态追踪:记录Saga每一步的执行状态(成功/失败),便于排查问题和手动恢复。
- 避免长流程:拆分大Saga为多个小Saga,减少补偿范围,降低一致性风险。
内容的提问来源于stack exchange,提问作者rizesky
相关产品推荐
相关产品推荐

