
本文介绍如何仅用python标准库构建支持多域名独立速率限制的并发http请求调度系统,核心在于为每个域名维护独立的等待队列与时间感知的唤醒机制,避免全局阻塞和轮询浪费。
本文介绍如何仅用python标准库构建支持多域名独立速率限制的并发http请求调度系统,核心在于为每个域名维护独立的等待队列与时间感知的唤醒机制,避免全局阻塞和轮询浪费。
在仅依赖 Python 标准库的前提下,实现多域名、高并发、带动态速率限制的 Web 请求调度器,关键不在于复现复杂算法(如漏桶或令牌桶),而在于设计一种可伸缩、无忙等、保序(可选)、低开销的调度结构。核心挑战是:当某个域名因限流进入暂停期时,其他域名的任务必须能立即被消费,而非被卡在共享队列中被动等待。
✅ 推荐方案:threading.Condition + 按域名分片的优先队列(heapq)+ 全局调度线程
该方案完全使用 threading, queue, heapq, time, urllib.request 等标准模块,无需第三方依赖,兼顾效率、清晰性与可维护性。
图片提示词生成器?不止如此。 马甲系统 —— 把脑海中的画面,翻译成AI能理解的专业表达。 用得越多,它越懂你:首次需要多问几句确认方向,用久了几乎一说就懂。 用得越多,它越快:缓存机制让后续对话越来越省。 RAG进化:成功案例持续入库,越跑越聪明。 输入「新手指南」查看完整功能介绍
? 数据结构设计
- 每个域名对应一个最小堆(heapq):存储待发请求,按「最早允许发送时间」排序(即 next_allowed_at: float)。
- 全局 domain_queues: dict[str, list]:键为域名,值为该域名的最小堆(list),配合 heapq 原地操作。
- 全局 condition: threading.Condition:用于协调“有新任务就绪”或“某域名解除限流”的通知。
- 全局 backoff_end: dict[str, float]:记录各域名当前限流截止时间戳(秒级浮点数),用于判断是否可出队。
✅ 优势:heapq 支持 O(1) 查看堆顶(最早可执行时间),O(log n) 插入/弹出;Condition 实现零轮询唤醒;无额外线程池爆炸风险。
⚙️ 调度逻辑(核心伪代码)
import heapq
import threading
import time
import urllib.request
from typing import NamedTuple, Dict, List, Optional
class DownloadItem(NamedTuple):
url: str
domain: str
# 其他元数据...
class RateLimitedScheduler:
def __init__(self, max_per_domain: int = 5, per_second: float = 1.0):
self.max_per_domain = max_per_domain
self.per_second = per_second # 即每秒最多 1 请求 → 间隔 ≥ 1.0s
self.domain_queues: Dict[str, List] = {}
self.backoff_end: Dict[str, float] = {}
self.condition = threading.Condition()
self.running = True
def add(self, item: DownloadItem):
# 计算该域名下一个允许时间(考虑当前限流状态)
now = time.time()
domain = item.domain
next_allowed = max(now, self.backoff_end.get(domain, 0))
# 将任务推入对应域名的最小堆:(next_allowed_time, item)
heap = self.domain_queues.setdefault(domain, [])
heapq.heappush(heap, (next_allowed, item))
with self.condition:
self.condition.notify_all() # 唤醒所有等待线程
def _get_next_ready_item(self) -> Optional[DownloadItem]:
now = time.time()
candidates = []
for domain, heap in self.domain_queues.items():
if not heap:
continue
# 查看堆顶(最早允许时间)
earliest_time, item = heap[0]
if earliest_time time.time():
next_wakeup = min(next_wakeup, end_time)
if next_wakeup == float('inf'):
# 所有限流已过期,但暂无任务 → 等待新任务
self.condition.wait(timeout=0.1)
else:
# 等待至最近 backoff 结束
timeout = max(0.01, next_wakeup - time.time())
self.condition.wait(timeout=timeout)
if item is not None:
try:
# 执行 HTTP 请求(示例)
with urllib.request.urlopen(item.url, timeout=10) as resp:
data = resp.read()
# 请求成功 → 更新该域名下次允许时间
now = time.time()
next_allowed = now + (1.0 / self.per_second)
self.backoff_end[item.domain] = next_allowed
except Exception as e:
# 失败时可延长 backoff(如指数退避)
self.backoff_end[item.domain] = time.time() + 2.0
finally:
# 无论成败,都通知可能有新任务可调度
with self.condition:
self.condition.notify_all()
# 使用示例
scheduler = RateLimitedScheduler(max_per_domain=3, per_second=0.5) # 每域名每2秒最多1次
# 启动多个工作线程
threads = []
for i in range(5):
t = threading.Thread(target=scheduler.worker, daemon=True)
t.start()
threads.append(t)
# 添加任务(混合域名)
scheduler.add(DownloadItem("https://example.com/a", "example.com"))
scheduler.add(DownloadItem("https://google.com/1", "google.com"))
scheduler.add(DownloadItem("https://example.com/b", "example.com"))
scheduler.add(DownloadItem("https://github.com/x", "github.com"))
# 主线程保持运行(实际中可 join 或监听信号)
try:
while True:
time.sleep(3600)
except KeyboardInterrupt:
scheduler.running = False
⚠️ 关键注意事项
- 非严格 FIFO,但强时间公平性:同一域名内任务按计划时间顺序执行,不同域名间以「最早可执行时间」为优先级,比随机轮询更合理。
- backoff 动态更新:每次请求后重置 backoff_end[domain],失败时可主动延长(如 +2.0 秒),天然支持退避策略。
- 无 busy-wait:Condition.wait(timeout=...) 是内核级等待,CPU 零占用;唤醒时机精准(新任务到达 或 backoff 到期)。
- 内存安全:所有对 domain_queues 和 backoff_end 的读写均受 Condition 保护,无需额外锁。
- 可扩展性:域名数量不受线程数限制(区别于“每域一线程”反模式),仅增加堆内存开销。
? 总结
放弃“单队列+轮询”或“多队列+sleep polling”的思路,转而采用 heapq 维护时间优先级 + threading.Condition 实现事件驱动唤醒,是标准库下最简洁、高效、可维护的解法。它将“何时能发请求”这一动态约束编码进数据结构本身,并由调度器统一协调,彻底解耦域名限流逻辑与并发执行模型。无需 asyncio、无需外部包,一行 pip install 都不需要——恰是生产环境对轻量、可控、审计友好的终极诉求。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!










