分布式事务实战:Saga模式在微服务架构中的落地指南

分布式事务为什么是微服务架构的必解难题

单体应用拆分为微服务后,原本一个数据库事务能完成的工作,现在分散到多个服务中,每个服务有自己的数据库。跨服务的数据一致性保障就是分布式事务要解决的核心问题。2PC(两阶段提交)虽然能保证强一致性,但在微服务场景下性能损耗过大且存在单点阻塞风险。Saga模式作为补偿事务方案,在保证最终一致性的前提下,允许服务间松耦合协作,是目前微服务架构中分布式事务落地的首选方案。

Saga模式的核心原理与两种实现

Saga将一个长事务拆分为多个本地事务,每个本地事务完成后发布事件触发下一个本地事务。如果某个步骤失败,则反向执行已完成步骤的补偿事务。Saga有两种编排方式:

Choreography(事件驱动编排):每个服务完成本地事务后发布事件,下游服务监听事件执行。没有中心协调器。

Orchestration(中心协调器编排):由一个Saga协调器(Orchestrator)统一控制事务流程,向各服务发送命令。协调器维护全局状态机。

两种方式对比:

  • Choreography:服务间松耦合,适合简单流程(3-5步),但流程不可见、调试困难
  • Orchestration:流程集中管理、状态清晰、易调试,但协调器有单点风险,适合复杂业务流程

生产环境推荐Orchestration模式——流程可视化和调试能力比松耦合更重要。

Saga Orchestration实战:订单创建流程

以电商订单创建为例,涉及三个服务:订单服务、库存服务、支付服务。流程如下:

// Saga定义:订单创建流程
const createOrderSaga = new SagaDefinition('create-order')
  .step('create-order')
    .invoke((ctx) => orderService.createOrder(ctx.payload))
    .compensate((ctx) => orderService.cancelOrder(ctx.result.orderId))
  .step('reserve-inventory')
    .invoke((ctx) => inventoryService.reserve(ctx.result.orderId, ctx.payload.items))
    .compensate((ctx) => inventoryService.release(ctx.result.orderId))
  .step('process-payment')
    .invoke((ctx) => paymentService.charge(ctx.result.orderId, ctx.payload.amount))
    .compensate((ctx) => paymentService.refund(ctx.result.orderId))
  .build()

Saga协调器核心实现(TypeScript/Node.js)

// saga/SagaOrchestrator.ts
interface SagaStep<T> {
  name: string
  invoke: (context: SagaContext<T>) => Promise<any>
  compensate?: (context: SagaContext<T>) => Promise<void>
}

interface SagaContext<T> {
  payload: T
  results: Map<string, any>
  status: 'running' | 'compensating' | 'completed' | 'failed'
}

class SagaOrchestrator<T> {
  private steps: SagaStep<T>[] = []
  private executed: number[] = []

  addStep(step: SagaStep<T>): this {
    this.steps.push(step)
    return this
  }

  async execute(payload: T): Promise<void> {
    const context: SagaContext<T> = {
      payload,
      results: new Map(),
      status: 'running'
    }

    // 正向执行
    for (let i = 0; i < this.steps.length; i++) {
      const step = this.steps[i]
      try {
        const result = await step.invoke(context)
        context.results.set(step.name, result)
        this.executed.push(i)
        
        // 持久化进度(关键:保证崩溃恢复)
        await this.saveCheckpoint(context, i)
      } catch (error) {
        context.status = 'compensating'
        await this.compensate(context, i)
        throw new SagaExecutionError(
          `Step "${step.name}" failed, compensation executed`, error
        )
      }
    }

    context.status = 'completed'
    await this.clearCheckpoint(context)
  }

  private async compensate(context: SagaContext<T>, failedAt: number): Promise<void> {
    // 反向执行已成功步骤的补偿事务
    const completedSteps = this.executed.reverse()
    for (const i of completedSteps) {
      const step = this.steps[i]
      if (step.compensate) {
        try {
          await step.compensate(context)
        } catch (compensateError) {
          // 补偿失败:记录日志,人工介入
          await this.logCompensationFailure(step.name, compensateError)
          // 继续尝试补偿后续步骤
        }
      }
    }
  }

  private async saveCheckpoint(context: SagaContext<T>, stepIndex: number): Promise<void> {
    // 将当前进度写入数据库,保证崩溃后能恢复
    await db.sagaCheckpoint.upsert({
      sagaId: context.results.get('sagaId'),
      currentStep: stepIndex,
      status: context.status,
      results: JSON.stringify(Object.fromEntries(context.results)),
      updatedAt: new Date()
    })
  }
}

