如何在NestJS中实现带补偿事务的Saga模式(MongoDB、gRPC环境)
NestJS CQRS Saga 补偿事务最佳实践(MongoDB + gRPC)
针对NestJS结合MongoDB、gRPC的分布式事务场景,Saga模式的核心是一阶段提交+补偿回滚,而非强一致性事务。下面是落地的具体实现和关键原则,解决你提到的补偿逻辑缺失问题。
核心思路
- 用NestJS CQRS的
@Saga()监听领域事件,按顺序执行分布式步骤 - 每个步骤执行成功后发布完成事件,失败则触发对应的补偿事件
- MongoDB单服务内跨集合操作利用4.0+支持的事务,跨服务(gRPC调用)依赖补偿逻辑
- 所有操作必须保证幂等性,避免重复执行或补偿多次
代码实现示例
1. 定义领域事件(含补偿事件)
// src/order/events/order.events.ts export class OrderCreatedEvent { constructor( public readonly orderId: string, public readonly productId: string, public readonly quantity: number, public readonly totalAmount: number, public readonly requestId: string // 全局幂等标识 ) {} } // 成功事件 export class PaymentCompletedEvent { constructor(public readonly orderId: string, public readonly requestId: string) {} } export class InventoryDeductedEvent { constructor(public readonly orderId: string, public readonly requestId: string) {} } // 补偿事件 export class PaymentFailedEvent { constructor(public readonly orderId: string, public readonly requestId: string) {} } export class InventoryRollbackNeededEvent { constructor(public readonly orderId: string, public readonly requestId: string) {} }
2. Saga核心逻辑(含try/catch与补偿触发)
// src/order/sagas/order.saga.ts import { Injectable } from '@nestjs/common'; import { Saga, ofType } from '@nestjs/cqrs'; import { Observable, of } from 'rxjs'; import { catchError, switchMap } from 'rxjs/operators'; import { OrderCreatedEvent, PaymentCompletedEvent, PaymentFailedEvent, InventoryRollbackNeededEvent } from '../events/order.events'; import { PaymentService } from '../services/payment.service'; import { InventoryService } from '../services/inventory.service'; import { CommandBus } from '@nestjs/cqrs'; import { RollbackPaymentCommand } from '../commands/rollback-payment.command'; @Injectable() export class OrderSaga { constructor( private readonly paymentService: PaymentService, private readonly inventoryService: InventoryService, private readonly commandBus: CommandBus ) {} @Saga() orderTransactionFlow = (events$: Observable<any>): Observable<any> => { return events$.pipe( // 监听订单创建事件,启动Saga流程 ofType(OrderCreatedEvent), switchMap((event) => { // 第一步:调用gRPC支付服务 return this.paymentService.processPayment( event.orderId, event.totalAmount, event.requestId ).pipe( map(() => new PaymentCompletedEvent(event.orderId, event.requestId)), catchError(() => { // 支付失败:直接触发订单取消事件 return of(new PaymentFailedEvent(event.orderId, event.requestId)); }) ); }), switchMap((event) => { if (event instanceof PaymentFailedEvent) { // 支付失败流程:发布订单取消事件,结束流程 return of(new OrderCancelledEvent(event.orderId, event.requestId)); } // 第二步:扣减库存(本地MongoDB事务或gRPC调用) return this.inventoryService.deductStock( event.orderId, event.productId, event.quantity, event.requestId ).pipe( map(() => new InventoryDeductedEvent(event.orderId, event.requestId)), catchError(() => { // 库存扣减失败:触发支付补偿命令 return this.commandBus.execute( new RollbackPaymentCommand(event.orderId, event.requestId) ).pipe( map(() => new InventoryRollbackNeededEvent(event.orderId, event.requestId)) ); }) ); }) ); }; }
3. 补偿命令与处理器
// src/order/commands/rollback-payment.command.ts export class RollbackPaymentCommand { constructor( public readonly orderId: string, public readonly requestId: string ) {} } // src/order/handlers/rollback-payment.handler.ts import { CommandHandler, ICommandHandler } from '@nestjs/cqrs'; import { RollbackPaymentCommand } from '../commands/rollback-payment.command'; import { PaymentService } from '../services/payment.service'; @CommandHandler(RollbackPaymentCommand) export class RollbackPaymentHandler implements ICommandHandler<RollbackPaymentCommand> { constructor(private readonly paymentService: PaymentService) {} async execute(command: RollbackPaymentCommand) { // 调用gRPC支付服务的退款接口,通过requestId保证幂等 await this.paymentService.refundPayment(command.orderId, command.requestId); } }
4. MongoDB事务处理示例(库存服务)
// src/inventory/services/inventory.service.ts import { Injectable } from '@nestjs/common'; import { InjectModel } from '@nestjs/mongoose'; import { Model } from 'mongoose'; import { Inventory } from '../schemas/inventory.schema'; @Injectable() export class InventoryService { constructor( @InjectModel(Inventory.name) private inventoryModel: Model<Inventory> ) {} async deductStock(orderId: string, productId: string, quantity: number, requestId: string) { // 开启MongoDB会话事务 const session = await this.inventoryModel.startSession(); session.startTransaction(); try { // 幂等校验:检查当前requestId是否已执行过扣减 const existingRecord = await this.inventoryModel.findOne({ productId, 'operationHistory.requestId': requestId }).session(session); if (existingRecord) return; // 检查库存充足性 const inventory = await this.inventoryModel.findOne({ productId }).session(session); if (!inventory || inventory.stock < quantity) { throw new Error('库存不足'); } // 扣减库存并记录操作历史 await this.inventoryModel.updateOne( { productId }, { $inc: { stock: -quantity }, $push: { operationHistory: { requestId, orderId, type: 'DEDUCT' } } }, { session } ); await session.commitTransaction(); } catch (err) { await session.abortTransaction(); throw err; } finally { session.endSession(); } } // 库存补偿逻辑(恢复库存) async rollbackStock(productId: string, quantity: number, requestId: string) { const session = await this.inventoryModel.startSession(); session.startTransaction(); try { // 幂等校验:避免重复补偿 const existingRecord = await this.inventoryModel.findOne({ productId, 'operationHistory.requestId': requestId, 'operationHistory.type': 'ROLLBACK' }).session(session); if (existingRecord) return; await this.inventoryModel.updateOne( { productId }, { $inc: { stock: quantity }, $push: { operationHistory: { requestId, type: 'ROLLBACK' } } }, { session } ); await session.commitTransaction(); } catch (err) { await session.abortTransaction(); throw err; } finally { session.endSession(); } } }
关键最佳实践
- 幂等性是核心:所有操作(包括补偿)必须通过
requestId或orderId做幂等校验,避免重复执行 - 补偿逻辑独立:补偿命令与主逻辑解耦,单独实现处理器,确保补偿操作本身是可逆且幂等的
- MongoDB事务边界:单服务内跨集合/文档操作使用MongoDB事务,跨服务场景必须依赖补偿(MongoDB不支持跨集群事务)
- 异常分级处理:区分业务异常(如库存不足)和系统异常(如gRPC超时),业务异常触发补偿,系统异常可结合重试机制(需保证幂等)
- 日志与监控:每个Saga步骤、补偿操作都要记录带
requestId的详细日志,方便追踪问题
补偿失败处理
如果补偿操作(如退款)失败,需加入:
- 重试机制:用定时任务轮询补偿状态,最多重试3-5次
- 告警触发:重试失败后推送告警,通知人工介入
- 状态标记:在数据库中标记补偿失败的订单,避免重复触发
内容的提问来源于stack exchange,提问作者lifegoeson
相关产品推荐
相关产品推荐

