首頁 / 博客 / Python 與代理開發
高並發採集如何降低封鎖率?代理池、限速與錯誤重試設計

高並發採集如何降低封鎖率?代理池、限速與錯誤重試設計

分類:Python 與代理開發 釋出時間:2026-07-22 00:30 作者:SolisProxy 瀏覽:9

高並發資料採集的瓶頸通常不在 Python 能開多少執行緒,而在排程器能否把流量限制在目標站與本地系統都能承受的範圍。真正穩定的架構會同時控制每個網域的併發、請求間隔、等待佇列、代理健康度與重試預算;當服務端發出限速或拒絕訊號時,系統應主動降載,而不是用更多請求放大問題。

先把併發預算按網域拆開

全域 worker 數不能代表安全速率。若 50 個 worker 同時遇到同一網域,即使總流量不高,也可能在短時間形成尖峰。排程器應先正規化 hostname,為每個網域維護獨立的 DomainGate:用 semaphore 限制同時請求數,用最小間隔平滑發送,再由近期錯誤率決定是否進入冷卻。

任務入口也要有界。可用 queue.Queue(maxsize=N) 搭配固定 worker,讓生產者在佇列滿時減速;同時記錄最舊任務等待時間,超過新鮮度要求便捨棄或降級。優先採用增量更新、URL 去重、內容快取與官方 API,通常比增加代理數更能降低來源站負載。產品選擇可參考住宅代理靜態資料中心代理,但代理不會取代流量治理。

代理健康度與目標站狀態要分開

代理池不應只有一個「成功/失敗」分數。連線失敗、代理驗證失敗與握手逾時可影響代理健康度;目標站回傳 429 或 503,首先反映該網域的流量或可用性狀態,不宜直接把代理判死。403、登入頁或 CAPTCHA 更應停止並人工檢查授權、條款與請求內容,不應自動換 IP 重送。

實務上可按「代理 × 目標網域 × 地區/會話」保留近期視窗或 EWMA,至少達到樣本門檻才調整權重。連續傳輸失敗的代理先冷卻,期滿只以少量探測恢復;成功率、延遲與成本則用於排序。住宅代理的 sticky 會話與輪換方式,應依會話控制說明設定,不能把每次請求都假設成新出口。

依 HTTP 訊號決定重試、降速或停止

訊號排程動作代理池動作
2xx完成任務,更新延遲增加有效成功樣本
403停止並人工檢查,不以相同條件自動重試不因單次 403 直接淘汰代理
429遵守 Retry-After;降低該網域預算記為目標限速,不當成純代理故障
500/502/503/504只對冪等請求有限重試;503 亦尊重 Retry-After與傳輸層錯誤分開統計
其他 4xx通常不自動重試,先修正請求或權限不盲目輪換

沒有 Retry-After 時,可採有上限的指數退避加 jitter,避免所有 worker 同時醒來。重試總次數與總等待時間都要設預算;一旦超出,任務進入失敗佇列或人工檢查,而不是無限循環。

一個緊湊的限速與冷卻核心

下例只展示排程策略:每個 worker 自己持有 requests.Session,排程器為每個正規化網域傳入一個共享 DomainGate。固定值只是起點,正式環境必須根據授權範圍、目標容量與監測資料調整。

