You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Celery子进程间API请求速率限制数据同步问题求助

解决Celery Worker间API速率限制数据同步问题

问题背景

你的Python应用基于Flask、APScheduler和Celery构建,调度器生成的任务由5个Celery Worker执行。由于调用的API存在速率限制,你尝试在父进程中创建全局RateLimits实例跟踪请求信息,但Worker只能更新各自的副本,无法实现跨进程数据共享。需要一种高效的内存缓存方案(非数据库)解决同步问题。

解决方案:使用Redis作为共享内存存储

Redis是内存型键值存储,性能高效,且天然支持跨进程/跨机器的数据共享,非常适合这类速率限制跟踪场景。Celery本身也常搭配Redis作为消息队列,无需额外引入复杂依赖。

步骤1:改造RateLimits类

将原来的本地内存存储替换为Redis,让所有Worker共享同一数据源:

import redis
from datetime import datetime
import time

class RateLimits:
    def __init__(self, redis_host="localhost", redis_port=6379, db=0):
        self.redis_client = redis.Redis(host=redis_host, port=redis_port, db=db, decode_responses=True)

    def record_request(self, api_key, url):
        """记录API请求时间"""
        redis_key = f"rate_limit:{api_key}:{url}"
        current_ts = datetime.now().timestamp()
        # 用有序集合存储请求时间,自动按时间排序
        self.redis_client.zadd(redis_key, {current_ts: current_ts})
        # 清理超出速率窗口的旧记录(示例窗口为60秒)
        self.redis_client.zremrangebyscore(redis_key, 0, current_ts - 60)

    def get_request_count(self, api_key, url, window_seconds=60):
        """获取速率窗口内的请求次数"""
        redis_key = f"rate_limit:{api_key}:{url}"
        current_ts = datetime.now().timestamp()
        return self.redis_client.zcount(redis_key, current_ts - window_seconds, current_ts)

    def wait_for_available_slot(self, api_key, url, max_requests=10, window_seconds=60):
        """等待直到有可用的请求槽位"""
        while True:
            count = self.get_request_count(api_key, url, window_seconds)
            if count < max_requests:
                break
            time.sleep(1)

步骤2:修改Celery任务代码

初始化Redis-backed的RateLimits实例,所有Worker复用同一实例访问Redis:

from .celery import app, Task
import requests
import time
from APIClass import APIClass
from RateLimits import RateLimits
from makeRequests import pollAPI

class taskBase(Task):
    def on_failure(self, *args, **kwargs):
        pass

# 初始化共享的RateLimits实例
ratelimits = RateLimits()

@app.task(base=taskBase, bind=True)
def getMoreData(self, api_keys, team):
    # 示例:执行请求前先检查速率限制
    target_url = "https://your-api-endpoint.com/team-data"
    ratelimits.wait_for_available_slot(api_keys, target_url)
    # 后续任务逻辑
    ## do stuff

@app.task(base=taskBase, bind=True)
def getData(self, api_keys):
    # 传入共享的ratelimits实例
    a = APIClass(api_key=api_keys, ratelimits=ratelimits)
    pollAPI(apis=a).run()
    print({'status':'done'})

步骤3:安装依赖并启动Redis

  1. 安装Redis服务(根据你的操作系统选择对应安装方式)
  2. 安装Python Redis客户端:
pip install redis

方案优势

  • 高效性:Redis基于内存操作,读写性能远高于磁盘数据库
  • 跨进程同步:所有Worker共享同一Redis实例,确保速率限制数据实时一致
  • 扩展性:后续如果扩展Worker到多台机器,Redis依然能支持分布式场景
  • 轻量集成:与Celery生态兼容,无需额外复杂配置

内容的提问来源于stack exchange,提问作者Stoick

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 19:08:12