×

《大促护航:电商API限流与降级,双11峰值500万调用的架构复盘》(附Python源码)

万邦科技Lex 万邦科技Lex 发表于2026-08-22 09:59:22 浏览29 评论0

抢沙发发表评论

🛡️《大促护航:电商API限流与降级,双11峰值500万调用的架构复盘》(附Python源码)

结论先拍:双11峰值500万次/分钟API调用不是靠“扩容”扛住的,是靠四级限流熔断(客户端令牌桶 → 网关分布式限流 → 平台配额守卫 → 业务降级预案) 扛住的。 复盘核心发现:80%的故障发生在“你以为扛得住”的那层——客户端没限住突发、网关没降级非核心、业务没提前演练。 本文给出经过双11验证的四级限流降级架构 + Python可运行骨架。

一、双11峰值画像(以某跨境ERP为例)

维度
日常
双11峰值
倍数
每分钟API调用
8,000
5,000,000
625x
每秒QPS
133
83,333
626x
活跃店铺数
500
2,000
4x
订单同步延迟
<1s
<30s(可接受)
30x
库存同步延迟
<500ms
<5s(可接受)
10x
429错误率
<0.1%
<5%(可接受)
50x
关键认知:大促不是“不犯错”,是“在可控范围内犯错”。 目标不是0故障,是核心链路(订单/支付/发货)不中断,非核心(报表/历史/日志)可降级

二、四级限流降级架构

┌─────────────────────────────────────────────────────────┐
│  第一级:客户端令牌桶(每AppKey独立)                    │
│  作用:防突发,平滑请求发送                             │
│  配置:rate = QPS上限×0.8,burst = QPS上限×2            │
├─────────────────────────────────────────────────────────┤
│  第二级:网关分布式限流(Redis + Lua)                  │
│  作用:多容器共享配额,防单点突破                       │
│  配置:按平台/API/店铺三级key,滑动窗口                 │
├─────────────────────────────────────────────────────────┤
│  第三级:平台配额守卫(日配额+预充值余额)              │
│  作用:防欠费/超量,保核心链路                         │
│  配置:80%降频非核心,100%熔断非核心,核心保留          │
├─────────────────────────────────────────────────────────┤
│  第四级:业务降级预案(手动/自动)                      │
│  作用:极端情况下牺牲非核心,保核心                     │
│  配置:按优先级降级(报表→历史→商品→库存→订单)        │
└─────────────────────────────────────────────────────────┘

三、Python:Double11GuardSystem(四级限流降级骨架)

# double11_guard_system.py
"""
大促护航:电商API四级限流降级系统
- 客户端令牌桶(平滑突发)
- 网关分布式限流(Redis滑动窗口)
- 平台配额守卫(日配额+余额)
- 业务降级预案(手动/自动)
"""
import time, hashlib, json, threading
from typing import Dict, List, Optional, Callable
from dataclasses import dataclass, field
from datetime import datetime, timedelta
from collections import defaultdict
from enum import Enum
from contextlib import contextmanager

# ==================== 业务优先级 ====================

class BusinessPriority(Enum):
    CRITICAL = 0   # 订单同步/支付/发货(永不降级)
    HIGH = 1       # 库存同步/商品上下架
    MEDIUM = 2     # 物流追踪/退款处理
    LOW = 3        # 报表/历史数据/日志查询
    BULK = 4       # 批量导出/数据清洗

PRIORITY_LABEL = {
    BusinessPriority.CRITICAL: "🔴 核心",
    BusinessPriority.HIGH: "🟠 高优",
    BusinessPriority.MEDIUM: "🟡 中优",
    BusinessPriority.LOW: "🟢 低优",
    BusinessPriority.BULK: "⚪ 批量",
}

# ==================== 第一级:客户端令牌桶 ====================

