本文目录导读:

我来详细介绍几种Python接口限流的实现方案:
基于时间窗口的简单限流
固定窗口限流
import time
from functools import wraps
from collections import defaultdict
class FixedWindowRateLimiter:
def __init__(self, max_requests=100, window_size=60):
self.max_requests = max_requests # 最大请求数
self.window_size = window_size # 时间窗口(秒)
self.requests = defaultdict(list) # {client_id: [timestamp1, timestamp2, ...]}
def is_allowed(self, client_id):
"""检查是否允许请求"""
current_time = time.time()
window_start = current_time - self.window_size
# 获取该客户端的请求记录
client_requests = self.requests[client_id]
# 清除过期记录
self.requests[client_id] = [
t for t in client_requests if t > window_start
]
# 检查是否超过限制
if len(self.requests[client_id]) >= self.max_requests:
return False
# 添加当前请求
self.requests[client_id].append(current_time)
return True
# 使用装饰器
def rate_limit(max_requests=100, window_size=60):
limiter = FixedWindowRateLimiter(max_requests, window_size)
def decorator(func):
@wraps(func)
def wrapper(client_id, *args, **kwargs):
if not limiter.is_allowed(client_id):
return {"error": "请求过于频繁", "status": 429}, 429
return func(client_id, *args, **kwargs)
return wrapper
return decorator
# 使用示例
@rate_limit(max_requests=5, window_size=10) # 10秒内最多5次请求
def api_endpoint(client_id, data):
return {"message": "success", "data": data}
滑动窗口限流(更精确)
import time
from collections import deque
class SlidingWindowRateLimiter:
def __init__(self, max_requests=100, window_size=60):
self.max_requests = max_requests
self.window_size = window_size
self.windows = defaultdict(lambda: deque())
def is_allowed(self, client_id):
current_time = time.time()
window = self.windows[client_id]
# 删除过期请求
while window and window[0] <= current_time - self.window_size:
window.popleft()
# 检查是否超过限制
if len(window) >= self.max_requests:
return False
# 添加当前请求
window.append(current_time)
return True
# 使用示例
limiter = SlidingWindowRateLimiter(max_requests=10, window_size=60)
def api_handler(client_id, request_data):
if not limiter.is_allowed(client_id):
return {"error": "速率限制", "retry_after": 60}, 429
# 处理请求
return {"status": "success"}
基于令牌桶算法
import time
import threading
from collections import defaultdict
class TokenBucket:
def __init__(self, capacity, fill_rate):
"""
capacity: 令牌桶容量
fill_rate: 令牌生成速率(个/秒)
"""
self.capacity = capacity
self.fill_rate = fill_rate
self.tokens = capacity
self.last_refill = time.time()
self.lock = threading.Lock()
def consume(self, tokens=1):
"""消费令牌"""
with self.lock:
self._refill()
if self.tokens >= tokens:
self.tokens -= tokens
return True
return False
def _refill(self):
"""补充令牌"""
now = time.time()
elapsed = now - self.last_refill
new_tokens = elapsed * self.fill_rate
self.tokens = min(self.capacity, self.tokens + new_tokens)
self.last_refill = now
class TokenBucketRateLimiter:
def __init__(self, capacity=100, fill_rate=10):
self.buckets = defaultdict(
lambda: TokenBucket(capacity, fill_rate)
)
def is_allowed(self, client_id):
return self.buckets[client_id].consume()
# 使用示例
limiter = TokenBucketRateLimiter(capacity=10, fill_rate=2) # 每秒2个令牌
@rate_limit_with_token_bucket
def protected_api(client_id, data):
if not limiter.is_allowed(client_id):
return {"error": "限流中", "status": 429}, 429
return {"data": data}
使用第三方库 - Flask-Limiter
# 安装: pip install flask-limiter
from flask import Flask, request, jsonify
from flask_limiter import Limiter
from flask_limiter.util import get_remote_address
app = Flask(__name__)
# 配置限流器
limiter = Limiter(
app,
key_func=get_remote_address,
default_limits=["200 per day", "50 per hour"]
)
# 全局限流
@app.route('/api/public')
@limiter.limit("10 per minute")
def public_api():
return jsonify({"status": "public endpoint"})
# 不同客户端不同限制
@app.route('/api/premium')
@limiter.limit("100 per minute")
def premium_api():
return jsonify({"status": "premium endpoint"})
# 基于用户认证的限流
@app.route('/api/user')
@limiter.limit(
"20 per minute",
key_func=lambda: request.headers.get('Authorization')
)
def user_api():
return jsonify({"status": "user endpoint"})
# 自定义错误处理
@app.errorhandler(429)
def ratelimit_handler(e):
return jsonify({
"error": "请求过于频繁",
"retry_after": e.description
}), 429
if __name__ == '__main__':
app.run()
使用Redis实现分布式限流
import redis
import time
from functools import wraps
# pip install redis
class RedisRateLimiter:
def __init__(self, redis_client=None):
self.redis = redis_client or redis.Redis(
host='localhost',
port=6379,
decode_responses=True
)
def fixed_window(self, key, max_requests=100, window_size=60):
"""
固定窗口限流(Redis实现)
"""
current_time = int(time.time())
window_key = f"rate_limit:{key}:{current_time // window_size}"
# 使用INCR和EXPIRE实现
count = self.redis.incr(window_key)
if count == 1:
self.redis.expire(window_key, window_size + 1)
return count <= max_requests
def sliding_window(self, key, max_requests=100, window_size=60):
"""
滑动窗口限流(Redis Sorted Set实现)
"""
current_time = time.time()
window_key = f"sliding_limit:{key}"
# 删除过期记录
self.redis.zremrangebyscore(
window_key,
0,
current_time - window_size
)
# 添加当前请求
self.redis.zadd(window_key, {str(current_time): current_time})
self.redis.expire(window_key, window_size + 1)
# 检查请求数量
request_count = self.redis.zcard(window_key)
return request_count <= max_requests
def token_bucket(self, key, capacity=100, refill_rate=10):
"""
令牌桶算法(Redis Lua脚本实现)
"""
lua_script = """
local key = KEYS[1]
local capacity = tonumber(ARGV[1])
local refill_rate = tonumber(ARGV[2])
local now = tonumber(ARGV[3])
local bucket = redis.call('hgetall', key)
local last_refill = 0
local tokens = capacity
if #bucket > 0 then
last_refill = tonumber(bucket[2])
tokens = tonumber(bucket[4])
end
local elapsed = now - last_refill
tokens = math.min(capacity, tokens + elapsed * refill_rate)
if tokens >= 1 then
redis.call('hmset', key, 'last_refill', now, 'tokens', tokens - 1)
redis.call('expire', key, 60)
return 1
else
return 0
end
"""
return self.redis.eval(
lua_script,
1,
f"token_bucket:{key}",
capacity,
refill_rate,
time.time()
)
# 使用示例
redis_client = redis.Redis()
limiter = RedisRateLimiter(redis_client)
def rate_limited_api(request, client_id):
# 检查限流
if not limiter.fixed_window(client_id, max_requests=100, window_size=60):
return {"error": "限流"}, 429
# 处理请求
return {"success": True}
高级限流策略 - 自适应限流
import time
from collections import defaultdict
import statistics
class AdaptiveRateLimiter:
def __init__(self, initial_limit=100, min_limit=10, max_limit=1000):
self.current_limit = initial_limit
self.min_limit = min_limit
self.max_limit = max_limit
self.response_times = defaultdict(list)
self.limit_history = []
def adjust_limit(self, client_id, response_time):
"""
根据响应时间动态调整限流阈值
"""
self.response_times[client_id].append(response_time)
# 只保留最近100个记录
if len(self.response_times[client_id]) > 100:
self.response_times[client_id].pop(0)
# 计算平均响应时间
avg_response_time = statistics.mean(
self.response_times[client_id][-10:] # 最近10个
)
# 根据响应时间调整限流
if avg_response_time > 1.0: # 响应时间超过1秒
self.current_limit = max(
self.min_limit,
self.current_limit * 0.9
)
elif avg_response_time < 0.1: # 响应时间小于100ms
self.current_limit = min(
self.max_limit,
self.current_limit * 1.1
)
return self.current_limit
# 使用示例
adaptive_limiter = AdaptiveRateLimiter()
def adaptive_api(client_id, request_data):
start_time = time.time()
# 检查当前限流
if not some_basic_limiter(client_id, adaptive_limiter.current_limit):
return {"error": "限流"}, 429
# 处理请求
result = process_request(request_data)
# 记录响应时间并调整限流
response_time = time.time() - start_time
adaptive_limiter.adjust_limit(client_id, response_time)
return result
最佳实践建议
-
选择合适的算法:
- 简单场景:固定窗口
- 精确控制:滑动窗口
- 突发流量:令牌桶
-
分布式部署:使用Redis实现集中式限流
-
监控和告警:
# 添加监控 class MonitoredRateLimiter: def __init__(self, limiter): self.limiter = limiter self.metrics = { 'allowed': 0, 'blocked': 0, 'current_rate': 0 } def is_allowed(self, client_id): allowed = self.limiter.is_allowed(client_id) if allowed: self.metrics['allowed'] += 1 else: self.metrics['blocked'] += 1 return allowed -
返回适当的HTTP状态码:
- 429 Too Many Requests
- 包含 Retry-After 头部
这些方案可以根据实际需求组合使用,实现灵活的接口限流系统。