You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在NestJS中实现带补偿事务的Saga模式(MongoDB、gRPC环境)

NestJS CQRS Saga 补偿事务最佳实践(MongoDB + gRPC)

针对NestJS结合MongoDB、gRPC的分布式事务场景,Saga模式的核心是一阶段提交+补偿回滚,而非强一致性事务。下面是落地的具体实现和关键原则,解决你提到的补偿逻辑缺失问题。

核心思路

  1. 用NestJS CQRS的@Saga()监听领域事件,按顺序执行分布式步骤
  2. 每个步骤执行成功后发布完成事件,失败则触发对应的补偿事件
  3. MongoDB单服务内跨集合操作利用4.0+支持的事务,跨服务(gRPC调用)依赖补偿逻辑
  4. 所有操作必须保证幂等性,避免重复执行或补偿多次

代码实现示例

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();
    }
  }
}

关键最佳实践

  1. 幂等性是核心:所有操作(包括补偿)必须通过requestId或orderId做幂等校验,避免重复执行
  2. 补偿逻辑独立:补偿命令与主逻辑解耦,单独实现处理器,确保补偿操作本身是可逆且幂等的
  3. MongoDB事务边界:单服务内跨集合/文档操作使用MongoDB事务,跨服务场景必须依赖补偿(MongoDB不支持跨集群事务)
  4. 异常分级处理:区分业务异常(如库存不足)和系统异常(如gRPC超时),业务异常触发补偿,系统异常可结合重试机制(需保证幂等)
  5. 日志与监控:每个Saga步骤、补偿操作都要记录带requestId的详细日志,方便追踪问题

补偿失败处理

如果补偿操作(如退款)失败,需加入:

  • 重试机制:用定时任务轮询补偿状态,最多重试3-5次
  • 告警触发:重试失败后推送告警,通知人工介入
  • 状态标记:在数据库中标记补偿失败的订单,避免重复触发

内容的提问来源于stack exchange,提问作者lifegoeson

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.09 03:16:21