一、为什么一加多线程就被封?

常见路径:for 循环太慢 → 上多线程/协程 → QPS 飙高 → 很快出现:

  • 429 Too Many Requests
  • 403 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、减少并发。


九、排查清单

  1. 是否全站共用一个限流? → 改为按 host 分桶。
  2. 间隔是否完全固定? → 加 random.uniform 抖动。
  3. 429 后是否立即重试? → 先 sleep 再试,并降低 QPS。
  4. 多线程是否每线程各自限流? → 应共享同一站点的桶。
  5. 是否忽略 200 风控页? → 用 classify_response 识别关键词。
  6. 登录态是否多线程乱用? → 会话与账号池见系列第 4 篇。

十、系列导航 & 关注说明

  • 第 1 篇:反爬全景图
  • 第 2 篇:请求层(Header、Session)
  • 第 3 篇(本篇):频率、令牌桶、多站分桶、队列 Worker、429 降速
  • 第 4 篇:会话与 Cookie / 账号池
  • 第 5 篇:Playwright + requests
Logo

腾讯云面向开发者汇聚海量精品云计算使用和开发经验,营造开放的云计算技术生态圈。

更多推荐