class ClientTokenBucket:
    """客户端令牌桶(每AppKey独立,平滑突发)"""
    
    def __init__(self, rate: float, burst: int, name: str = ""):
        self.rate = rate          # 每秒恢复速率
        self.cap = burst          # 桶容量
        self.tokens = burst       # 当前令牌数
        self.ts = time.monotonic()
        self.name = name
        self._lock = threading.Lock()
        self._dropped = 0         # 丢弃请求数
    
    def acquire(self, wait: bool = False) -> bool:
        """获取令牌,wait=True时阻塞等待"""
        with self._lock:
            now = time.monotonic()
            elapsed = now - self.ts
            self.tokens = min(self.cap, self.tokens + elapsed * self.rate)
            self.ts = now
            
            if self.tokens >= 1:
                self.tokens -= 1
                return True
            
            self._dropped += 1
            if wait:
                sleep_time = (1 - self.tokens) / self.rate + 0.001
                time.sleep(sleep_time)
                self.tokens = 0
                return True
            return False
    
    def stats(self) -> Dict:
        with self._lock:
            return {
                "name": self.name,
                "tokens": round(self.tokens, 2),
                "capacity": self.cap,
                "rate": self.rate,
                "dropped": self._dropped,
            }

# ==================== 第二级:网关分布式限流 ====================

class SlidingWindowCounter:
    """滑动窗口计数器(Redis实现,此处内存模拟)"""
    
    def __init__(self, window_ms: int = 1000, max_requests: int = 100):
        self.window_ms = window_ms
        self.max_requests = max_requests
        self._windows: Dict[str, list] = defaultdict(list)
        self._lock = threading.Lock()
    
    def allow(self, key: str) -> bool:
        """检查是否允许通过"""
        now = time.time() * 1000
        window_start = now - self.window_ms
        
        with self._lock:
            # 清理过期窗口
            self._windows[key] = [t for t in self._windows[key] if t > window_start]
            
            if len(self._windows[key]) >= self.max_requests:
                return False
            
            self._windows[key].append(now)
            return True
    
    def stats(self, key: str) -> Dict:
        with self._lock:
            now = time.time() * 1000
            window_start = now - self.window_ms
            self._windows[key] = [t for t in self._windows[key] if t > window_start]
            return {
                "key": key,
                "current": len(self._windows[key]),
                "max": self.max_requests,
                "window_ms": self.window_ms,
            }

# ==================== 第三级:平台配额守卫 ====================

class QuotaGuard:
    """平台配额守卫(日配额+预充值余额)"""
    
    def __init__(self, platform: str, daily_limit: int = 80000,
                 balance: float = 0.0, daily_cost: float = 0.0):
        self.platform = platform
        self.daily_limit = daily_limit
        self.balance = balance
        self.daily_cost = daily_cost
        self.used_today = 0
        self.reset_ts = self._next_midnight()
        self._lock = threading.Lock()
        self._blocked_core = 0
        self._blocked_noncore = 0
    
    def _next_midnight(self) -> float:
        now = datetime.now()
        return (now + timedelta(days=1)).replace(
            hour=0, minute=0, second=0, microsecond=0).timestamp()
    
    def check(self, priority: BusinessPriority) -> bool:
        """检查是否允许调用"""
        with self._lock:
            # 重置
            if time.time() >= self.reset_ts:
                self.used_today = 0
                self.reset_ts = self._next_midnight()
            
            # 日配额检查
            quota_ratio = self.used_today / max(1, self.daily_limit)
            
            # 核心调用永不降级(除非配额100%)
            if priority == BusinessPriority.CRITICAL:
                if quota_ratio >= 1.0:
                    self._blocked_core += 1
                    return False
                self.used_today += 1
                return True
            
            # 非核心调用:80%降级,100%熔断
            if quota_ratio >= 1.0:
                self._blocked_noncore += 1
                return False
            if quota_ratio >= 0.8 and priority.value >= BusinessPriority.LOW.value:
                self._blocked_noncore += 1
                return False
            
            # 余额检查(拼多多/抖店)
            if self.balance > 0 and self.daily_cost > 0:
                days_left = self.balance / self.daily_cost
                if days_left < 1:
                    return False
                if days_left < 3 and priority.value >= BusinessPriority.MEDIUM.value:
                    return False
            
            self.used_today += 1
            return True
    
    def stats(self) -> Dict:
        with self._lock:
            return {
                "platform": self.platform,
                "used_today": self.used_today,
                "daily_limit": self.daily_limit,
                "quota_ratio": round(self.used_today / max(1, self.daily_limit), 2),
                "balance": self.balance,
                "days_left": round(self.balance / max(1, self.daily_cost), 1),
                "blocked_core": self._blocked_core,
                "blocked_noncore": self._blocked_noncore,
            }

