使用Python手搓一個生產(chǎn)級限流器
“限流不是拒絕服務(wù),而是保護服務(wù)——讓系統(tǒng)在壓力下優(yōu)雅地活著,而不是突然崩塌。”
一、為什么你的服務(wù)需要限流
2023 年某知名 AI 平臺在新品發(fā)布當(dāng)晚,因未做好限流防護,短短 15 分鐘內(nèi)涌入的請求量超過正常峰值的 80 倍,數(shù)據(jù)庫連接池耗盡,整個服務(wù)雪崩。事后復(fù)盤,技術(shù)負責(zé)人苦澀地說:“如果當(dāng)時有一個限流器,最多損失部分用戶體驗,而不是全站癱瘓。”
限流(Rate Limiting),是系統(tǒng)穩(wěn)定性設(shè)計中最重要的防線之一。它的本質(zhì)是:在時間維度上控制資源的消費速率,保護下游系統(tǒng)不被壓垮,同時對用戶提供公平的服務(wù)保障。
三種經(jīng)典限流算法——令牌桶(Token Bucket)、漏桶(Leaky Bucket)、滑動窗口(Sliding Window)——各有側(cè)重,適用場景不同。本文將帶你從原理到代碼,逐一實現(xiàn)它們,并在文末給出選型建議。
二、固定窗口:最樸素的限流(以及它的致命缺陷)
在深入三大算法之前,先了解最簡單的固定窗口計數(shù)器,理解它的不足,才能體會后續(xù)算法的價值。
import time
import threading
from collections import defaultdict
class FixedWindowRateLimiter:
"""
固定窗口限流器
缺陷:窗口邊界處可能出現(xiàn) 2 倍流量突刺
"""
def __init__(self, max_requests: int, window_seconds: int):
self.max_requests = max_requests
self.window_seconds = window_seconds
self._counts = defaultdict(int)
self._window_starts = defaultdict(float)
self._lock = threading.Lock()
def is_allowed(self, key: str = "default") -> bool:
now = time.time()
with self._lock:
window_start = self._window_starts[key]
# 當(dāng)前時間超過窗口,重置計數(shù)
if now - window_start >= self.window_seconds:
self._window_starts[key] = now
self._counts[key] = 0
if self._counts[key] < self.max_requests:
self._counts[key] += 1
return True
return False
# 演示固定窗口的邊界問題
limiter = FixedWindowRateLimiter(max_requests=10, window_seconds=60)
# 假設(shè)窗口為 [0s, 60s),在第 59s 發(fā)送 10 個請求(全部通過)
# 然后在第 61s(新窗口開始)再發(fā)送 10 個請求(全部通過)
# 實際上在 [59s, 61s] 這 2 秒內(nèi)通過了 20 個請求——是上限的 2 倍!
print("固定窗口在邊界處存在 2 倍流量突刺風(fēng)險")
這個"邊界突刺"問題在高并發(fā)場景下可能讓下游系統(tǒng)瞬間承受雙倍壓力。這正是滑動窗口要解決的問題。
三、令牌桶算法:允許突發(fā),平滑限速
原理圖解
令牌桶示意圖:
┌─────────────────┐
以固定速率 │ ?? ?? ?? ?? ?? │ 桶容量上限
補充令牌 → │ ?? ?? ?? ?? │ (burst capacity)
└────────┬────────┘
│
請求到來時取令牌
↓
有令牌 → 放行請求
無令牌 → 拒絕/等待
令牌桶的核心特性:允許一定程度的流量突發(fā)(桶里積累的令牌可以被一次性消費),但長期平均速率不超過令牌補充速率。這非常適合 API 調(diào)用:用戶可能偶爾有突發(fā)需求,但不能持續(xù)高頻。
Python 實現(xiàn)
import time
import threading
from typing import Optional
class TokenBucketRateLimiter:
"""
令牌桶限流器
特性:
- 允許短時突發(fā)(桶滿時可消費全部令牌)
- 長期速率受 rate 控制
- 線程安全
"""
def __init__(self, rate: float, capacity: int):
"""
:param rate: 令牌補充速率(個/秒)
:param capacity: 桶的最大容量(突發(fā)上限)
"""
self.rate = rate
self.capacity = capacity
self._tokens = float(capacity) # 初始滿桶
self._last_refill = time.monotonic()
self._lock = threading.Lock()
def _refill(self):
"""根據(jù)時間流逝補充令牌(惰性計算,不需要后臺線程)"""
now = time.monotonic()
elapsed = now - self._last_refill
new_tokens = elapsed * self.rate
self._tokens = min(self.capacity, self._tokens + new_tokens)
self._last_refill = now
def acquire(self, tokens: int = 1) -> bool:
"""
嘗試獲取令牌
:param tokens: 需要的令牌數(shù)(支持批量消費)
:return: True=放行, False=拒絕
"""
with self._lock:
self._refill()
if self._tokens >= tokens:
self._tokens -= tokens
return True
return False
def acquire_with_wait(self, tokens: int = 1,
timeout: float = None) -> bool:
"""
等待直到獲取到令牌(阻塞版本)
:param timeout: 最長等待時間(秒),None 表示永久等待
"""
start = time.monotonic()
while True:
with self._lock:
self._refill()
if self._tokens >= tokens:
self._tokens -= tokens
return True
# 計算需要等待多久才能獲得足夠令牌
with self._lock:
deficit = tokens - self._tokens
wait_time = deficit / self.rate
if timeout is not None:
elapsed = time.monotonic() - start
if elapsed + wait_time > timeout:
return False
time.sleep(min(wait_time, 0.01)) # 最多睡 10ms 避免過度等待
@property
def current_tokens(self) -> float:
"""查看當(dāng)前令牌數(shù)(調(diào)試用)"""
with self._lock:
self._refill()
return self._tokens
# ── 使用示例 ──────────────────────────────────────────────
def demo_token_bucket():
# 每秒補充 5 個令牌,桶容量 10
limiter = TokenBucketRateLimiter(rate=5, capacity=10)
print("=== 令牌桶演示 ===")
print(f"初始令牌數(shù):{limiter.current_tokens:.1f}")
# 模擬突發(fā)請求:一次性消費 8 個令牌
results = []
for i in range(12):
allowed = limiter.acquire()
results.append("?" if allowed else "?")
print(f"請求 {i+1:2d}:{'放行' if allowed else '拒絕'} "
f"(剩余令牌:{limiter.current_tokens:.2f})")
# 等待 2 秒后令牌恢復(fù)
print("\n等待 2 秒,令牌恢復(fù)中...")
time.sleep(2)
print(f"2 秒后令牌數(shù):{limiter.current_tokens:.1f}(補充了約 10 個)")
print(f"再次請求:{'放行' if limiter.acquire() else '拒絕'}")
demo_token_bucket()
四、漏桶算法:絕對平滑,削峰填谷
原理圖解
漏桶示意圖:
請求涌入(任意速率)
↓↓↓↓↓↓↓↓
┌──────────┐
│ │ ← 桶滿則溢出(拒絕請求)
│ 隊列 │
│ 緩沖 │
│ │
└────┬─────┘
│ 以固定速率漏出(處理請求)
↓
恒定速率輸出
漏桶與令牌桶的核心區(qū)別:漏桶的輸出速率是嚴格恒定的,不允許突發(fā)。無論入流量多大,處理速率始終如一。這對下游服務(wù)的保護最為徹底,但用戶體驗上感知更明顯——即使桶里有空間,也要排隊等待。
Python 實現(xiàn)
import time
import threading
from collections import deque
from dataclasses import dataclass
from typing import Any, Optional
@dataclass
class Request:
"""封裝請求及其元數(shù)據(jù)"""
data: Any
arrive_time: float
key: str = "default"
class LeakyBucketRateLimiter:
"""
漏桶限流器
特性:
- 嚴格恒定的輸出速率
- 超出桶容量的請求直接丟棄
- 天然實現(xiàn)請求整形(Traffic Shaping)
"""
def __init__(self, rate: float, capacity: int):
"""
:param rate: 漏出速率(請求數(shù)/秒)
:param capacity: 桶容量(最大排隊數(shù))
"""
self.rate = rate
self.capacity = capacity
self._interval = 1.0 / rate # 兩次漏出之間的間隔
self._queue: deque = deque()
self._lock = threading.Lock()
self._last_leak_time = time.monotonic()
# 啟動后臺漏出線程
self._running = True
self._worker = threading.Thread(target=self._leak_worker, daemon=True)
self._worker.start()
def _leak_worker(self):
"""后臺工作線程:以固定速率處理請求"""
while self._running:
time.sleep(self._interval)
with self._lock:
if self._queue:
request = self._queue.popleft()
wait_time = time.monotonic() - request.arrive_time
print(f" ?? 漏出處理:key={request.key}, "
f"等待={wait_time*1000:.1f}ms")
def submit(self, data: Any, key: str = "default") -> bool:
"""
提交請求到漏桶
:return: True=入隊成功, False=桶已滿(丟棄)
"""
with self._lock:
if len(self._queue) >= self.capacity:
print(f" ?? 桶已滿,丟棄請求:key={key}")
return False
request = Request(data=data, arrive_time=time.monotonic(), key=key)
self._queue.append(request)
print(f" ?? 請求入隊:key={key}, 隊列深度={len(self._queue)}")
return True
def stop(self):
self._running = False
# ── 使用示例 ──────────────────────────────────────────────
def demo_leaky_bucket():
print("\n=== 漏桶演示(2 req/s,桶容量 5)===")
limiter = LeakyBucketRateLimiter(rate=2, capacity=5)
# 模擬突發(fā):瞬間提交 8 個請求
for i in range(8):
limiter.submit(f"request-{i}", key=f"user_{i % 3}")
# 等待漏桶處理完畢
time.sleep(4)
limiter.stop()
print("漏桶處理完成,輸出速率嚴格為 2 req/s")
demo_leaky_bucket()
五、滑動窗口算法:精準統(tǒng)計,消除邊界突刺
原理圖解
固定窗口(存在邊界問題):
|← 窗口1(0-60s)→|← 窗口2(60-120s)→|
| 9 requests | 9 requests |
↑
邊界時刻:前后 2 秒內(nèi)實際有 18 個請求!
滑動窗口(精確統(tǒng)計):
當(dāng)前時間 t=75s,窗口大小 60s
統(tǒng)計 [15s, 75s] 內(nèi)的請求數(shù) → 精確到每一秒
滑動窗口分為兩種實現(xiàn):滑動日志(精確但內(nèi)存大)和滑動計數(shù)器(近似但高效)。
Python 實現(xiàn)
import time
import threading
from collections import deque
from typing import Dict
class SlidingWindowLogLimiter:
"""
滑動窗口日志限流器(精確版)
記錄每個請求的時間戳,精確統(tǒng)計窗口內(nèi)的請求數(shù)
適合:低流量、高精度要求的場景
缺陷:內(nèi)存占用隨請求量線性增長
"""
def __init__(self, max_requests: int, window_seconds: float):
self.max_requests = max_requests
self.window_seconds = window_seconds
# key → deque of timestamps
self._logs: Dict[str, deque] = {}
self._lock = threading.Lock()
def is_allowed(self, key: str = "default") -> bool:
now = time.monotonic()
window_start = now - self.window_seconds
with self._lock:
if key not in self._logs:
self._logs[key] = deque()
log = self._logs[key]
# 清除窗口外的舊記錄
while log and log[0] <= window_start:
log.popleft()
# 判斷窗口內(nèi)請求數(shù)
if len(log) < self.max_requests:
log.append(now)
return True
return False
def current_count(self, key: str = "default") -> int:
"""當(dāng)前窗口內(nèi)的請求數(shù)"""
now = time.monotonic()
window_start = now - self.window_seconds
with self._lock:
if key not in self._logs:
return 0
return sum(1 for t in self._logs[key] if t > window_start)
class SlidingWindowCounterLimiter:
"""
滑動窗口計數(shù)器限流器(高效版)
將窗口劃分為多個小格子,用計數(shù)器近似統(tǒng)計
適合:高流量、對精度要求不苛刻的場景
Redis 分布式限流的主流方案
"""
def __init__(self, max_requests: int, window_seconds: float,
precision: int = 10):
"""
:param precision: 將窗口劃分為多少個子格子(越多越精確,越耗內(nèi)存)
"""
self.max_requests = max_requests
self.window_seconds = window_seconds
self.precision = precision
self.slot_duration = window_seconds / precision
# 每個 key 的滑動計數(shù)器:{slot_index: count}
self._counters: Dict[str, Dict[int, int]] = {}
self._lock = threading.Lock()
def _current_slot(self) -> int:
"""當(dāng)前時間對應(yīng)的格子索引"""
return int(time.monotonic() / self.slot_duration)
def is_allowed(self, key: str = "default") -> bool:
current_slot = self._current_slot()
# 有效的格子范圍
valid_slots = set(range(current_slot - self.precision + 1,
current_slot + 1))
with self._lock:
if key not in self._counters:
self._counters[key] = {}
counter = self._counters[key]
# 清除過期格子
expired = [s for s in counter if s not in valid_slots]
for s in expired:
del counter[s]
# 統(tǒng)計窗口內(nèi)總請求數(shù)
total = sum(counter.values())
if total < self.max_requests:
counter[current_slot] = counter.get(current_slot, 0) + 1
return True
return False
# ── 對比演示 ──────────────────────────────────────────────
def demo_sliding_window():
print("\n=== 滑動窗口演示(10 req/10s)===")
log_limiter = SlidingWindowLogLimiter(max_requests=10, window_seconds=10)
counter_limiter = SlidingWindowCounterLimiter(max_requests=10,
window_seconds=10,
precision=10)
# 模擬在 5 秒內(nèi)發(fā)送 15 個請求
results_log = []
results_counter = []
for i in range(15):
results_log.append(log_limiter.is_allowed("user_1"))
results_counter.append(counter_limiter.is_allowed("user_1"))
time.sleep(0.2)
print(f"滑動日志結(jié)果:{['?' if r else '?' for r in results_log]}")
print(f"計數(shù)器結(jié)果: {['?' if r else '?' for r in results_counter]}")
print(f"日志版窗口內(nèi)請求數(shù):{log_limiter.current_count('user_1')}")
demo_sliding_window()
六、整合實戰(zhàn):帶限流的 FastAPI 中間件
將三種限流器整合為一個統(tǒng)一接口,集成到 FastAPI 中:
from abc import ABC, abstractmethod
from fastapi import FastAPI, Request, HTTPException
from fastapi.responses import JSONResponse
import time
class RateLimiter(ABC):
"""限流器抽象基類"""
@abstractmethod
def is_allowed(self, key: str) -> bool:
pass
@abstractmethod
def get_info(self, key: str) -> dict:
"""返回限流狀態(tài)信息(用于響應(yīng)頭)"""
pass
class RateLimitMiddleware:
"""
FastAPI 限流中間件
支持按 IP、用戶 ID、API Key 等維度限流
"""
def __init__(self, app, limiter: RateLimiter,
key_func=None, exclude_paths: list = None):
self.app = app
self.limiter = limiter
self.key_func = key_func or self._default_key_func
self.exclude_paths = exclude_paths or ["/health", "/docs"]
@staticmethod
def _default_key_func(request: Request) -> str:
"""默認按客戶端 IP 限流"""
forwarded_for = request.headers.get("X-Forwarded-For")
if forwarded_for:
return forwarded_for.split(",")[0].strip()
return request.client.host
async def __call__(self, scope, receive, send):
if scope["type"] == "http":
request = Request(scope, receive)
# 排除不限流的路徑
if request.url.path not in self.exclude_paths:
key = self.key_func(request)
if not self.limiter.is_allowed(key):
response = JSONResponse(
status_code=429,
content={
"error": "Too Many Requests",
"message": "請求過于頻繁,請稍后再試",
"retry_after": 1
},
headers={
"Retry-After": "1",
"X-RateLimit-Limit": "100",
}
)
await response(scope, receive, send)
return
await self.app(scope, receive, send)
# 構(gòu)建應(yīng)用
app = FastAPI(title="限流演示 API")
# 選擇限流算法(可熱切換)
# 方案1:令牌桶(允許突發(fā))
rate_limiter = TokenBucketRateLimiter(rate=10, capacity=20)
# 方案2:滑動窗口(精確統(tǒng)計)
# rate_limiter = SlidingWindowLogLimiter(max_requests=100, window_seconds=60)
app.add_middleware(RateLimitMiddleware, limiter=rate_limiter)
@app.get("/api/data")
async def get_data(request: Request):
client_ip = request.client.host
return {
"message": "請求成功",
"client": client_ip,
"timestamp": time.time()
}
@app.get("/health")
async def health():
return {"status": "ok"}
七、Redis 分布式限流:生產(chǎn)環(huán)境的標準方案
單機限流器在多實例部署時會失效——每個實例各自維護狀態(tài),總請求數(shù)可能超出預(yù)期。生產(chǎn)環(huán)境必須用 Redis 實現(xiàn)分布式限流:
import redis
import time
class RedisTokenBucketLimiter:
"""
基于 Redis 的分布式令牌桶限流器
使用 Lua 腳本保證原子性
"""
# Lua 腳本:原子化地補充令牌并判斷是否放行
LUA_SCRIPT = """
local key = KEYS[1]
local rate = tonumber(ARGV[1])
local capacity = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local requested = tonumber(ARGV[4])
local last_tokens = tonumber(redis.call('hget', key, 'tokens') or capacity)
local last_refill = tonumber(redis.call('hget', key, 'last_refill') or now)
-- 補充令牌
local elapsed = now - last_refill
local new_tokens = math.min(capacity, last_tokens + elapsed * rate)
local allowed = 0
if new_tokens >= requested then
new_tokens = new_tokens - requested
allowed = 1
end
redis.call('hset', key, 'tokens', new_tokens)
redis.call('hset', key, 'last_refill', now)
redis.call('expire', key, math.ceil(capacity / rate) * 2)
return {allowed, math.floor(new_tokens)}
"""
def __init__(self, redis_url: str, rate: float, capacity: int):
self.redis = redis.from_url(redis_url)
self.rate = rate
self.capacity = capacity
self._script = self.redis.register_script(self.LUA_SCRIPT)
def is_allowed(self, key: str, tokens: int = 1) -> tuple[bool, int]:
"""
:return: (是否允許, 剩余令牌數(shù))
"""
now = time.time()
result = self._script(
keys=[f"ratelimit:{key}"],
args=[self.rate, self.capacity, now, tokens]
)
allowed = bool(result[0])
remaining = int(result[1])
return allowed, remaining
# 使用示例(需要 Redis 服務(wù))
# redis_limiter = RedisTokenBucketLimiter(
# redis_url="redis://localhost:6379",
# rate=10,
# capacity=50
# )
# allowed, remaining = redis_limiter.is_allowed("user:12345")
# print(f"放行:{allowed},剩余令牌:{remaining}")
八、三大算法橫向?qū)Ρ?/h2>
| 維度 | 令牌桶 | 漏桶 | 滑動窗口 |
|---|---|---|---|
| 突發(fā)流量 | ? 允許(桶容量內(nèi)) | ? 嚴格排隊 | ? 有限允許 |
| 輸出平滑度 | 中 | 極高 | 高 |
| 實現(xiàn)復(fù)雜度 | 低 | 中 | 中 |
| 內(nèi)存占用 | 極低 | 中(隊列) | 低(計數(shù)) |
| 精確度 | 高 | 高 | 近似(計數(shù)器版) |
| 適用場景 | API 調(diào)用限額 | 流量整形 | 精細化統(tǒng)計 |
| 典型應(yīng)用 | GitHub API | 網(wǎng)絡(luò)帶寬限速 | Nginx limit_req |
選型建議:
如果你的場景是對外 API 服務(wù),用戶可能有合理的突發(fā)需求(如批量導(dǎo)出),選令牌桶,既能控制長期速率,又不會讓正常突發(fā)請求體驗太差。
如果你的場景是保護下游服務(wù)(如數(shù)據(jù)庫、第三方 API),需要嚴格控制請求速率避免過載,選漏桶,它能將任意突發(fā)流量整形為恒定輸出。
如果你需要精確統(tǒng)計時間窗口內(nèi)的請求次數(shù)(如"每分鐘最多 100 次"的 SLA 保障),選滑動窗口,配合 Redis 可輕松實現(xiàn)分布式精確限流。
九、總結(jié)
限流是系統(tǒng)穩(wěn)定性設(shè)計的基石,三種算法各有側(cè)重:令牌桶平衡突發(fā)與限速,漏桶追求絕對平滑,滑動窗口精確感知流量。在 Python 實戰(zhàn)中,本地場景用線程安全的純 Python 實現(xiàn)即可;多實例部署時,務(wù)必引入 Redis + Lua 腳本的分布式方案,保證全局一致性。
好的限流器不只是"拒絕請求",它還應(yīng)該做到:返回清晰的 429 Too Many Requests、在響應(yīng)頭中告知 Retry-After 時間、對不同用戶等級設(shè)置差異化配額。這些細節(jié),才是用戶友好型 API 的體現(xiàn)。
限流的本質(zhì)是保護——保護你的服務(wù),也保護你的用戶。把這道門守好,系統(tǒng)才能在風(fēng)雨中巋然不動。
以上就是使用Python手搓一個生產(chǎn)級限流器的詳細內(nèi)容,更多關(guān)于Python限流器的資料請關(guān)注腳本之家其它相關(guān)文章!
相關(guān)文章
python安裝numpy&安裝matplotlib& scipy的教程
下面小編就為大家?guī)硪黄猵ython安裝numpy&安裝matplotlib& scipy的教程。小編覺得挺不錯的,現(xiàn)在就分享給大家,也給大家做個參考。一起跟隨小編過來看看吧2017-11-11
使用Python構(gòu)建Markdown轉(zhuǎn)Word文檔轉(zhuǎn)換器
在當(dāng)今的文檔處理中,Markdown因其簡潔的語法和易讀性而廣受歡迎,而Microsoft Word(DOCX格式)則因其廣泛的兼容性和專業(yè)的排版效果成為商業(yè)文檔的標準,本文將介紹如何使用Python構(gòu)建一個帶有圖形界面的Markdown轉(zhuǎn)Word文檔轉(zhuǎn)換器,需要的朋友可以參考下2025-02-02
初學(xué)者學(xué)習(xí)Python好還是Java好
在本篇文章里小編給大家分享的是關(guān)于初學(xué)者學(xué)習(xí)Python好還是Java好的相關(guān)內(nèi)容,需要的朋友們可以學(xué)習(xí)下。2020-05-05
Python實現(xiàn)批量識別銀行卡號碼以及自動寫入Excel表格步驟詳解
這篇文章主要介紹了使用Python實現(xiàn)高效摸魚,批量識別銀行卡號碼并且自動寫入Excel表格,文中通過示例代碼介紹的非常詳細,對大家的學(xué)習(xí)或者工作具有一定的參考學(xué)習(xí)價值,需要的朋友們下面隨著小編來一起學(xué)習(xí)吧2023-01-01
python+opencv實現(xiàn)移動偵測(幀差法)
這篇文章主要為大家詳細介紹了python+opencv實現(xiàn)移動偵測,文中示例代碼介紹的非常詳細,具有一定的參考價值,感興趣的小伙伴們可以參考一下2020-03-03

