Skip to content

边界 - 极端并发穿透 ​

概述 ​

在极端高并发场景下,即使实现了完善的幂等性机制,仍可能出现"穿透"问题:

  1. 缓存穿透:大量请求同时绕过缓存
  2. 锁竞争:分布式锁成为瓶颈
  3. 数据库压力:唯一索引检查导致性能下降
  4. 雪崩效应:系统级联失败

典型场景 ​

场景 1:秒杀活动 ​

时间 T0: 10000 个用户同时点击"抢购"按钮
         ↓
时间 T1: 10000 个请求同时到达网关
         ↓
时间 T2: 10000 个请求同时检查 Token
         ↓
时间 T3: Redis 负载飙升,响应变慢
         ↓
时间 T4: 部分请求超时,重试加剧问题
         ↓
结果: 系统崩溃或响应极慢

场景 2:热点商品库存扣减 ​

商品 A 只剩 1 件库存
100 个用户同时下单
   ↓
所有请求都通过 Token 验证
   ↓
所有请求都尝试扣减库存
   ↓
数据库行锁竞争激烈
   ↓
大量请求超时或失败

解决方案 ​

1. 令牌桶限流(网关层) ​

csharp
public class TokenBucketRateLimiter
{
    private readonly double _capacity;
    private readonly double _refillRate;
    private double _tokens;
    private DateTime _lastRefill;
    private readonly object _lock = new();
    
    public TokenBucketRateLimiter(double capacity, double refillRate)
    {
        _capacity = capacity;
        _refillRate = refillRate;
        _tokens = capacity;
        _lastRefill = DateTime.UtcNow;
    }
    
    public bool AllowRequest()
    {
        lock (_lock)
        {
            Refill();
            
            if (_tokens >= 1)
            {
                _tokens -= 1;
                return true;
            }
            
            return false;
        }
    }
    
    private void Refill()
    {
        var now = DateTime.UtcNow;
        var elapsed = (now - _lastRefill).TotalSeconds;
        
        _tokens = Math.Min(_capacity, _tokens + elapsed * _refillRate);
        _lastRefill = now;
    }
}

// 使用:每秒最多 100 个请求,桶容量 200
var limiter = new TokenBucketRateLimiter(200, 100);

2. 排队机制 ​

csharp
public class RequestQueue
{
    private readonly Channel<RequestContext> _queue;
    private readonly int _maxConcurrency;
    private readonly SemaphoreSlim _semaphore;
    
    public RequestQueue(int maxConcurrency)
    {
        _maxConcurrency = maxConcurrency;
        _semaphore = new SemaphoreSlim(maxConcurrency);
        _queue = Channel.CreateBounded<RequestContext>(new BoundedChannelOptions(10000)
        {
            FullMode = BoundedChannelFullMode.DropOldest
        });
    }
    
    public async Task StartProcessingAsync(CancellationToken cancellationToken)
    {
        await foreach (var context in _queue.Reader.ReadAllAsync(cancellationToken))
        {
            // 等待可用槽位
            await _semaphore.WaitAsync(cancellationToken);
            
            _ = Task.Run(async () =>
            {
                try
                {
                    await ProcessRequestAsync(context);
                }
                finally
                {
                    _semaphore.Release();
                }
            }, cancellationToken);
        }
    }
    
    public async Task<bool> EnqueueAsync(RequestContext context)
    {
        return await _queue.Writer.WriteAsync(context);
    }
    
    private async Task ProcessRequestAsync(RequestContext context)
    {
        // 处理请求逻辑
        await context.Process();
    }
}

3. 库存预扣减(Redis) ​

csharp
public class InventoryReservationService
{
    private readonly IDatabase _redis;
    
