分布式事务与 Saga 编排

讲解 .NET 微服务中跨服务事务的解决方案,涵盖 Saga 补偿事务、编排与协同的取舍、幂等与去重、Outbox 模式可靠投递,以及状态机实现与故障恢复。

1. 分布式事务的困境

一句话总结: 跨服务的事务无法靠数据库的本地事务保证,两阶段提交在微服务里代价过高,现实选择是「最终一致 + 补偿」。

单体应用里,一次下单是 BEGIN TRANSACTION 到 COMMIT 的一段代码:扣库存、建订单、写支付单,要么全成功要么全回滚。拆成微服务后,这三件事分属三个数据库、三个进程,本地事务保护不了跨进程的原子性。

两阶段提交(2PC)理论上能解决,但它在微服务里几乎不可用:协调者成为单点,参与者持有锁直到全局提交,网络分区时长期阻塞,且很多现代存储(Kafka、部分 NoSQL、外部支付网关)根本不支持 XA 协议。

方案原子性可用性复杂度现实可行性
本地事务强高低单库内
2PC/XA强低高跨库但锁阻塞
TCC强(业务层)中高资金类核心链路
Saga最终一致高中主流选择
事件驱动 + 对账最终一致高中异步业务
// 单体时代的写法:一个事务保护三张表
using var tx = await db.Database.BeginTransactionAsync(ct);
db.Orders.Add(order);
db.Stock.Add(new StockChange(order.Id, -1));
db.Payments.Add(payment);
await db.SaveChangesAsync(ct);
await tx.CommitAsync(ct);
// 拆分后,这三件事分属三个服务,无法共享事务

避坑: 分布式事务的难点从来不是「怎么提交」,而是「怎么处理中间态」。订单已创建但库存扣减失败时,那笔订单是「待确认」还是「已失败」?必须为每个步骤定义中间状态,并让业务能容忍它短暂存在。设计时若假设「不会失败」,补偿逻辑就永远写不对。

2. Saga 模型与补偿事务

一句话总结: Saga 把长事务拆成一串本地事务,每一步都配一个语义相反的补偿动作,任一步失败就逆序执行已成功步骤的补偿。

Saga 的核心概念只有两个:正向步骤(T1, T2, ..., Tn)与补偿步骤(C1, C2, ..., Cn)。若 T3 失败,则依次执行 C2、C1 把系统恢复到一致状态。

关键认知:补偿不是回滚。回滚是物理撤销(数据回到没发生过的样子),补偿是语义抵消(再发一笔反向操作)。扣了 100 元库存,补偿是加回 100 元,而不是「假装没扣过」——中间可能有别的请求也动过这份库存。

// 每一步都实现正向与补偿两个动作
public interface ISagaStep
{
    string Name { get; }
    Task ExecuteAsync(SagaContext ctx, CancellationToken ct);
    Task CompensateAsync(SagaContext ctx, CancellationToken ct);
}

public sealed class ReserveStockStep(IStockService stock) : ISagaStep
{
    public string Name => "reserve-stock";

    public async Task ExecuteAsync(SagaContext ctx, CancellationToken ct)
    {
        var reservationId = await stock.ReserveAsync(ctx.OrderId, ctx.Items, ct);
        ctx.Set("reservationId", reservationId);   // 供补偿使用
    }

    public Task CompensateAsync(SagaContext ctx, CancellationToken ct)
        => stock.ReleaseAsync(ctx.Get<Guid>("reservationId"), ct);
}
正向动作补偿动作补偿是否可逆
创建订单取消订单可逆
预占库存释放库存可逆
扣减余额退款可逆但需时间
发送短信无法撤回不可逆,放最后
调用外部支付发起退款不可逆,需幂等
写入审计日志不补偿只追加

避坑: 补偿动作本身也可能失败(网络抖动、下游宕机)。Saga 必须定义补偿的重试策略与最终兜底——通常是把补偿也做成幂等操作并无限重试,直到成功或转入人工处理队列。另外,不可逆动作要放在 Saga 最后:发短信、寄快递、调用第三方支付这类无法撤销的步骤,应排在所有可补偿步骤之后,把失败窗口压到最小。

3. 编排与协同的取舍

一句话总结: 编排(Orchestration)由中心协调器发号施令,流程清晰可观测;协同(Choreography)由服务各自监听事件自发推进,耦合低但流程散落难追踪。

两种模式不是优劣之分,而是复杂度放在哪里的选择。编排把流程逻辑集中在一个状态机里,容易理解和排错,但协调器承担了额外的可用性责任;协同没有中心,服务间通过事件解耦,但「这个订单卡在哪一步」需要跨多个服务拼日志才能回答。

