边界 - 极端并发穿透
概述
在极端高并发场景下,即使实现了完善的幂等性机制,仍可能出现"穿透"问题:
- 缓存穿透:大量请求同时绕过缓存
- 锁竞争:分布式锁成为瓶颈
- 数据库压力:唯一索引检查导致性能下降
- 雪崩效应:系统级联失败
典型场景
场景 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
- 多层防护:网关限流 + 应用排队 + 数据库保护
- 快速失败:超出容量立即拒绝
- 异步处理:削峰填谷
- 降级策略:保证核心功能可用
- 监控告警:及时发现异常
❌ DON'T
- 不要硬扛:超过容量要拒绝
- 不要无限重试:加重系统负担
- 不要忽略背压:防止内存溢出
- 不要忘记清理:定期清理过期数据
总结
极端并发场景下的幂等性保障需要综合策略:
✅ 限流:令牌桶、漏桶算法
✅ 排队:异步化处理,削峰填谷
✅ 预扣减:Redis 原子操作
✅ 降级:断路器模式
✅ 监控:实时告警
通过这些技术,可以在高并发场景下保持系统的稳定性和可用性。