分布式系统局部失败
概述
在分布式系统中,局部失败(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
- 假设失败会发生:设计时考虑所有可能的失败场景
- 实现幂等操作:所有重试都必须是安全的
- 使用 Saga 模式:处理跨服务的长事务
- 实施断路器:防止级联失败
- 添加监控和追踪:快速定位问题
- 设计补偿机制:能够回滚部分成功的操作
❌ DON'T
- 不要假设网络可靠:总是处理超时和失败
- 不要忽略部分成功:记录并处理中间状态
- 不要同步等待所有服务:使用异步和超时
- 不要忘记清理资源:释放锁、关闭连接
总结
分布式系统的局部失败是不可避免的,关键策略包括:
- 幂等性:保证重试安全
- 最终一致性:接受短暂的不一致
- Saga 模式:管理分布式事务
- 重试和断路器:提高容错能力
- 监控和追踪:快速发现和解决问题
通过这些技术,可以构建出高度可靠的分布式系统。