分布式事务为什么是微服务架构的必解难题
单体应用拆分为微服务后,原本一个数据库事务能完成的工作,现在分散到多个服务中,每个服务有自己的数据库。跨服务的数据一致性保障就是分布式事务要解决的核心问题。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/