// 编排:协调器显式调用每一步,流程一目了然
public sealed class OrderSagaOrchestrator(
    IEnumerable<ISagaStep> steps,
    ISagaRepository repo,
    ILogger<OrderSagaOrchestrator> logger)
{
    public async Task RunAsync(SagaContext ctx, CancellationToken ct)
    {
        var executed = new Stack<ISagaStep>();

        foreach (var step in steps)
        {
            try
            {
                await step.ExecuteAsync(ctx, ct);
                executed.Push(step);
                await repo.SaveAsync(ctx, step.Name, ct);   // 持久化进度
            }
            catch (Exception ex)
            {
                logger.LogError(ex, "步骤 {Step} 失败,开始补偿", step.Name);
                await CompensateAsync(executed, ctx, ct);
                throw;
            }
        }
    }

    private static async Task CompensateAsync(
        Stack<ISagaStep> executed, SagaContext ctx, CancellationToken ct)
    {
        while (executed.TryPop(out var step))
        {
            try { await step.CompensateAsync(ctx, ct); }
            catch (Exception ex)
            {
                // 补偿失败必须记录并交给重试机制,不能吞掉
                ctx.RecordCompensationFailure(step.Name, ex);
            }
        }
    }
}
// 协同:服务各自订阅事件,自发推进,无中心协调者
public sealed class StockEventConsumer(IStockService stock, IEventBus bus)
    : IConsumer<OrderCreated>
{
    public async Task ConsumeAsync(OrderCreated e, CancellationToken ct)
    {
        try
        {
            await stock.ReserveAsync(e.OrderId, e.Items, ct);
            await bus.PublishAsync(new StockReserved(e.OrderId), ct);
        }
        catch (Exception ex)
        {
            await bus.PublishAsync(new StockReservationFailed(e.OrderId, ex.Message), ct);
        }
    }
}
维度编排协同
流程可见性集中,易追踪分散,需拼日志
服务耦合服务依赖协调器服务只依赖事件
单点风险协调器需高可用无中心
循环依赖天然避免容易出现事件环
适用步数步骤多、分支多步骤少、线性
调试难度低高

避坑: 协同模式最危险的是事件环:服务 A 发事件触发 B,B 又发事件触发 A,循环放大流量直至系统崩溃。防御手段是给事件带上「因果链 ID」并限制链长。另一个坑是责任真空:没人负责整体流程,超时无人处理,订单永远卡在中间态——必须有一个超时扫描任务兜底。

4. 幂等与去重

一句话总结: 在至少一次投递的消息系统里,重复消费是常态而非异常,每个消息处理器都必须幂等,去重键要落在业务语义上。

消息队列普遍提供「至少一次」投递:网络超时后重发、消费者处理完但 ACK 丢失、分区重平衡都会导致同一条消息被处理两次。若处理器不幂等,就会出现扣两次库存、发两笔退款。

// 幂等处理器:用业务去重键 + 唯一索引保证只生效一次
public sealed class PaymentConsumer(
    AppDbContext db,
    IPaymentGateway gateway) : IConsumer<OrderPlaced>
{
    public async Task ConsumeAsync(OrderPlaced e, CancellationToken ct)
    {
        var key = $"payment:{e.OrderId}";     // 业务去重键

        // 依赖数据库唯一索引:并发下只有一个能插入成功
        if (await db.ProcessedMessages.AnyAsync(m => m.Key == key, ct))
            return;                            // 已处理,直接跳过

        await gateway.ChargeAsync(e.OrderId, e.Amount, key, ct);

        db.ProcessedMessages.Add(new ProcessedMessage(key, DateTimeOffset.UtcNow));
        await db.SaveChangesAsync(ct);         // 唯一索引冲突则说明并发重复
    }
}
// 更稳的做法:业务写入与去重记录放进同一个本地事务
await using var tx = await db.Database.BeginTransactionAsync(ct);
db.Payments.Add(new Payment(e.OrderId, e.Amount, PaymentStatus.Pending));
db.ProcessedMessages.Add(new ProcessedMessage(key, DateTimeOffset.UtcNow));
await db.SaveChangesAsync(ct);              // 业务与去重同时提交,原子
await tx.CommitAsync(ct);
await gateway.ChargeAsync(e.OrderId, e.Amount, ct);   // 外部调用放事务外
去重键来源稳定性适用
消息 ID高(同一条消息重投时不变)通用去重
业务 ID高同一业务只应生效一次
业务 ID + 动作高一个业务多个动作
内容哈希中无业务 ID 时兜底
时间戳低不推荐

避坑: 去重表会无限增长,必须有清理策略:保留最近 7~30 天的记录,更早的删除(因为消息重投不会跨越那么久)。另一个坑是先调用外部接口再写去重记录——如果外部调用成功但写记录失败,重试时会重复扣款。正确顺序是把外部调用的幂等键传给对方(多数支付网关支持 Idempotency-Key),让对方去重。

