Skip to content

分布式系统局部失败 ​

概述 ​

在分布式系统中,局部失败(Partial Failure)是常态而非异常。系统的某个组件可能失败,而其他部分仍然正常运行。这种特性使得幂等性设计变得至关重要。

什么是局部失败 ​

典型场景 ​

微服务架构:
┌─────────────┐     ┌──────────────┐     ┌─────────────┐
│  API Gateway │────>│ Order Service │────>│ Payment Svc │
└─────────────┘     └──────────────┘     └─────────────┘
                           │                      │
                      ✓ 正常                 ❌ 失败
                           
用户看到: "订单创建成功,但支付失败"
实际问题: 订单已创建,支付超时,需要重试

局部失败的原因 ​

类型示例概率
网络分区服务间通信中断中
节点故障单个服务器宕机低
资源耗尽数据库连接池满中
依赖失败第三方支付接口超时高
软件 Bug内存泄漏、死锁低

CAP 定理与幂等性 ​

CAP 权衡 ​

一致性 (Consistency)    可用性 (Availability)
         \                    /
          \                  /
           \                /
            分区容忍 (Partition Tolerance)

在分布式系统中:

  • P(分区容忍) 是必须的
  • 选择 C(强一致):可能降低可用性
  • 选择 A(高可用):可能出现不一致,需要幂等性修复

最终一致性 + 幂等性 ​

csharp
// 最终一致性方案
public class OrderService
{
    public async Task<Order> CreateOrderAsync(CreateOrderRequest request)
    {
        // 1. 创建订单(本地事务)
        var order = await _dbContext.Orders.AddAsync(new Order 
        { 
            UserId = request.UserId,
            Status = "pending_payment"
        });
        await _dbContext.SaveChangesAsync();
        
        // 2. 调用支付服务(可能失败)
        try
        {
            await _paymentClient.ChargeAsync(order.Id, request.Amount);
            order.Status = "paid";
        }
        catch (Exception ex)
        {
            // 支付失败,订单仍然是 pending
            // 通过定时任务或消息队列重试
            _logger.LogWarning(ex, "Payment failed for order {OrderId}", order.Id);
            
            // 发送延迟消息,稍后重试
            await _messageQueue.PublishAsync(new RetryPaymentEvent 
            { 
                OrderId = order.Id,
                RetryAfter = DateTime.UtcNow.AddMinutes(5)
            });
        }
        
        await _dbContext.SaveChangesAsync();
        
        return order;
    }
}

Saga 模式与幂等性 ​

什么是 Saga ​

Saga 是一种处理分布式事务的模式,将长事务拆分为多个本地事务,每个步骤都有补偿操作。

订单流程 Saga:
1. 创建订单 → 补偿:取消订单
2. 扣减库存 → 补偿:恢复库存  
3. 创建支付 → 补偿:取消支付
4. 发送通知 → 补偿:发送取消通知

C# 实现 ​

csharp
public class OrderSaga
{
    private readonly IOrderRepository _orders;
    private readonly IInventoryService _inventory;
    private readonly IPaymentService _payments;
    
    public async Task ExecuteAsync(OrderSagaData sagaData)
    {
        try
        {
            // Step 1: 创建订单(幂等)
            var order = await _orders.GetOrCreateAsync(sagaData.OrderId, () => 
                new Order 
                { 
                    Id = sagaData.OrderId,
                    Status = "created"
                });
            
            // Step 2: 扣减库存(幂等)
            await _inventory.ReserveStockAsync(
                sagaData.ProductId, 
                sagaData.Quantity,
                sagaData.ReservationId); // 幂等键
            
            // Step 3: 创建支付(幂等)
            var payment = await _payments.CreatePaymentAsync(new PaymentRequest
            {
                OrderId = sagaData.OrderId,
                Amount = sagaData.Amount,
                IdempotencyKey = sagaData.PaymentIdempotencyKey
            });
            
            // Step 4: 确认订单
            order.Status = "confirmed";
            await _orders.UpdateAsync(order);
            
            // Step 5: 发送通知
            await _notificationService.SendAsync(order.UserId, "Order confirmed");
        }
        catch (Exception ex)
        {
            _logger.LogError(ex, "Saga failed, starting compensation");
            
            // 执行补偿操作(也需要幂等)
            await CompensateAsync(sagaData);
            
            throw;
        }
    }
    