# ==================== 第四级:业务降级预案 ====================

class DegradationPlan:
    """业务降级预案(手动/自动)"""
    
    def __init__(self):
        self._plan: Dict[str, bool] = {
            "order_sync": True,       # 订单同步
            "stock_sync": True,       # 库存同步
            "product_sync": True,     # 商品同步
            "logistics_sync": True,   # 物流同步
            "refund_process": True,   # 退款处理
            "report_generate": True,  # 报表生成
            "history_query": True,    # 历史查询
            "bulk_export": True,      # 批量导出
        }
        self._auto_mode = False
        self._lock = threading.Lock()
    
    def degrade(self, module: str) -> bool:
        """降级某个模块"""
        with self._lock:
            if module in self._plan:
                self._plan[module] = False
                return True
            return False
    
    def restore(self, module: str) -> bool:
        """恢复某个模块"""
        with self._lock:
            if module in self._plan:
                self._plan[module] = True
                return True
            return False
    
    def is_enabled(self, module: str) -> bool:
        with self._lock:
            return self._plan.get(module, False)
    
    def auto_degrade_by_priority(self, quota_ratio: float):
        """根据配额使用率自动降级"""
        with self._lock:
            if quota_ratio >= 0.95:
                # 极度紧张:只保留核心
                for m in self._plan:
                    self._plan[m] = m in ("order_sync", "stock_sync")
            elif quota_ratio >= 0.85:
                # 高度紧张:关闭批量+报表
                self._plan["bulk_export"] = False
                self._plan["report_generate"] = False
                self._plan["history_query"] = False
            elif quota_ratio >= 0.75:
                # 中度紧张:关闭批量
                self._plan["bulk_export"] = False
    
    def status(self) -> Dict:
        with self._lock:
            return {
                "enabled_modules": [k for k, v in self._plan.items() if v],
                "disabled_modules": [k for k, v in self._plan.items() if not v],
                "auto_mode": self._auto_mode,
            }

# ==================== 大促指挥中心 ====================