5. Outbox 模式与可靠投递

一句话总结: 业务数据写入与消息发布必须原子,Outbox 模式把「要发的消息」先写进同库的 outbox 表,再由后台任务投递,从根本上消除「数据写了消息没发」。

「更新数据库 + 发消息」是经典的双写问题:数据库提交成功但发消息失败,下游永远收不到;发消息成功但数据库回滚,下游收到不该发生的消息。Outbox 把这两步变成一步:消息作为一行记录,与业务数据在同一个本地事务里提交。

// 同一个事务里写业务数据与 outbox 记录
public async Task PlaceOrderAsync(Order order, CancellationToken ct)
{
    await using var tx = await db.Database.BeginTransactionAsync(ct);

    db.Orders.Add(order);
    db.Outbox.Add(new OutboxMessage
    {
        Id = Guid.NewGuid(),
        Type = nameof(OrderPlaced),
        Payload = JsonSerializer.Serialize(new OrderPlaced(order.Id, order.Items)),
        OccurredOn = DateTimeOffset.UtcNow,
        ProcessedOn = null,
    });

    await db.SaveChangesAsync(ct);   // 业务与消息一起提交,原子
    await tx.CommitAsync(ct);
}
// 后台投递器:扫描未投递的 outbox 记录,发送后标记已处理
public sealed class OutboxPublisher(
    IDbContextFactory<AppDbContext> factory,
    IEventBus bus) : BackgroundService
{
    protected override async Task ExecuteAsync(CancellationToken ct)
    {
        using var timer = new PeriodicTimer(TimeSpan.FromSeconds(2));
        while (await timer.WaitForNextTickAsync(ct))
        {
            await using var db = await factory.CreateDbContextAsync(ct);
            var batch = await db.Outbox.Where(m => m.ProcessedOn == null)
                .OrderBy(m => m.OccurredOn).Take(100).ToListAsync(ct);

            foreach (var msg in batch)
            {
                await bus.PublishAsync(msg.Type, msg.Payload, ct);   // 至少一次
                msg.ProcessedOn = DateTimeOffset.UtcNow;
            }
            await db.SaveChangesAsync(ct);
        }
    }
}
投递语义实现方式结果
至多一次先标记后发送可能丢消息
至少一次先发送后标记可能重复,需幂等
事务性发件箱同库事务 + 后台投递不丢,可能重复
变更数据捕获订阅数据库日志不丢,低延迟

避坑: Outbox 天然是「至少一次」,重复投递无法避免,所以消费端必须幂等——这是第 4 节与第 5 节必须成对出现的原因。另一个坑是投递顺序:多实例并行投递时,同一实体的两个事件可能乱序到达。若业务对顺序敏感,需要在消息里带版本号或序列号,由消费端丢弃过期消息。用 CDC(如 Debezium)订阅 binlog 可以进一步降低延迟,但要额外维护一套采集链路。

6. 状态机实现

一句话总结: 把 Saga 的每一步与状态显式建模成状态机,用持久化状态驱动推进,服务重启后能从断点继续。

编排器如果只存在于内存里,进程崩溃后流程就丢了。正确做法是把 Saga 状态持久化,每次状态变更都落库,恢复时从数据库读回状态继续推进。MassTransit 的 MassTransitStateMachine 提供了成熟的实现,自研时至少要定义清楚状态、事件与转换。

// MassTransit 状态机:声明状态与事件驱动的转换
public sealed class OrderStateMachine : MassTransitStateMachine<OrderSagaState>
{
    public State StockReserved { get; private set; } = null!;
    public State PaymentCharged { get; private set; } = null!;
    public Event<OrderPlaced> OrderPlaced { get; private set; } = null!;
    public Event<StockReserved> StockReservedEvent { get; private set; } = null!;
    public Event<PaymentFailed> PaymentFailed { get; private set; } = null!;

    public OrderStateMachine()
    {
        InstanceState(x => x.CurrentState);
        Event(() => OrderPlaced, e => e.CorrelateById(ctx => ctx.Message.OrderId));
        Event(() => StockReservedEvent, e => e.CorrelateById(ctx => ctx.Message.OrderId));
        Event(() => PaymentFailed, e => e.CorrelateById(ctx => ctx.Message.OrderId));
        Initially(
            When(OrderPlaced)
                .Then(ctx => ctx.Saga.OrderId = ctx.Message.OrderId)
                .TransitionTo(StockReserved)
                .Publish(ctx => new ReserveStock(ctx.Saga.OrderId, ctx.Message.Items)));

        During(StockReserved,
            When(StockReservedEvent)
                .TransitionTo(PaymentCharged)
                .Publish(ctx => new ChargePayment(ctx.Saga.OrderId)),
            When(PaymentFailed)
                .TransitionTo(Compensating)
                .Publish(ctx => new ReleaseStock(ctx.Saga.OrderId)));   // 补偿
    }
}
// 状态持久化:EF Core 存储,Saga 状态随事务落库
services.AddMassTransit(x =>
    x.AddSagaStateMachine<OrderStateMachine, OrderSagaState>()
     .EntityFrameworkRepository(r =>
     {
         r.ExistingDbContext<AppDbContext>();
         r.UsePostgres();
     }));
