diff --git a/core/douyu/proxy_manager.py b/core/douyu/proxy_manager.py index d59d273..aaad9c7 100644 --- a/core/douyu/proxy_manager.py +++ b/core/douyu/proxy_manager.py @@ -14,6 +14,11 @@ from .proxy_whitelist import DouyuWhitelistSyncer class ProxyManager: """代理管理器(带已验证代理池缓存,批次内共享复用)""" + # 代理失败软隔离配置 + _COOLDOWN_BASE = 30 # 基础冷却秒数(失败1次) + _COOLDOWN_MAX = 120 # 最大冷却秒数 + _MAX_FAIL_COUNT = 3 # 连续失败此次数后永久移出池 + def __init__(self, api_url: str = "", whitelist_uid: str = "", whitelist_ukey: str = ""): self.api_url = api_url self.whitelist_uid = whitelist_uid @@ -24,6 +29,8 @@ class ProxyManager: self._verified_pool: dict[str, float] = {} # 正在使用中的代理(取走但未归还),避免并发账号用同一个代理 self._in_use: set[str] = set() + # 代理失败记录: {proxy_url: {"count": 失败次数, "cooldown_until": 冷却到期时间戳}} + self._fail_records: dict[str, dict] = {} self._pool_ttl = 90 self._fetching = False self._whitelist_syncer = ( @@ -32,23 +39,35 @@ class ProxyManager: else None ) + def _is_cooling_down(self, proxy_url: str) -> bool: + """检查代理是否在冷却期内(调用前需持有锁)。""" + record = self._fail_records.get(proxy_url) + if not record: + return False + return time.time() < record.get("cooldown_until", 0) + def _pick_from_pool_locked(self) -> Optional[str]: - """从池中取一个未过期且未在使用的代理(调用前需持有锁)。""" + """从池中取一个未过期、未冷却、未在使用的代理(调用前需持有锁)。""" now = time.time() + # 清理过期代理 expired = [proxy for proxy, ts in self._verified_pool.items() if now - ts > self._pool_ttl] for proxy in expired: del self._verified_pool[proxy] self._in_use.discard(proxy) + # 优先选未冷却、未在用的代理 + for proxy in self._verified_pool: + if proxy not in self._in_use and not self._is_cooling_down(proxy): + self._in_use.add(proxy) + self.current_proxy = proxy + return proxy + + # 退而求其次:所有未在用的(含冷却中的),也比没有强 for proxy in self._verified_pool: if proxy not in self._in_use: self._in_use.add(proxy) self.current_proxy = proxy return proxy - - for proxy in self._verified_pool: - self.current_proxy = proxy - return proxy return None def get_proxy(self, max_attempts: int = 5) -> Optional[str]: @@ -112,20 +131,47 @@ class ProxyManager: return None, msg def mark_bad(self, proxy_url: str) -> None: - """标记代理为不可用,从池中移除(极验失败/代理连接失败时调用)。""" + """ + 标记代理失败:软隔离(冷却期),而非永久移除。 + + - 失败 < _MAX_FAIL_COUNT 次:进入冷却期,冷却后可重新入池 + - 连续失败 >= _MAX_FAIL_COUNT 次:永久移出池 + """ with self._cond: - removed = self._verified_pool.pop(proxy_url, None) self._in_use.discard(proxy_url) if self.current_proxy == proxy_url: self.current_proxy = None - if removed: - logger.info(f"代理标记为不可用并移出池: {proxy_url} (池剩余 {len(self._verified_pool)})") + + record = self._fail_records.get(proxy_url, {"count": 0, "cooldown_until": 0}) + record["count"] = record.get("count", 0) + 1 + fail_count = record["count"] + + if fail_count >= self._MAX_FAIL_COUNT: + # 连续失败超限,永久移出 + self._verified_pool.pop(proxy_url, None) + self._fail_records.pop(proxy_url, None) + logger.info( + f"代理连续失败 {fail_count} 次,永久移出池: {proxy_url} " + f"(池剩余 {len(self._verified_pool)})" + ) + else: + # 软隔离:指数退避冷却 + cooldown = min(self._COOLDOWN_BASE * (2 ** (fail_count - 1)), self._COOLDOWN_MAX) + record["cooldown_until"] = time.time() + cooldown + self._fail_records[proxy_url] = record + logger.info( + f"代理失败 {fail_count} 次,冷却 {cooldown}s: {proxy_url} " + f"(池剩余 {len(self._verified_pool)})" + ) + self._cond.notify_all() def release_proxy(self, proxy_url: str) -> None: """归还代理到池(登录完成后调用,让其他账号可以复用)。""" with self._cond: self._in_use.discard(proxy_url) + # 成功归还时清除失败记录 + self._fail_records.pop(proxy_url, None) self._cond.notify_all() def get_proxies_dict(self, proxy: str = None) -> dict: