Python接口限流案例如何限制接口请求频率

wen python案例 22

本文目录导读:

Python接口限流案例如何限制接口请求频率

  1. 基于时间窗口的简单限流
  2. 基于令牌桶算法
  3. 使用第三方库 - Flask-Limiter
  4. 使用Redis实现分布式限流
  5. 高级限流策略 - 自适应限流
  6. 最佳实践建议

我来详细介绍几种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

最佳实践建议

  1. 选择合适的算法

    • 简单场景:固定窗口
    • 精确控制:滑动窗口
    • 突发流量:令牌桶
  2. 分布式部署:使用Redis实现集中式限流

  3. 监控和告警

    # 添加监控
    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
  4. 返回适当的HTTP状态码

    • 429 Too Many Requests
    • 包含 Retry-After 头部

这些方案可以根据实际需求组合使用,实现灵活的接口限流系统。

抱歉,评论功能暂时关闭!