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

Python多脚本调用限频REST API遇503错误,求统一调度方案

问题描述

单脚本调用限制每秒1次的REST API时,以下代码可正常运行:

for something in ListOfSomethings:
    conn = http.client.HTTPSConnection("example.com") 
    url = f"/api/somethings/{ListOfSomethings['id']}"
    conn.request("GET", url, payload, headers)
    res = conn.getresponse()
    doMagic(res)
    time.sleep(1)

但20多个不同调度频率的脚本同时调用时,触发HTTP 503 Rate Limit Exceeded错误。直接延长sleep时间效率太低,因此想编写常驻的MainApiHandler.py作为请求队列管理器,每秒发送一个请求并返回响应,其他客户端脚本可将请求加入队列并异步等待响应。但测试现有方案时多实例仍并行执行未排队,请求时间重叠,需要实现按时间顺序统一调度所有请求的逻辑,期望时序如下:

t0 - nothing happens
--> A wants to send 5 requests
t1 - A1 sent & received
t2 - A2 sent & received
--> B wants to send 3 requests
t3 - A3 sent & received
t4 - A4 sent & received
--> C wants to send 1 requests
t5 - A5 sent & received
t6 - B1 sent & received
t7 - B2 sent & received
t8 - B3 sent & received
t9 - C1 sent & received
t10 - nothing happens

求实现方案或其他解决思路。


实现方案与替代思路

一、常驻队列管理器实现

1. 基于Redis的分布式队列方案

适合跨独立脚本的全局调度,Redis提供可靠的队列存储和结果缓存:

服务端(MainApiHandler.py)

常驻进程负责从队列取任务、控制请求频率、执行API调用并存储结果:

import time
import redis
import http.client
import json

# 初始化Redis连接
r = redis.Redis(host='localhost', port=6379, db=0, decode_responses=True)
QUEUE_KEY = 'api_request_queue'
RESULT_PREFIX = 'api_result_'

def process_api_request(request_data):
    """执行API调用并返回响应数据"""
    req_data = json.loads(request_data)
    conn = http.client.HTTPSConnection("example.com")
    conn.request(
        method=req_data.get('method', 'GET'),
        url=req_data['url'],
        body=req_data.get('payload'),
        headers=req_data['headers']
    )
    res = conn.getresponse()
    return json.dumps({
        'status': res.status,
        'content': res.read().decode('utf-8'),
        'headers': dict(res.getheaders())
    })

# 记录上次请求时间,控制每秒1次的频率
last_request_time = 0

while True:
    # 阻塞式获取队列任务(无任务时等待)
    _, task_data = r.blpop(QUEUE_KEY)
    task_id, request_data = json.loads(task_data)
    
    # 控制请求频率
    current_time = time.time()
    if current_time - last_request_time < 1:
        time.sleep(1 - (current_time - last_request_time))
    
    # 处理请求并存储结果(设置5分钟过期,避免内存占用)
    result = process_api_request(request_data)
    r.setex(RESULT_PREFIX + task_id, 300, result)
    
    # 更新上次请求时间
    last_request_time = time.time()

客户端脚本

提交请求到队列,轮询获取结果:

import uuid
import redis
import json
import time

r = redis.Redis(host='localhost', port=6379, db=0, decode_responses=True)
QUEUE_KEY = 'api_request_queue'
RESULT_PREFIX = 'api_result_'

def submit_api_request(url, headers, method='GET', payload=None):
    """提交API请求到队列,等待并返回结果"""
    task_id = str(uuid.uuid4())
    request_data = json.dumps({
        'url': url,
        'headers': headers,
        'method': method,
        'payload': payload
    })
    # 将任务ID和请求数据序列化后推入队列
    r.rpush(QUEUE_KEY, json.dumps([task_id, request_data]))
    
    # 轮询获取结果
    while True:
        result = r.get(RESULT_PREFIX + task_id)
        if result:
            r.delete(RESULT_PREFIX + task_id)
            return json.loads(result)
        time.sleep(0.1)

# 客户端调用示例
headers = {'Authorization': 'Bearer YOUR_TOKEN'}
for something in ListOfSomethings:
    url = f"/api/somethings/{something['id']}"
    api_result = submit_api_request(url, headers)
    doMagic(api_result)

2. 基于FastAPI的HTTP队列服务

如果不想让客户端直接操作Redis,可以封装HTTP接口,降低客户端依赖:

服务端(MainApiHandler.py)

from fastapi import FastAPI, BackgroundTasks
import asyncio
import time
import http.client
import json
from uuid import uuid4
from pydantic import BaseModel