状态触发事件下一状态动作
InitialOrderPlacedStockReserved预占库存
StockReservedStockReservedPaymentCharged扣款
StockReservedStockFailedCompensating取消订单
PaymentChargedPaymentFailedCompensating释放库存 + 取消订单
CompensatingCompensatedCancelled通知用户
PaymentChargedCompletedCompleted完成订单

避坑: 状态机的超时必须显式处理。如果库存服务挂了,Saga 永远停在 StockReserved 状态,订单永远卡住。MassTransit 支持在状态上挂 Schedule 超时,超时后触发补偿。另一个坑是状态机版本升级:新增状态或转换后,正在飞行中的老实例可能无法匹配新定义,需要保证状态机的向后兼容,或对存量实例做数据迁移。

7. 可观测性与故障恢复

一句话总结: Saga 的排错依赖「一条链路能看全」,需要把相关性 ID 贯穿所有消息与日志,并提供查询「这个订单卡在哪一步」的能力。

分布式流程的调试体验取决于两件事:能否按业务 ID 找到全链路,以及能否看到当前处于哪一步。前者靠相关性 ID 与分布式追踪,后者靠 Saga 状态的持久化查询接口。

// 相关性 ID 贯穿消息与日志:从请求头一路透传到每个服务
public async Task InvokeAsync(HttpContext ctx)
{
    var correlationId = ctx.Request.Headers["X-Correlation-Id"].FirstOrDefault()
                        ?? Guid.NewGuid().ToString("N");
    ctx.Response.Headers["X-Correlation-Id"] = correlationId;

    using (LogContext.PushProperty("CorrelationId", correlationId))
    using (Activity.Current?.SetTag("correlation.id", correlationId))
    {
        await next(ctx);
    }
}
// 运维查询:直接回答「这个订单的 Saga 走到哪一步」
var state = await db.OrderSagaStates
    .Where(s => s.OrderId == orderId)
    .Select(s => new { s.CurrentState, s.UpdatedAt, s.CompensationFailures })
    .FirstOrDefaultAsync(ct);
观测维度手段回答的问题
链路相关性 ID + 分布式追踪这个请求经过了哪些服务
状态Saga 状态表查询当前卡在哪一步
耗时每步耗时直方图哪一步慢
失败补偿次数与原因哪一步常失败
积压Outbox 未投递数量消息投递是否正常
兜底超时扫描任务是否有实例长期停滞

避坑: 相关性 ID 必须跨消息边界传递,只在 HTTP 头里传是不够的——发消息时要把它写进消息头,消费端读出来再放进日志上下文,否则链路在第一个异步跳转处就断了。另一个坑是缺少「卡住实例」的告警:没有告警,Saga 停滞只能靠用户投诉发现,而这时业务损失已经发生。

8. 总结

环节要点
现实选择放弃 2PC,用最终一致 + 补偿换可用性
Saga 模型正向步骤配补偿动作,补偿是语义抵消不是回滚
编排与协同步骤多选编排易追踪,步骤少选协同低耦合
幂等至少一次投递下重复是常态,业务去重键 + 唯一索引
Outbox业务与消息同事务提交,后台投递,消费端幂等
状态机状态持久化,事件驱动推进,显式处理超时
可观测相关性 ID 贯穿链路,Saga 状态可查询,卡住要告警

分布式事务没有银弹,Saga 的全部价值在于把「不可能做到的事」(跨服务原子提交)换成「可以做到的事」(可补偿、可重试、可观测的最终一致)。代价是系统里会长期存在中间状态,业务必须接受并处理它们。真正决定成败的不是 Saga 框架,而是幂等是否做扎实、Outbox 是否可靠、超时是否有兜底——这三件事做到了,最终一致就是可靠的;漏掉任何一件,它就会在某个深夜变成对账事故。下一篇我们回到单机视角,讨论当内存持续增长、GC 压力居高不下时,如何用 dump 分析定位真正的根因。

延伸阅读

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「csharp」更多文章

  1. .NET 机器学习实战
  2. 内存剖析与 dump 分析
  3. GraphQL 服务端开发