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
相关产品推荐
相关产品推荐

