一致性 - 分布式事务幂等
概述
在分布式系统中,一个业务操作可能涉及多个服务的数据修改。如何保证这些跨服务的操作要么全部成功,要么全部失败,并且在重试时保持幂等性,是分布式事务的核心挑战。
典型场景
场景 1:订单创建流程
订单服务 库存服务 支付服务 通知服务
| | | |
|-- 创建订单 ---------->| | |
| |-- 扣减库存 --------->| |
| | |-- 创建支付 -------->|
| | | |-- 发送通知
| | | |
问题:如果任何一步失败,如何回滚?重试时如何保证幂等?解决方案对比
1. Saga 模式(推荐)
Saga 将长事务拆分为多个本地事务,每个步骤都有补偿操作。
csharp
public class OrderSaga : ISaga
{
private readonly IOrderService _orderService;
private readonly IInventoryService _inventoryService;
private readonly IPaymentService _paymentService;
private readonly INotificationService _notificationService;
public async Task ExecuteAsync(OrderSagaData data)
{
var compensations = new Stack<Func<Task>>();
try
{
// Step 1: 创建订单(幂等)
var order = await _orderService.CreateOrderAsync(data.OrderRequest);
compensations.Push(async () => await _orderService.CancelOrderAsync(order.Id));
// Step 2: 扣减库存(幂等)
await _inventoryService.ReserveStockAsync(
data.ProductId,
data.Quantity,
data.ReservationId); // 幂等键
compensations.Push(async () => await _inventoryService.ReleaseStockAsync(
data.ProductId,
data.Quantity,
data.ReservationId));
// Step 3: 创建支付(幂等)
var payment = await _paymentService.CreatePaymentAsync(new PaymentRequest
{
OrderId = order.Id,
Amount = order.TotalAmount,
IdempotencyKey = data.PaymentIdempotencyKey
});
compensations.Push(async () => await _paymentService.CancelPaymentAsync(payment.Id));
// Step 4: 发送通知(幂等)
await _notificationService.SendOrderConfirmationAsync(order.Id);
// 通知不需要补偿
// 所有步骤成功,提交 Saga
await MarkSagaCompletedAsync(data.SagaId);
}
catch (Exception ex)
{
// 执行补偿操作(反向执行,幂等)
await ExecuteCompensationsAsync(compensations);
await MarkSagaFailedAsync(data.SagaId, ex.Message);
throw;
}
}
private async Task ExecuteCompensationsAsync(Stack<Func<Task>> compensations)
{
foreach (var compensation in compensations)
{
try
{
await compensation();
}
catch (Exception ex)
{
// 记录补偿失败,需要人工介入
_logger.LogError(ex, "Compensation failed");
}
}
}
}
public class OrderSagaData
{
public Guid SagaId { get; set; }
public CreateOrderRequest OrderRequest { get; set; }
public Guid ProductId { get; set; }
public int Quantity { get; set; }
public string ReservationId { get; set; }
public string PaymentIdempotencyKey { get; set; }
}2. TCC 模式(Try-Confirm-Cancel)
TCC 将操作分为三个阶段:尝试、确认、取消。
csharp
public interface ITccResource
{
Task<bool> TryAsync(TccContext context);
Task ConfirmAsync(TccContext context);
Task CancelAsync(TccContext context);
}
public class InventoryTccResource : ITccResource
{
public async Task<bool> TryAsync(TccContext context)
{
var request = context.GetData<ReserveStockRequest>();
// 预占库存(不实际扣减)
var reserved = await _dbContext.Database.ExecuteSqlRawAsync(
"UPDATE products SET reserved_stock = reserved_stock + @quantity WHERE id = @productId AND stock - reserved_stock >= @quantity",
new NpgsqlParameter("@quantity", request.Quantity),
new NpgsqlParameter("@productId", request.ProductId));
return reserved > 0;
}
public async Task ConfirmAsync(TccContext context)
{
var request = context.GetData<ReserveStockRequest>();
// 确认扣减:将预占转为实际扣减
await _dbContext.Database.ExecuteSqlRawAsync(
"UPDATE products SET stock = stock - @quantity, reserved_stock = reserved_stock - @quantity WHERE id = @productId",
new NpgsqlParameter("@quantity", request.Quantity),
new NpgsqlParameter("@productId", request.ProductId));
}
public async Task CancelAsync(TccContext context)
{
var request = context.GetData<ReserveStockRequest>();
// 取消预占:释放预留的库存
await _dbContext.Database.ExecuteSqlRawAsync(
"UPDATE products SET reserved_stock = reserved_stock - @quantity WHERE id = @productId",
new NpgsqlParameter("@quantity", request.Quantity),
new NpgsqlParameter("@productId", request.ProductId));
}
}
public class TccCoordinator
{
private readonly List<ITccResource> _resources;
public async Task ExecuteAsync(TccContext context)
{
// Phase 1: Try
var tryResults = new List<bool>();
foreach (var resource in _resources)
{
var result = await resource.TryAsync(context);
tryResults.Add(result);
if (!result)
{
// Try 失败,执行 Cancel
await CancelAllAsync(context);
throw new TccException("Try phase failed");
}
}
// Phase 2: Confirm or Cancel
try
{
// 所有 Try 成功,执行 Confirm
foreach (var resource in _resources)
{
await resource.ConfirmAsync(context);
}
}
catch (Exception)
{
// Confirm 失败,执行 Cancel
await CancelAllAsync(context);
throw;
}
}
private async Task CancelAllAsync(TccContext context)
{
foreach (var resource in _resources)
{
try
{
await resource.CancelAsync(context);
}
catch (Exception ex)
{
_logger.LogError(ex, "Cancel failed");
}
}
}
}3. 本地消息表(最终一致性)
通过本地消息表保证消息必达,实现最终一致性。
sql
-- 本地消息表
CREATE TABLE outbox_messages (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
aggregate_id UUID NOT NULL,
aggregate_type VARCHAR(50) NOT NULL,
event_type VARCHAR(50) NOT NULL,
payload JSONB NOT NULL,
status VARCHAR(20) NOT NULL DEFAULT 'pending', -- pending, sent, confirmed
retry_count INTEGER NOT NULL DEFAULT 0,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
sent_at TIMESTAMP WITH TIME ZONE
);
CREATE INDEX idx_outbox_messages_status ON outbox_messages(status);
CREATE INDEX idx_outbox_messages_created_at ON outbox_messages(created_at);csharp
public class OutboxService
{
private readonly AppDbContext _dbContext;
private readonly IMessagePublisher _messagePublisher;
/// <summary>
/// 在事务中保存业务数据和消息
/// </summary>
public async Task SaveWithOutboxAsync<TAggregate>(
TAggregate aggregate,
DomainEvent domainEvent) where TAggregate : class
{
await using var transaction = await _dbContext.Database.BeginTransactionAsync();
try
{
// 1. 保存业务数据
_dbContext.Set<TAggregate>().Update(aggregate);
await _dbContext.SaveChangesAsync();
// 2. 保存出箱消息(同一事务)
var message = new OutboxMessage
{
Id = Guid.NewGuid(),
AggregateId = GetAggregateId(aggregate),
AggregateType = typeof(TAggregate).Name,
EventType = domainEvent.GetType().Name,
Payload = JsonSerializer.Serialize(domainEvent),
Status = "pending",
CreatedAt = DateTime.UtcNow
};
_dbContext.OutboxMessages.Add(message);
await _dbContext.SaveChangesAsync();
await transaction.CommitAsync();
}
catch
{
await transaction.RollbackAsync();
throw;
}
}
/// <summary>
/// 后台任务:处理待发送的消息
/// </summary>
public async Task ProcessPendingMessagesAsync(int batchSize = 100)
{
var messages = await _dbContext.OutboxMessages
.Where(m => m.Status == "pending" && m.RetryCount < 5)
.OrderBy(m => m.CreatedAt)
.Take(batchSize)
.ToListAsync();
foreach (var message in messages)
{
try
{
// 发布消息
await _messagePublisher.PublishAsync(
message.EventType,
message.Payload);
// 标记为已发送
message.Status = "sent";
message.SentAt = DateTime.UtcNow;
message.UpdatedAt = DateTime.UtcNow;
}
catch (Exception ex)
{
// 增加重试次数
message.RetryCount++;
message.UpdatedAt = DateTime.UtcNow;
_logger.LogError(ex, "Failed to publish message {MessageId}", message.Id);
}
}
await _dbContext.SaveChangesAsync();
}
}幂等性保证
1. 每个步骤的幂等性
csharp
public class IdempotentOrderService
{
public async Task<Order> CreateOrderAsync(CreateOrderRequest request, string idempotencyKey)
{
// 检查是否已存在
var existing = await _dbContext.Orders
.FirstOrDefaultAsync(o => o.IdempotencyKey == idempotencyKey);
if (existing != null)
{
return existing; // 幂等返回
}
// 创建新订单
var order = new Order
{
Id = Guid.NewGuid(),
IdempotencyKey = idempotencyKey,
// ...
};
_dbContext.Orders.Add(order);
await _dbContext.SaveChangesAsync();
return order;
}
}2. Saga 状态的幂等性
sql
-- Saga 状态表
CREATE TABLE saga_instances (
saga_id UUID PRIMARY KEY,
saga_type VARCHAR(50) NOT NULL,
status VARCHAR(20) NOT NULL, -- running, completed, failed, compensating
current_step INTEGER NOT NULL DEFAULT 0,
saga_data JSONB NOT NULL,
created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(),
updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
-- Saga 步骤历史
CREATE TABLE saga_step_history (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
saga_id UUID NOT NULL REFERENCES saga_instances(saga_id),
step_name VARCHAR(50) NOT NULL,
status VARCHAR(20) NOT NULL, -- pending, completed, failed, compensated
executed_at TIMESTAMP WITH TIME ZONE DEFAULT NOW()
);
CREATE UNIQUE INDEX idx_saga_step_unique ON saga_step_history(saga_id, step_name);csharp
public class SagaStateManager
{
public async Task<bool> IsStepCompletedAsync(Guid sagaId, string stepName)
{
return await _dbContext.SagaStepHistory
.AnyAsync(s => s.SagaId == sagaId &&
s.StepName == stepName &&
s.Status == "completed");
}
public async Task MarkStepCompletedAsync(Guid sagaId, string stepName)
{
// 使用唯一索引保证幂等
var history = new SagaStepHistory
{
Id = Guid.NewGuid(),
SagaId = sagaId,
StepName = stepName,
Status = "completed",
ExecutedAt = DateTime.UtcNow
};
try
{
_dbContext.SagaStepHistory.Add(history);
await _dbContext.SaveChangesAsync();
}
catch (DbUpdateException ex) when (IsUniqueViolation(ex))
{
// 步骤已标记为完成,忽略
}
}
}监控与告警
指标收集
csharp
public class DistributedTransactionMetrics
{
private readonly Counter<long> _sagasStarted;
private readonly Counter<long> _sagasCompleted;
private readonly Counter<long> _sagasFailed;
private readonly Counter<long> _compensationsExecuted;
private readonly Histogram<double> _sagaDuration;
public void RecordSagaStarted(string sagaType)
{
_sagasStarted.Add(1, new KeyValuePair<string, object?>("type", sagaType));
}
public void RecordSagaCompleted(string sagaType, double durationMs)
{
_sagasCompleted.Add(1, new KeyValuePair<string, object?>("type", sagaType));
_sagaDuration.Record(durationMs);
}
public void RecordSagaFailed(string sagaType, string reason)
{
_sagasFailed.Add(1,
new KeyValuePair<string, object?>("type", sagaType),
new KeyValuePair<string, object?>("reason", reason));
}
public void RecordCompensationExecuted(string sagaType, string step)
{
_compensationsExecuted.Add(1,
new KeyValuePair<string, object?>("type", sagaType),
new KeyValuePair<string, object?>("step", step));
}
}最佳实践总结
✅ DO
- 每个步骤幂等:所有操作都支持重试
- 记录状态:持久化 Saga 状态和步骤历史
- 实现补偿:每个步骤都有对应的补偿操作
- 监控告警:跟踪 Saga 执行情况和补偿次数
- 超时控制:设置合理的超时时间
❌ DON'T
- 不要长时间持有锁:尽快释放资源
- 不要忽略补偿失败:记录并告警
- 不要假设网络可靠:处理所有可能的失败
- 不要忘记清理:定期清理完成的 Saga 记录
总结
分布式事务的幂等性是构建可靠分布式系统的关键:
✅ Saga 模式:灵活、实用、易于实现
✅ TCC 模式:强一致性、适合金融场景
✅ 本地消息表:最终一致性、解耦服务
✅ 幂等保证:每个步骤都支持重试
通过合理选择模式和完善的幂等性设计,可以构建出高可用的分布式系统。