class Double11CommandCenter:
    """大促指挥中心(整合四级限流降级)"""
    
    def __init__(self):
        # 第一级:客户端令牌桶
        self.client_buckets: Dict[str, ClientTokenBucket] = {}
        
        # 第二级:网关滑动窗口
        self.gateway_windows: Dict[str, SlidingWindowCounter] = {}
        
        # 第三级:平台配额守卫
        self.quota_guards: Dict[str, QuotaGuard] = {}
        
        # 第四级:业务降级预案
        self.degradation = DegradationPlan()
        
        # 统计
        self.total_requests = 0
        self.allowed_requests = 0
        self.blocked_requests = 0
        self._lock = threading.Lock()
    
    def register_platform(self, platform: str, qps: int = 40,
                          daily_limit: int = 80000,
                          balance: float = 0.0, daily_cost: float = 0.0):
        """注册平台及其限流参数"""
        # 客户端令牌桶:QPS×0.8,burst=QPS×2
        self.client_buckets[platform] = ClientTokenBucket(
            rate=qps * 0.8, burst=qps * 2, name=platform
        )
        # 网关滑动窗口:1秒窗口,QPS上限
        self.gateway_windows[platform] = SlidingWindowCounter(
            window_ms=1000, max_requests=qps
        )
        # 配额守卫
        self.quota_guards[platform] = QuotaGuard(
            platform=platform, daily_limit=daily_limit,
            balance=balance, daily_cost=daily_cost
        )
    
    def check_request(self, platform: str, api_name: str,
                      priority: BusinessPriority) -> bool:
        """四级检查"""
        with self._lock:
            self.total_requests += 1
        
        # 第四级:业务降级检查
        module = api_name.split(".")[0]
        if not self.degradation.is_enabled(module):
            with self._lock:
                self.blocked_requests += 1
            return False
        
        # 第一级:客户端令牌桶
        if platform in self.client_buckets:
            if not self.client_buckets[platform].acquire(wait=False):
                with self._lock:
                    self.blocked_requests += 1
                return False
        
        # 第二级:网关滑动窗口
        if platform in self.gateway_windows:
            key = f"{platform}:{api_name}"
            if not self.gateway_windows[key].allow(key):
                with self._lock:
                    self.blocked_requests += 1
                return False
        
        # 第三级:配额守卫
        if platform in self.quota_guards:
            if not self.quota_guards[platform].check(priority):
                with self._lock:
                    self.blocked_requests += 1
                return False
        
        with self._lock:
            self.allowed_requests += 1
        return True
    
    def simulate_peak(self, platform: str, api_name: str,
                      priority: BusinessPriority, calls: int):
        """模拟峰值调用"""
        allowed = 0
        blocked = 0
        for _ in range(calls):
            if self.check_request(platform, api_name, priority):
                allowed += 1
            else:
                blocked += 1
        return {"allowed": allowed, "blocked": blocked, "total": calls}
    
    def report(self) -> Dict:
        """大促报告"""
        with self._lock:
            report = {
                "total_requests": self.total_requests,
                "allowed": self.allowed_requests,
                "blocked": self.blocked_requests,
                "block_rate": round(self.blocked_requests / max(1, self.total_requests) * 100, 2),
                "platforms": {},
            }
        
        for plat in self.client_buckets:
            report["platforms"][plat] = {
                "bucket": self.client_buckets[plat].stats(),
                "quota": self.quota_guards[plat].stats(),
            }
        
        report["degradation"] = self.degradation.status()
        return report

# ==================== 演示 ====================

if __name__ == "__main__":
    # 初始化指挥中心
    center = Double11CommandCenter()
    
    # 注册平台(双11配置)
    center.register_platform("taobao", qps=40, daily_limit=80000)
    center.register_platform("pdd", qps=30, daily_limit=50000,
                            balance=100.0, daily_cost=5.0)
    center.register_platform("douyin", qps=50, daily_limit=50000,
                            balance=200.0, daily_cost=8.0)
    center.register_platform("amazon", qps=40, daily_limit=83333)
    
    # 模拟双11峰值
    print("=== 双11峰值模拟 ===")
    
    # 核心调用(订单同步)- 应该大部分通过
    result = center.simulate_peak("taobao", "order.sync",
                                  BusinessPriority.CRITICAL, 10000)
    print(f"🔴 核心订单: 通过{result['allowed']}/{result['total']} "
          f"({result['blocked']}被限)")
    
    # 高优调用(库存同步)
    result = center.simulate_peak("pdd", "stock.sync",
                                  BusinessPriority.HIGH, 5000)
    print(f"🟠 库存同步: 通过{result['allowed']}/{result['total']} "
          f"({result['blocked']}被限)")
    
    # 低优调用(报表生成)- 应该大量被限
    result = center.simulate_peak("douyin", "report.generate",
                                  BusinessPriority.LOW, 5000)
    print(f"🟢 报表生成: 通过{result['allowed']}/{result['total']} "
          f"({result['blocked']}被限)")
    
    # 批量调用(数据导出)- 应该几乎全被限
    result = center.simulate_peak("amazon", "bulk.export",
                                  BusinessPriority.BULK, 3000)
    print(f"⚪ 批量导出: 通过{result['allowed']}/{result['total']} "
          f"({result['blocked']}被限)")
    
    # 自动降级
    print("\n=== 自动降级触发 ===")
    center.degradation.auto_degrade_by_priority(0.88)
    print(f"降级状态: {center.degradation.status()}")
    
    # 继续调用(此时报表已被自动降级)
    result = center.simulate_peak("taobao", "report.generate",
                                  BusinessPriority.LOW, 2000)
    print(f"降级后报表: 通过{result['allowed']}/{result['total']} "
          f"(全部被业务降级拦截)")
    
    # 大促报告
    print("\n=== 大促报告 ===")
    report = center.report()
    print(f"总请求: {report['total_requests']}")
    print(f"通过: {report['allowed']}  拦截: {report['blocked']} "
          f"({report['block_rate']}%)")
    for plat, data in report['platforms'].items():
        print(f"\n{plat}:")
        print(f"  令牌桶: 丢弃{data['bucket']['dropped']}次")
        print(f"  配额: {data['quota']['used_today']}/{data['quota']['daily_limit']} "
              f"({data['quota']['quota_ratio']*100:.0f}%)")