    private async Task CompensateAsync(OrderSagaData sagaData)
    {
        // 补偿步骤(反向执行,幂等)
        
        // 1. 取消支付
        await _payments.CancelPaymentAsync(sagaData.OrderId);
        
        // 2. 恢复库存
        await _inventory.ReleaseStockAsync(
            sagaData.ProductId, 
            sagaData.Quantity,
            sagaData.ReservationId);
        
        // 3. 取消订单
        var order = await _orders.GetAsync(sagaData.OrderId);
        if (order != null && order.Status != "cancelled")
        {
            order.Status = "cancelled";
            await _orders.UpdateAsync(order);
        }
        
        // 4. 发送取消通知
        await _notificationService.SendAsync(sagaData.UserId, "Order cancelled");
    }
}

public class OrderSagaData
{
    public Guid OrderId { get; set; }
    public Guid ProductId { get; set; }
    public int Quantity { get; set; }
    public decimal Amount { get; set; }
    public Guid UserId { get; set; }
    public string ReservationId { get; set; } // 库存预留 ID
    public string PaymentIdempotencyKey { get; set; }
}

重试策略 ​

1. 指数退避重试 ​

csharp
public class ResilientDistributedCall
{
    public async Task<T> ExecuteWithRetry<T>(
        Func<Task<T>> operation,
        string operationName,
        int maxRetries = 5)
    {
        for (int attempt = 1; attempt <= maxRetries; attempt++)
        {
            try
            {
                return await operation();
            }
            catch (Exception ex) when (IsTransient(ex) && attempt < maxRetries)
            {
                var delay = TimeSpan.FromSeconds(Math.Pow(2, attempt));
                
                _logger.LogWarning(ex, 
                    "{Operation} failed (attempt {Attempt}/{Max}), retrying in {Delay}s",
                    operationName, attempt, maxRetries, delay.TotalSeconds);
                
                await Task.Delay(delay);
            }
        }
        
        throw new Exception($"{operationName} failed after {maxRetries} retries");
    }
    
    private bool IsTransient(Exception ex)
    {
        return ex is HttpRequestException httpEx &&
               (httpEx.StatusCode >= 500 || httpEx.StatusCode == 408 || httpEx.StatusCode == 429);
    }
}

2. 断路器模式 ​

csharp
using Polly;

public class CircuitBreakerService
{
    private readonly AsyncPolicy _circuitBreaker;
    
    public CircuitBreakerService()
    {
        _circuitBreaker = Policy
            .Handle<HttpRequestException>()
            .CircuitBreakerAsync(
                exceptionsAllowedBeforeBreaking: 5,
                durationOfBreak: TimeSpan.FromMinutes(1),
                onBreak: (outcome, breakDelay) =>
                {
                    _logger.LogWarning("Circuit broken for {Service}", breakDelay);
                },
                onReset: () =>
                {
                    _logger.LogInformation("Circuit reset");
                });
    }
    
    public async Task<T> ExecuteAsync<T>(Func<Task<T>> operation)
    {
        return await _circuitBreaker.ExecuteAsync(operation);
    }
}

事件溯源与幂等性 ​

事件溯源模式 ​

csharp
public class EventSourcedOrder
{
    private readonly List<OrderEvent> _events = new();
    
    public Guid Id { get; private set; }
    public string Status { get; private set; }
    
    // 应用事件(幂等)
    public void ApplyEvent(OrderEvent @event)
    {
        switch (@event)
        {
            case OrderCreatedEvent e:
                Id = e.OrderId;
                Status = "created";
                break;
            
            case PaymentCompletedEvent e:
                Status = "paid";
                break;
            
            case OrderCancelledEvent e:
                Status = "cancelled";
                break;
        }
        
        _events.Add(@event);
    }
    