    /// <summary>
    /// 预扣减库存(原子操作)
    /// </summary>
    public async Task<Result<string>> ReserveStockAsync(
        Guid productId,
        int quantity,
        string userId)
    {
        var stockKey = $"stock:{productId}";
        var reservedKey = $"reserved:{productId}:{userId}";
        
        // Lua 脚本:原子性地检查和扣减库存
        const string script = @"
            local stock_key = KEYS[1]
            local reserved_key = KEYS[2]
            local quantity = tonumber(ARGV[1])
            local user_id = ARGV[2]
            local ttl = tonumber(ARGV[3])
            
            -- 检查是否已预留
            local existing = redis.call('GET', reserved_key)
            if existing then
                return -2  -- 已预留
            end
            
            -- 检查库存
            local stock = tonumber(redis.call('GET', stock_key) or '0')
            if stock < quantity then
                return -1  -- 库存不足
            end
            
            -- 扣减库存
            redis.call('DECRBY', stock_key, quantity)
            
            -- 设置预留标记
            redis.call('SET', reserved_key, quantity, 'EX', ttl)
            
            -- 返回剩余库存
            return stock - quantity
        """;
        
        var result = await _redis.ScriptEvaluateAsync(
            script,
            new RedisKey[] { stockKey, reservedKey },
            new RedisValue[] { quantity, userId, 300 }); // 5 分钟过期
        
        var remainingStock = (int)result;
        
        if (remainingStock == -2)
        {
            return Result<string>.Failure("Already reserved");
        }
        
        if (remainingStock == -1)
        {
            return Result<string>.Failure("Out of stock");
        }
        
        var reservationId = Guid.NewGuid().ToString("N");
        return Result<string>.Success(reservationId);
    }
    
    /// <summary>
    /// 确认订单(将预留转为实际扣减)
    /// </summary>
    public async Task ConfirmReservationAsync(Guid productId, string userId)
    {
        var reservedKey = $"reserved:{productId}:{userId}";
        
        // 删除预留标记(库存已在预留时扣减)
        await _redis.KeyDeleteAsync(reservedKey);
    }
    
    /// <summary>
    /// 取消预留(恢复库存)
    /// </summary>
    public async Task CancelReservationAsync(Guid productId, string userId)
    {
        var stockKey = $"stock:{productId}";
        var reservedKey = $"reserved:{productId}:{userId}";
        
        const string script = @"
            local stock_key = KEYS[1]
            local reserved_key = KEYS[2]
            
            local quantity = tonumber(redis.call('GET', reserved_key))
            if quantity then
                redis.call('DEL', reserved_key)
                redis.call('INCRBY', stock_key, quantity)
                return 1
            end
            
            return 0
        """;
        
        await _redis.ScriptEvaluateAsync(
            script,
            new RedisKey[] { stockKey, reservedKey });
    }
}

4. 异步化处理 ​

csharp
public class AsyncOrderProcessor
{
    private readonly Channel<OrderRequest> _orderChannel;
    private readonly ILogger<AsyncOrderProcessor> _logger;
    
    public AsyncOrderProcessor()
    {
        // 有界通道,背压控制
        _orderChannel = Channel.CreateBounded<OrderRequest>(new BoundedChannelOptions(50000)
        {
            FullMode = BoundedChannelFullMode.Wait
        });
    }
    
    /// <summary>
    /// 快速接受订单请求
    /// </summary>
    public async Task<string> SubmitOrderAsync(OrderRequest request)
    {
        var orderId = Guid.NewGuid().ToString("N");
        
        // 快速入队(不阻塞)
        await _orderChannel.Writer.WriteAsync(new OrderRequest
        {
            Id = orderId,
            Data = request,
            SubmittedAt = DateTime.UtcNow
        });
        
        // 立即返回订单 ID
        return orderId;
    }
    
    /// <summary>
    /// 后台处理订单
    /// </summary>
    public async Task StartProcessingAsync(int concurrency = 10)
    {
        var tasks = Enumerable.Range(0, concurrency)
            .Select(_ => ProcessOrdersAsync())
            .ToList();
        
        await Task.WhenAll(tasks);
    }
    
    private async Task ProcessOrdersAsync()
    {
        await foreach (var order in _orderChannel.Reader.ReadAllAsync())
        {
            try
            {
                await ProcessSingleOrderAsync(order);
            }
            catch (Exception ex)
            {
                _logger.LogError(ex, "Failed to process order {OrderId}", order.Id);
            }
        }
    }
    
    private async Task ProcessSingleOrderAsync(OrderRequest order)
    {
        // 模拟处理延迟
        await Task.Delay(Random.Shared.Next(100, 500));
        
        // 实际订单处理逻辑
        // ...
    }
}

5. 降级策略 ​