补偿事务设计的三个原则

Saga模式下,补偿事务(Compensating Transaction)的设计质量直接决定业务安全性。

原则1:补偿事务必须幂等

网络超时、重试等场景下,补偿事务可能被执行多次。幂等性保证重复执行不会产生额外副作用:

// 库存释放:幂等实现
async release(orderId: string): Promise<void> {
  // 先查询再操作,而非直接扣减
  const reservation = await db.inventoryReservation.findUnique({
    where: { orderId }
  })
  if (!reservation || reservation.status === 'released') {
    return  // 已释放或不存在,直接返回
  }
  
  await db.$transaction([
    db.inventory.update({
      where: { skuId: reservation.skuId },
      data: { available: { increment: reservation.quantity } }
    }),
    db.inventoryReservation.update({
      where: { orderId },
      data: { status: 'released' }
    })
  ])
}

原则2:补偿事务尽量语义补偿而非物理回滚

物理回滚(如恢复删除的记录)在高并发场景下容易冲突。语义补偿(如退款、发优惠券)更安全:

// 支付补偿:退款而非删除支付记录
async refund(orderId: string): Promise<void> {
  const payment = await db.payment.findUnique({ where: { orderId } })
  if (!payment || payment.status === 'refunded') return

  // 调用支付网关退款
  await paymentGateway.refund(payment.transactionId, payment.amount)
  
  // 更新状态而非删除
  await db.payment.update({
    where: { orderId },
    data: { status: 'refunded', refundedAt: new Date() }
  })
}

原则3:补偿失败必须有兜底机制

补偿事务自身也可能失败(如支付网关不可用)。必须设计兜底机制:

  • 重试队列:补偿失败后写入重试队列,由定时任务周期性重试
  • 告警:连续重试3次仍失败,触发人工告警
  • 对账:每日对账程序自动检测不一致数据并修复

Saga状态机与消息队列集成

生产环境中Saga协调器需要与消息队列(如Kafka/RabbitMQ)集成,实现异步执行和可靠传递:

// Kafka + Saga集成
class KafkaSagaOrchestrator extends SagaOrchestrator {
  private producer: KafkaProducer
  private consumer: KafkaConsumer

  async startSaga(sagaType: string, payload: any): Promise<string> {
    const sagaId = uuid()
    
    // 发布Saga启动事件
    await this.producer.send({
      topic: `saga.${sagaType}.start`,
      messages: [{
        key: sagaId,
        value: JSON.stringify({ sagaId, payload, timestamp: Date.now() })
      }]
    })
    
    return sagaId
  }

  async handleStepResult(message: KafkaMessage): Promise<void> {
    const { sagaId, stepName, result, error } = JSON.parse(message.value)
    
    if (error) {
      // 步骤失败,触发补偿流程
      await this.executeCompensation(sagaId, stepName)
    } else {
      // 步骤成功,推进到下一步
      await this.advanceSaga(sagaId, stepName, result)
    }
  }
}

消息队列保证Sage步骤的可靠传递和幂等消费,是生产环境的关键基础设施。

隔离级别与脏读问题

Saga的AB隔离性(缺乏隔离性)是最大的业务风险。中间状态可能被其他事务读到,导致脏读。解决方案:

语义锁:在每个步骤中给业务对象加状态标记,其他事务读取时检查状态:

// 订单创建Saga中的语义锁
async createOrder(payload: OrderPayload): Promise<Order> {
  return db.order.create({
    data: {
      ...payload,
      status: 'PENDING',  // 语义锁:其他事务看到PENDING状态时知道未完成
      createdAt: new Date()
    }
  })
}

// 其他服务查询订单时过滤PENDING状态
const activeOrders = await db.order.findMany({
  where: { status: { not: 'PENDING' } }
})

交换律设计:保证Saga步骤的执行顺序不影响最终结果。例如库存扣减和支付两个步骤,先执行哪个不影响最终一致性。

Saga模式不是银弹——它牺牲了强一致性换取可用性和性能。在电商、金融等业务场景中,理解补偿事务设计三原则(幂等、语义补偿、兜底机制)并配合消息队列实现可靠传递,才能在生产环境中稳定运行。

原创文章,作者:小编,如若转载,请注明出处:https://www.yunthe.com/fen-bu-shi-shi-wu-shi-zhan-saga-mo-shi-zai-wei-fu-wu-jia/

(0)
小编小编
上一篇 17小时前
下一篇 17小时前

相关推荐