    // 从事件重建状态
    public static EventSourcedOrder FromEvents(IEnumerable<OrderEvent> events)
    {
        var order = new EventSourcedOrder();
        foreach (var evt in events)
        {
            order.ApplyEvent(evt);
        }
        return order;
    }
}

// 事件存储(幂等追加)
public class EventStore
{
    public async Task AppendEventsAsync(Guid aggregateId, IEnumerable<OrderEvent> events)
    {
        await using var transaction = await _dbContext.Database.BeginTransactionAsync();
        
        foreach (var evt in events)
        {
            var record = new EventRecord
            {
                AggregateId = aggregateId,
                EventType = evt.GetType().Name,
                Data = JsonSerializer.Serialize(evt),
                OccurredAt = DateTime.UtcNow
            };
            
            _dbContext.Events.Add(record);
        }
        
        await _dbContext.SaveChangesAsync();
        await transaction.CommitAsync();
    }
}

监控与可观测性 ​

1. 分布式追踪 ​

csharp
using OpenTelemetry.Trace;

public class TracedDistributedOperation
{
    private readonly TracerProvider _tracerProvider;
    
    public async Task ProcessOrderAsync(Guid orderId)
    {
        using var activity = _tracerProvider.StartActivity("ProcessOrder");
        activity?.SetTag("order.id", orderId.ToString());
        
        try
        {
            // 调用其他服务时传播 trace context
            await _inventoryService.ReserveStockAsync(orderId);
            await _paymentService.ChargeAsync(orderId);
            
            activity?.SetStatus(ActivityStatusCode.Ok);
        }
        catch (Exception ex)
        {
            activity?.SetStatus(ActivityStatusCode.Error, ex.Message);
            activity?.RecordException(ex);
            throw;
        }
    }
}

2. 健康检查 ​

csharp
public class DistributedSystemHealthCheck : IHealthCheck
{
    private readonly IInventoryService _inventory;
    private readonly IPaymentService _payment;
    
    public async Task<HealthCheckResult> CheckHealthAsync(
        HealthCheckContext context,
        CancellationToken cancellationToken = default)
    {
        var checks = new Dictionary<string, HealthStatus>();
        
        // 检查依赖服务
        try
        {
            await _inventory.PingAsync(cancellationToken);
            checks["inventory"] = HealthStatus.Healthy;
        }
        catch
        {
            checks["inventory"] = HealthStatus.Unhealthy;
        }
        
        try
        {
            await _payment.PingAsync(cancellationToken);
            checks["payment"] = HealthStatus.Healthy;
        }
        catch
        {
            checks["payment"] = HealthStatus.Unhealthy;
        }
        
        var overallStatus = checks.Any(v => v.Value == HealthStatus.Unhealthy)
            ? HealthStatus.Degraded
            : HealthStatus.Healthy;
        
        return new HealthCheckResult(
            overallStatus,
            data: checks);
    }
}

最佳实践总结 ​

✅ DO ​

  1. 假设失败会发生:设计时考虑所有可能的失败场景
  2. 实现幂等操作:所有重试都必须是安全的
  3. 使用 Saga 模式:处理跨服务的长事务
  4. 实施断路器:防止级联失败
  5. 添加监控和追踪:快速定位问题
  6. 设计补偿机制:能够回滚部分成功的操作

❌ DON'T ​

  1. 不要假设网络可靠:总是处理超时和失败
  2. 不要忽略部分成功:记录并处理中间状态
  3. 不要同步等待所有服务:使用异步和超时
  4. 不要忘记清理资源:释放锁、关闭连接

总结 ​

分布式系统的局部失败是不可避免的,关键策略包括:

  1. 幂等性:保证重试安全
  2. 最终一致性:接受短暂的不一致
  3. Saga 模式:管理分布式事务
  4. 重试和断路器:提高容错能力
  5. 监控和追踪:快速发现和解决问题

通过这些技术,可以构建出高度可靠的分布式系统。

Released under the MIT License.