csharp
public class CircuitBreakerService
{
    private readonly IOrderService _orderService;
    private readonly IFallbackService _fallbackService;
    private CircuitState _state = CircuitState.Closed;
    private int _failureCount = 0;
    private DateTime _lastFailureTime;
    
    public async Task<Result<Order>> CreateOrderAsync(OrderRequest request)
    {
        // 检查断路器状态
        if (_state == CircuitState.Open)
        {
            if (DateTime.UtcNow - _lastFailureTime > TimeSpan.FromMinutes(1))
            {
                _state = CircuitState.HalfOpen;
            }
            else
            {
                // 快速失败,使用降级方案
                return await _fallbackService.QueueOrderAsync(request);
            }
        }
        
        try
        {
            var result = await _orderService.CreateOrderAsync(request);
            
            if (_state == CircuitState.HalfOpen)
            {
                _state = CircuitState.Closed;
                _failureCount = 0;
            }
            
            return result;
        }
        catch (Exception ex)
        {
            _failureCount++;
            _lastFailureTime = DateTime.UtcNow;
            
            if (_failureCount > 10)
            {
                _state = CircuitState.Open;
            }
            
            // 降级:队列化订单
            return await _fallbackService.QueueOrderAsync(request);
        }
    }
}

public enum CircuitState
{
    Closed,     // 正常
    Open,       // 断开(快速失败)
    HalfOpen    // 半开(试探)
}

压力测试 ​

k6 极限测试脚本 ​

javascript
import http from 'k6/http';
import { check, sleep } from 'k6';
import { uuidv4 } from 'https://jslib.k6.io/k6-utils/1.2.0/index.js';

export const options = {
  stages: [
    { duration: '10s', target: 1000 },   // 快速上升到 1000 用户
    { duration: '30s', target: 5000 },   // 继续上升到 5000 用户
    { duration: '1m', target: 10000 },   // 峰值 10000 用户
    { duration: '30s', target: 0 },      // 下降
  ],
  thresholds: {
    http_req_duration: ['p(95)<1000'],   // 95% 请求 < 1s
    http_req_failed: ['rate<0.05'],      // 错误率 < 5%
  },
};

export default function () {
  const userId = uuidv4();
  const productId = 'hot-product-123';
  
  const payload = JSON.stringify({
    productId: productId,
    quantity: 1,
    userId: userId
  });
  
  const params = {
    headers: {
      'Content-Type': 'application/json',
      'Idempotency-Key': uuidv4()
    },
  };
  
  const res = http.post('http://gateway:5000/api/orders', payload, params);
  
  check(res, {
    'accepted or queued': (r) => r.status === 201 || r.status === 202,
  });
  
  sleep(0.1); // 少量延迟
}

监控与告警 ​

Prometheus 告警规则 ​

yaml
groups:
  - name: extreme_concurrency_alerts
    rules:
      # 队列积压过多
      - alert: HighQueueBacklog
        expr: order_queue_size > 10000
        for: 1m
        labels:
          severity: warning
        annotations:
          summary: "High order queue backlog"
      
      # 断路器打开
      - alert: CircuitBreakerOpen
        expr: circuit_breaker_state == 1
        for: 30s
        labels:
          severity: critical
        annotations:
          summary: "Circuit breaker is open"
      
      # Redis 负载过高
      - alert: HighRedisLoad
        expr: redis_used_memory_percentage > 80
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Redis memory usage too high"

最佳实践总结 ​

✅ DO ​

  1. 多层防护:网关限流 + 应用排队 + 数据库保护
  2. 快速失败:超出容量立即拒绝
  3. 异步处理:削峰填谷
  4. 降级策略:保证核心功能可用
  5. 监控告警:及时发现异常

❌ DON'T ​

  1. 不要硬扛:超过容量要拒绝
  2. 不要无限重试:加重系统负担
  3. 不要忽略背压:防止内存溢出
  4. 不要忘记清理:定期清理过期数据

总结 ​

极端并发场景下的幂等性保障需要综合策略:

✅ 限流:令牌桶、漏桶算法
✅ 排队:异步化处理,削峰填谷
✅ 预扣减:Redis 原子操作
✅ 降级:断路器模式
✅ 监控:实时告警

通过这些技术,可以在高并发场景下保持系统的稳定性和可用性。

Released under the MIT License.