Skip to content

一致性 - 分布式事务幂等 ​

概述 ​

在分布式系统中,一个业务操作可能涉及多个服务的数据修改。如何保证这些跨服务的操作要么全部成功,要么全部失败,并且在重试时保持幂等性,是分布式事务的核心挑战。

典型场景 ​

场景 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 ​

  1. 每个步骤幂等:所有操作都支持重试
  2. 记录状态:持久化 Saga 状态和步骤历史
  3. 实现补偿:每个步骤都有对应的补偿操作
  4. 监控告警:跟踪 Saga 执行情况和补偿次数
  5. 超时控制:设置合理的超时时间

❌ DON'T ​

  1. 不要长时间持有锁:尽快释放资源
  2. 不要忽略补偿失败:记录并告警
  3. 不要假设网络可靠:处理所有可能的失败
  4. 不要忘记清理:定期清理完成的 Saga 记录

总结 ​

分布式事务的幂等性是构建可靠分布式系统的关键:

✅ Saga 模式:灵活、实用、易于实现
✅ TCC 模式:强一致性、适合金融场景
✅ 本地消息表:最终一致性、解耦服务
✅ 幂等保证:每个步骤都支持重试

通过合理选择模式和完善的幂等性设计,可以构建出高可用的分布式系统。

Released under the MIT License.