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();
}));
| 状态 | 触发事件 | 下一状态 | 动作 |
|---|---|---|---|
| Initial | OrderPlaced | StockReserved | 预占库存 |
| StockReserved | StockReserved | PaymentCharged | 扣款 |
| StockReserved | StockFailed | Compensating | 取消订单 |
| PaymentCharged | PaymentFailed | Compensating | 释放库存 + 取消订单 |
| Compensating | Compensated | Cancelled | 通知用户 |
| PaymentCharged | Completed | Completed | 完成订单 |
避坑: 状态机的超时必须显式处理。如果库存服务挂了,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 分析定位真正的根因。
延伸阅读
继续阅读
探索更多技术文章
浏览归档,发现更多关于系统设计、工具链和工程实践的内容。