第 3 篇:频率控制与限流——别因为太贪心先被封
·
一、为什么一加多线程就被封?
常见路径:for 循环太慢 → 上多线程/协程 → QPS 飙高 → 很快出现:
429 Too Many Requests403 Forbidden- 200 但内容是风控页/验证码提示
站点往往按 IP / 账号 / 接口 做时间窗口统计。你一分钟打几百次同一接口,和真人行为差太远,先封频率是最便宜的反爬手段。
结论:对抗频率层,目标不是「最快」,而是 在可接受时间内平稳拉完数据。
二、三个词:限速、抖动、退避
| 概念 | 含义 |
|---|---|
| 限速 | 控制单位时间内请求次数(如每域名 2 QPS) |
| 抖动 | 间隔加随机数,避免固定 1.000s、2.000s 的「机器节奏」 |
| 退避 | 遇 429/风控先降温(sleep、降 QPS),不要死循环重试 |
三、最简单:按固定间隔 + 抖动
适合单进程、快速止血:
import time
import random
class RateLimiter:
def __init__(self, min_interval: float = 1.5):
self.min_interval = min_interval
self._last_ts = 0.0
def wait(self):
now = time.time()
gap = now - self._last_ts
if gap < self.min_interval:
time.sleep(self.min_interval - gap + random.uniform(0, 0.3))
self._last_ts = time.time()
用法:每次 requests 前调用 limiter.wait()。
四、令牌桶:控制平均 QPS + 允许小幅突发
比固定间隔更灵活:长期平均不超过 rate,短时允许 capacity 个突发。
import time
import threading
class TokenBucket:
def __init__(self, rate: float, capacity: int):
"""
rate: 每秒补充令牌数(期望 QPS)
capacity: 桶容量(最大突发)
"""
self._rate = max(rate, 1e-6)
self._capacity = capacity
self._tokens = float(capacity)
self._last_ts = time.time()
self._lock = threading.Lock()
def _refill(self):
now = time.time()
delta = now - self._last_ts
self._last_ts = now
self._tokens = min(self._capacity, self._tokens + delta * self._rate)
def consume(self, tokens: float = 1.0):
while True:
with self._lock:
self._refill()
if self._tokens >= tokens:
self._tokens -= tokens
return
need = tokens - self._tokens
wait_time = need / self._rate
time.sleep(min(wait_time, 2.0))
多线程共用一个桶,整体 QPS 会被压在约 rate 附近。
五、多站点分桶:A 站限流别拖死 B 站
按 域名 各放一个桶,避免一个站打满连累别的站。
from urllib.parse import urlparse
class MultiSiteLimiter:
def __init__(self, default_rate=1.5, default_capacity=5):
self.default_rate = default_rate
self.default_capacity = default_capacity
self._buckets = {}
self._lock = threading.Lock()
def _host(self, url: str) -> str:
h = urlparse(url).netloc
return h if h else "default"
def wait(self, url: str, rate: float = None, capacity: int = None):
host = self._host(url)
with self._lock:
if host not in self._buckets:
r = rate if rate is not None else self.default_rate
c = capacity if capacity is not None else self.default_capacity
self._buckets[host] = TokenBucket(r, c)
self._buckets[host].consume(1)
可选:为特定域名单独配置,例如 per_host={"api.xxx.com": (0.5, 3)},在 wait 里查表即可(可自行扩展)。
六、限流参数怎么选(干货表)
| 场景 | 建议 rate(QPS) | capacity | 说明 |
|---|---|---|---|
| 个人博客、小站 | 0.3~1 | 3~5 | 宁可慢 |
| 一般内容站 | 1~2 | 5~8 | 先看 robots/用户协议 |
| 明显有风控的站 | 0.2~0.5 | 2~4 | 配合更长冷却 |
| 需登录的接口 | 每账号再限一层 | - | 账号池各带一桶 |
原则:从保守开始,根据成功率再慢慢调高;一旦出现 429/风控关键词,先减半 rate,别硬顶。
七、响应分类 + 遇 429/风控自动降温
def is_block_page(html: str) -> bool:
if not html:
return False
keys = ["验证码", "安全验证", "访问过于频繁", "访问受限", "verify you are human", "robot check"]
low = html.lower()
return any(k.lower() in low for k in keys)
def classify_response(status_code: int, text: str) -> str:
if status_code == 429:
return "rate_limited"
if status_code == 403:
return "forbidden"
if status_code == 200:
if not (text or "").strip():
return "empty"
if is_block_page(text):
return "blocked"
return "ok"
return "unknown"
自适应降速:某 host 连续多次 rate_limited/blocked,对该 host 临时把 rate *= 0.5 或进入冷却 60s(见下节完整 Worker)。
八、完整示例:队列 + 多 Worker + 多站限流 + 429 降速
下面是一段可直接改 URL 跑通结构的完整脚本(请替换为你自己的合法目标与频率)。
import queue
import random
import threading
import time
import requests
from dataclasses import dataclass
from urllib.parse import urlparse
# ========== TokenBucket & MultiSiteLimiter(同上,此处合并) ==========
class TokenBucket:
def __init__(self, rate: float, capacity: int):
self._rate = max(rate, 1e-6)
self._capacity = capacity
self._tokens = float(capacity)
self._last_ts = time.time()
self._lock = threading.Lock()
def _refill(self):
now = time.time()
delta = now - self._last_ts
self._last_ts = now
self._tokens = min(self._capacity, self._tokens + delta * self._rate)
def consume(self, tokens: float = 1.0):
while True:
with self._lock:
self._refill()
if self._tokens >= tokens:
self._tokens -= tokens
return
need = tokens - self._tokens
wait_time = need / self._rate
time.sleep(min(wait_time, 2.0))
def slow_down(self, factor: float = 0.5):
with self._lock:
self._rate = max(self._rate * factor, 0.05)
class HostRateState:
def __init__(self, rate: float, capacity: int):
self.bucket = TokenBucket(rate, capacity)
self.bad_streak = 0
self.cooldown_until = 0.0
self.lock = threading.Lock()
def on_bad(self):
with self.lock:
self.bad_streak += 1
if self.bad_streak >= 3:
self.bucket.slow_down(0.5)
self.bad_streak = 0
self.cooldown_until = time.time() + 30
def on_ok(self):
with self.lock:
self.bad_streak = 0
def in_cooldown(self) -> bool:
return time.time() < self.cooldown_until
class MultiSiteLimiter:
def __init__(self, default_rate=1.0, default_capacity=5):
self.default_rate = default_rate
self.default_capacity = default_capacity
self._hosts = {}
self._lock = threading.Lock()
def _get_state(self, url: str) -> HostRateState:
host = urlparse(url).netloc or "default"
with self._lock:
if host not in self._hosts:
self._hosts[host] = HostRateState(self.default_rate, self.default_capacity)
return self._hosts[host]
def wait(self, url: str):
st = self._get_state(url)
if st.in_cooldown():
time.sleep(min(st.cooldown_until - time.time(), 10) + random.uniform(0, 1))
st.bucket.consume(1)
def report(self, url: str, kind: str):
st = self._get_state(url)
if kind in ("rate_limited", "blocked", "forbidden"):
st.on_bad()
elif kind == "ok":
st.on_ok()
def classify_response(status_code: int, text: str) -> str:
if status_code == 429:
return "rate_limited"
if status_code == 403:
return "forbidden"
if status_code == 200:
t = (text or "").strip()
if not t:
return "empty"
low = t.lower()
for k in ["验证码", "安全验证", "访问过于频繁", "访问受限", "verify you are human"]:
if k.lower() in low:
return "blocked"
return "ok"
return "unknown"
@dataclass
class Task:
url: str
def worker(name: str, q: queue.Queue, limiter: MultiSiteLimiter, session: requests.Session):
headers = {
"User-Agent": "Mozilla/5.0 (Windows NT 10.0; Win64; x64) AppleWebKit/537.36 (KHTML, like Gecko) Chrome/122.0.0.0 Safari/537.36",
"Accept-Language": "zh-CN,zh;q=0.9",
}
while True:
try:
task = q.get(timeout=5)
except queue.Empty:
print(f"[{name}] 队列空,退出")
return
try:
limiter.wait(task.url)
r = session.get(task.url, headers=headers, timeout=15)
kind = classify_response(r.status_code, r.text)
limiter.report(task.url, kind)
print(f"[{name}] {task.url[:60]}... -> {kind} {r.status_code} len={len(r.text)}")
if kind in ("rate_limited", "blocked"):
time.sleep(5 + random.uniform(0, 3))
except Exception as e:
print(f"[{name}] err {e}")
finally:
q.task_done()
if __name__ == "__main__":
# 示例:请换成你有权访问的 URL;禁止对未授权站点高频访问
urls = [
"https://httpbin.org/get",
"https://httpbin.org/status/429",
] * 5
task_q = queue.Queue()
for u in urls:
task_q.put(Task(url=u))
limiter = MultiSiteLimiter(default_rate=1.5, default_capacity=4)
sess = requests.Session()
threads = []
for i in range(3):
t = threading.Thread(target=worker, args=(f"W{i}", task_q, limiter, sess), daemon=True)
t.start()
threads.append(t)
task_q.join()
print("done")
说明:httpbin.org/status/429 用于本地看 429 → report → slow_down 行为;真实站点请降低 rate、减少并发。
九、排查清单
- 是否全站共用一个限流? → 改为按 host 分桶。
- 间隔是否完全固定? → 加
random.uniform抖动。 - 429 后是否立即重试? → 先 sleep 再试,并降低 QPS。
- 多线程是否每线程各自限流? → 应共享同一站点的桶。
- 是否忽略 200 风控页? → 用
classify_response识别关键词。 - 登录态是否多线程乱用? → 会话与账号池见系列第 4 篇。
十、系列导航 & 关注说明
- 第 1 篇:反爬全景图
- 第 2 篇:请求层(Header、Session)
- 第 3 篇(本篇):频率、令牌桶、多站分桶、队列 Worker、429 降速
- 第 4 篇:会话与 Cookie / 账号池
- 第 5 篇:Playwright + requests
更多推荐
所有评论(0)