本文目录导读:

- 使用
time.sleep()简单限流 - 使用
time.sleep()+ 固定速率控制 - 使用
time模块的令牌桶算法 - 使用
decorator装饰器实现 - 基于时间窗口的精确限流
- 使用第三方库
ratelimit - 实际应用示例:API 请求限流
- 选择建议
在Python中限制同步任务的执行频率(限流/频率控制)有多种方法,以下是几种常见的实现方式:
使用 time.sleep() 简单限流
最简单的限流方法,控制每次任务执行后的等待时间。
import time
from datetime import datetime
def task():
print(f"执行任务: {datetime.now()}")
# 你的任务代码
time.sleep(1) # 模拟任务执行
def rate_limited_loop(interval=1.0):
"""每 interval 秒执行一次任务"""
while True:
start_time = time.time()
task()
elapsed = time.time() - start_time
sleep_time = interval - elapsed
if sleep_time > 0:
time.sleep(sleep_time)
# 每秒执行不超过1次
rate_limited_loop(1.0)
使用 time.sleep() + 固定速率控制
import time
class RateLimiter:
def __init__(self, max_calls, period=1.0):
"""
max_calls: 在 period 秒内允许的最大调用次数
period: 时间窗口(秒)
"""
self.max_calls = max_calls
self.period = period
self.calls = []
def wait_if_needed(self):
now = time.time()
# 移除过期记录
self.calls = [t for t in self.calls if now - t < self.period]
if len(self.calls) >= self.max_calls:
# 等待到最早的那个调用过期
wait_time = self.period - (now - self.calls[0])
print(f"限流中,等待 {wait_time:.2f} 秒")
time.sleep(wait_time)
self.calls.append(time.time())
# 使用示例
limiter = RateLimiter(max_calls=5, period=1.0)
for i in range(10):
limiter.wait_if_needed()
print(f"执行任务 {i}")
time.sleep(0.1) # 模拟任务耗时
使用 time 模块的令牌桶算法
import time
import threading
class TokenBucket:
def __init__(self, rate, capacity):
"""
rate: 每秒生成的令牌数
capacity: 桶的最大容量(突发流量能力)
"""
self.rate = rate
self.capacity = capacity
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 wait_and_consume(self, tokens=1):
while True:
if self.consume(tokens):
return
# 计算需要等待的时间
with self.lock:
wait_time = (tokens - self.tokens) / self.rate
time.sleep(min(wait_time, 0.1)) # 最多等待0.1秒
def _refill(self):
now = time.time()
elapsed = now - self.last_refill
self.tokens = min(self.capacity, self.tokens + elapsed * self.rate)
self.last_refill = now
# 使用示例
bucket = TokenBucket(rate=10, capacity=10) # 每秒10个令牌
for i in range(20):
bucket.wait_and_consume()
print(f"执行任务 {i}")
使用 decorator 装饰器实现
import time
import functools
def rate_limit(max_calls, period=1.0):
"""装饰器:限制函数调用频率"""
def decorator(func):
calls = []
@functools.wraps(func)
def wrapper(*args, **kwargs):
nonlocal calls
now = time.time()
calls = [t for t in calls if now - t < period]
if len(calls) >= max_calls:
wait_time = period - (now - calls[0])
print(f"限流中,等待 {wait_time:.2f} 秒")
time.sleep(wait_time)
calls = [t for t in calls if now - t < period]
calls.append(time.time())
return func(*args, **kwargs)
return wrapper
return decorator
# 使用示例
@rate_limit(max_calls=3, period=1.0) # 每秒最多3次
def my_task(task_id):
print(f"执行任务 {task_id}")
for i in range(10):
my_task(i)
time.sleep(0.1) # 模拟其他操作
基于时间窗口的精确限流
import time
from collections import deque
class SlidingWindow:
def __init__(self, max_calls, window_size):
"""
max_calls: 窗口内最大调用次数
window_size: 窗口大小(秒)
"""
self.max_calls = max_calls
self.window_size = window_size
self.timestamps = deque()
def allow_request(self):
now = time.time()
# 移除窗口外的记录
while self.timestamps and now - self.timestamps[0] > self.window_size:
self.timestamps.popleft()
if len(self.timestamps) < self.max_calls:
self.timestamps.append(now)
return True
return False
def wait_for_slot(self):
while not self.allow_request():
if self.timestamps:
# 计算需要等待的时间
wait_time = self.window_size - (time.time() - self.timestamps[0])
time.sleep(min(wait_time, 0.1))
else:
time.sleep(0.1)
# 使用示例
limiter = SlidingWindow(max_calls=5, window_size=1)
for i in range(10):
limiter.wait_for_slot()
print(f"执行任务 {i}")
使用第三方库 ratelimit
pip install ratelimit
from ratelimit import limits, sleep_and_retry
import time
@sleep_and_retry
@limits(calls=5, period=1) # 每秒最多5次
def my_task():
print(f"执行任务 {time.time()}")
# 测试
for i in range(15):
my_task()
实际应用示例:API 请求限流
import time
import requests
from urllib.parse import urlencode
class APIRateLimiter:
def __init__(self, max_requests_per_second=10):
self.max_requests = max_requests_per_second
self.last_request_time = 0
self.request_times = []
def wait_and_request(self, url, params=None, method='GET', **kwargs):
"""限流发送HTTP请求"""
self._wait_for_slot()
if method == 'GET':
response = requests.get(url, params=params, **kwargs)
elif method == 'POST':
response = requests.post(url, data=params, **kwargs)
else:
raise ValueError(f"Unsupported method: {method}")
self.request_times.append(time.time())
return response
def _wait_for_slot(self):
now = time.time()
# 清理超过1秒的记录
self.request_times = [t for t in self.request_times if now - t < 1.0]
while len(self.request_times) >= self.max_requests:
# 等待最早请求过期
wait_time = 1.0 - (now - self.request_times[0])
if wait_time > 0:
time.sleep(wait_time)
self.request_times = [t for t in self.request_times if time.time() - t < 1.0]
# 使用示例
api_limiter = APIRateLimiter(max_requests_per_second=5)
urls = [
"https://api.example.com/endpoint1",
"https://api.example.com/endpoint2",
# ... 更多URL
]
for url in urls:
try:
response = api_limiter.wait_and_request(url)
print(f"成功请求: {url}, 状态码: {response.status_code}")
except Exception as e:
print(f"请求失败: {url}, 错误: {e}")
选择建议
- 简单场景:使用
time.sleep()或装饰器方式 - 需要精确控制:使用令牌桶或滑动窗口算法
- 复杂生产环境:考虑使用专业的限流库如
ratelimit、pyrate-limiter - 分布式系统:考虑使用 Redis 等外部存储实现分布式限流
这些方法可以根据你的具体需求进行调整和组合使用。