app = FastAPI()
request_queue = asyncio.Queue()
results = {}
last_request_time = 0

class ApiRequest(BaseModel):
    url: str
    headers: dict
    method: str = "GET"
    payload: str = ""

async def process_queue():
    """后台任务:处理队列中的请求,控制每秒1次频率"""
    global last_request_time
    while True:
        task_id, req = await request_queue.get()
        
        # 控制请求间隔
        now = time.time()
        if now - last_request_time < 1:
            await asyncio.sleep(1 - (now - last_request_time))
        
        # 执行API调用
        conn = http.client.HTTPSConnection("example.com")
        conn.request(req.method, req.url, req.payload, req.headers)
        res = conn.getresponse()
        response_data = json.dumps({
            'status': res.status,
            'content': res.read().decode('utf-8'),
            'headers': dict(res.getheaders())
        })
        
        # 存储结果
        results[task_id] = response_data
        last_request_time = time.time()
        request_queue.task_done()

@app.on_event("startup")
async def startup():
    """启动时启动队列处理任务"""
    asyncio.create_task(process_queue())

@app.post("/submit")
async def submit_request(req: ApiRequest):
    """提交API请求接口"""
    task_id = str(uuid4())
    await request_queue.put((task_id, req))
    return {"task_id": task_id}

@app.get("/result/{task_id}")
async def get_result(task_id: str):
    """查询请求结果接口"""
    if task_id in results:
        result = json.loads(results.pop(task_id))
        return result
    return {"status": "pending"}

客户端调用示例

import requests
import time

def submit_api_request(url, headers, method='GET', payload=None):
    req_data = {
        "url": url,
        "headers": headers,
        "method": method,
        "payload": payload or ""
    }
    # 提交请求
    resp = requests.post("http://localhost:8000/submit", json=req_data)
    task_id = resp.json()['task_id']
    
    # 轮询结果
    while True:
        result_resp = requests.get(f"http://localhost:8000/result/{task_id}")
        result = result_resp.json()
        if result.get('status') != 'pending':
            return result
        time.sleep(0.1)

# 调用示例
headers = {'Authorization': 'Bearer YOUR_TOKEN'}
for something in ListOfSomethings:
    url = f"/api/somethings/{something['id']}"
    api_result = submit_api_request(url, headers)
    doMagic(api_result)

二、替代思路

1. 分布式锁控制请求频率

无需常驻服务,所有客户端通过Redis分布式锁竞争请求权限:

import time
import redis
import http.client

r = redis.Redis(host='localhost', port=6379, db=0)
LOCK_KEY = 'api_rate_limit_lock'
LOCK_EXPIRE = 1.1  # 锁过期时间略大于1秒,避免客户端崩溃导致锁无法释放

def send_api_request(url, headers, payload=None):
    while True:
        # 尝试获取锁
        if r.set(LOCK_KEY, 'locked', nx=True, ex=LOCK_EXPIRE):
            # 拿到锁,发送请求
            conn = http.client.HTTPSConnection("example.com")
            conn.request("GET", url, payload, headers)
            res = conn.getresponse()
            result = {
                'status': res.status,
                'content': res.read().decode('utf-8')
            }
            # 主动释放锁(可选,因为锁会自动过期)
            r.delete(LOCK_KEY)
            return result
        # 没拿到锁,等待后重试
        time.sleep(0.1)

缺点:请求顺序无法严格保证,可能出现后提交的请求先执行的情况。

2. 分布式速率限制工具

使用pyrate-limiter配合Redis实现全局速率限制,所有客户端共享同一规则:

from pyrate_limiter import Duration, RequestRate, Limiter
from pyrate_limiter.redis import RedisBucket
import redis
import http.client

# 初始化Redis存储
r = redis.Redis(host='localhost', port=6379, db=0)
bucket = RedisBucket(r, 'api_rate_limit')
# 设置每秒1次请求的规则
rate = RequestRate(1, Duration.SECOND)
limiter = Limiter(rate, bucket)

def send_api_request(url, headers, payload=None):
    # 等待获取速率许可
    with limiter.ratelimit('global_api_call', delay=True):
        conn = http.client.HTTPSConnection("example.com")
        conn.request("GET", url, payload, headers)
        res = conn.getresponse()
        return {
            'status': res.status,
            'content': res.read().decode('utf-8')
        }

优点:无需自己实现队列,依赖成熟库;缺点:同样无法严格保证请求的提交顺序。


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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 14:35:00