import random, threading, time, requests
from collections import deque
from contextlib import contextmanager
from urllib3.exceptions import InvalidHeader
from urllib3.util import Retry
RETRYABLE, MAX_ATTEMPTS, MAX_WAIT = {429, 500, 502, 503, 504}, 3, 10
class StopRequest(Exception): pass
class DomainGate:
    def __init__(self, concurrency=2, min_gap=0.75):
        self.slots = threading.BoundedSemaphore(concurrency)
        self.min_gap, self.lock = min_gap, threading.Lock()
        self.next_at = self.cooldown_until = 0.0
        self.recent = deque(maxlen=10)
    @contextmanager
    def permit(self):
        self.slots.acquire()
        try:
            with self.lock:
                now = time.monotonic()
                if self.cooldown_until and now >= self.cooldown_until:
                    self.cooldown_until = 0.0
                    self.recent.clear()
                if now < self.cooldown_until:
                    raise StopRequest("domain is cooling down")
                delay = max(0.0, self.next_at - now)
                self.next_at = max(now, self.next_at) + self.min_gap
            time.sleep(delay)
            yield
        finally:
            self.slots.release()
    def record(self, healthy):
        with self.lock:
            self.recent.append(healthy)
            if len(self.recent) == 10 and self.recent.count(False) >= 5:
                self.cooldown_until = max(self.cooldown_until, time.monotonic() + 30)
def retry_delay(response, attempt):
    raw = response.headers.get("Retry-After")
    if raw is None:
        base = 0.5 * (2 ** attempt)
        return min(MAX_WAIT, base + random.uniform(0, 0.25))
    try:
        delay = Retry(total=0).parse_retry_after(raw)
    except InvalidHeader:
        raise StopRequest("invalid Retry-After") from None
    if delay > MAX_WAIT:
        raise StopRequest("Retry-After exceeds local budget")
    return delay
def fetch(session, url, gate):
    for attempt in range(MAX_ATTEMPTS):
        with gate.permit():
            response = session.get(url, timeout=(3.05, 20), allow_redirects=False)
            try:
                status = response.status_code
                gate.record(status not in RETRYABLE and status != 403)
                if status == 403:
                    raise StopRequest("403 requires review")
                if 300 <= status < 400:
                    raise StopRequest("redirect requires review")
                if status not in RETRYABLE:
                    response.raise_for_status()
                    return response.content
                if attempt + 1 == MAX_ATTEMPTS:
                    raise StopRequest("retry budget exhausted")
                delay = retry_delay(response, attempt)
            finally:
                response.close()
        time.sleep(delay)
    raise AssertionError("unreachable")

這段程式不會自動重試連線錯誤,也不包含完整工作佇列、代理選擇器與回應正文大小限制。正式採集器應使用串流讀取及累計位元組上限;若需要整體 deadline,也要在單次 connect/read timeout 之外另行實作。Session 不應在多個 worker 間任意共享,代理 URL 與帳密必須來自密鑰管理器,日誌不得輸出完整 URL、查詢字串或憑證。

熔斷器不要只有「暫停 30 秒」

範例的近期視窗只是冷卻閘門。完整熔斷器應有 closed、open、half-open 三個狀態:正常時收集錯誤率;達到樣本與門檻後打開,暫停新請求;冷卻期結束後只允許少量探測,成功才逐步恢復。若半開時一次放行全部積壓任務,舊尖峰會立刻重現。

熔斷與有界佇列要一起運作。監測佇列深度與最舊任務年齡,對過期的全量任務做 load shedding,保留高價值增量工作;當 429 比率、p95 延遲或有效成功率惡化時,同步降低該網域的 concurrency 與速率,而不是只替換代理。

用可量化指標判斷是否真的更穩

  • 有效成功率:通過內容契約驗證的資料筆數/唯一任務數,而非只有 HTTP 200。
  • 限制與拒絕率:429、403、登入頁與 CAPTCHA 分開計數,按網域和時間窗觀察。
  • 延遲與佇列:p50/p95/p99 回應時間、佇列深度及最舊任務年齡。
  • 重試放大:總嘗試次數/唯一任務數;數值上升代表重試正放大流量。
  • 代理效率:可用代理比例、冷卻數量,以及每 GB 或每單位成本取得的有效資料。

調整時一次只改一個變數,保留基準線與回滾條件。只在已授權目標上測試,遵守網站條款、API/robots 政策、個資與所在地法律;403、登入/CAPTCHA、連續逾時或預算門檻一旦觸發就停止並人工檢查,不以換身分或提高流量規避限制。

官方參考