Python脚本如何限流同步任务执行频率

wen python案例 35

本文目录导读:

Python脚本如何限流同步任务执行频率

  1. 使用 time.sleep() 简单限流
  2. 使用 time.sleep() + 固定速率控制
  3. 使用 time 模块的令牌桶算法
  4. 使用 decorator 装饰器实现
  5. 基于时间窗口的精确限流
  6. 使用第三方库 ratelimit
  7. 实际应用示例:API 请求限流
  8. 选择建议

在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}")

选择建议

  1. 简单场景:使用 time.sleep() 或装饰器方式
  2. 需要精确控制:使用令牌桶或滑动窗口算法
  3. 复杂生产环境:考虑使用专业的限流库如 ratelimitpyrate-limiter
  4. 分布式系统:考虑使用 Redis 等外部存储实现分布式限流

这些方法可以根据你的具体需求进行调整和组合使用。

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