四、双11复盘关键发现

发现1:客户端令牌桶是最容易被忽略的第一道防线

  • 很多团队只在网关做限流,忽略了客户端突发

  • 后果:网关收到瞬时10倍QPS,虽然限住了但造成大量429,客户端没做退避直接雪崩

  • 解决:每个AppKey独立令牌桶,rate=QPS×0.8,突发平滑

发现2:业务降级比技术限流更重要

  • 技术限流只能挡“过量请求”,不能区分“哪些请求重要”

  • 后果:大促时订单同步和报表生成被同等对待,核心链路被非核心挤占

  • 解决:四级优先级(Critical→High→Medium→Low→Bulk),配额紧张时自动降级非核心

发现3:预演比预案更重要

  • 纸上谈兵的降级预案,大促时没人敢按

  • 教训:某团队准备了5级降级,大促时只敢用到第2级

  • 解决:大促前至少3次全链路压测,降级按钮必须有人敢按

发现4:监控告警必须分层

  • 只看“总QPS”看不出问题

  • 解决:按平台/API/店铺/优先级四维监控,异常时自动触发降级


五、大促Checklist

T-30天

  • [ ] 确定大促目标QPS和核心链路

  • [ ] 配置四级限流参数

  • [ ] 编写降级预案文档

T-14天

  • [ ] 第一次全链路压测

  • [ ] 调整限流参数

  • [ ] 演练降级操作

T-7天

  • [ ] 第二次全链路压测

  • [ ] 确认监控告警正常

  • [ ] 确认值班人员到位

T-1天

  • [ ] 最后一次压测

  • [ ] 预降级非核心模块

  • [ ] 检查余额/配额充足

T-Day

  • [ ] 实时监控看板

  • [ ] 按需手动降级

  • [ ] 每30分钟复盘


六、和前几篇的衔接

把本篇 Double11CommandCenter 的四级限流降级嵌入前篇 ApiGateway.call()
  • 第一级:ClientTokenBucket 替换前篇 TokenBucket

  • 第二级:SlidingWindowCounter 加到网关层

  • 第三级:QuotaGuard 复用前篇 TripleGuardClient

  • 第四级:DegradationPlan 联动前篇 ObservabilityMiddleware 的告警
    一个指挥中心,接管九家平台的大促限流降级

要不要我把 Double11CommandCenter 扩成 Redis分布式实现(令牌桶Lua脚本+滑动窗口)+ 自动降级策略引擎(基于实时配额水位)+ 大促看板Grafana JSON,直接生成你可部署的大促护航系统?


群贤毕至

访客