使用asyncio批量请求URL迭代后出现连接超时错误的解决方法
问题解决思路与代码修改方案
核心问题分析
- 无限制并发引发资源耗尽:设置
max_connections=None意味着同时发起无上限请求,容易触发目标服务器限流、本地端口耗尽,进而导致连接超时。 - 单个请求失败牵连整批任务:
asyncio.gather默认会因单个任务抛出异常终止所有任务,只要批次里有一个URL超时,整批都会崩溃。 - 异常处理逻辑粗糙:全局
except:捕获所有异常,且重试仅执行一次,没有针对超时请求做针对性的重试策略。 - 依赖全局变量:
testing方法依赖外部text变量,代码耦合度高,易引发不可预期问题。
具体修改方案
1. 限制并发连接数
调整httpx.AsyncClient的连接池参数,设置合理的并发上限,避免资源耗尽:
async with httpx.AsyncClient(limits=httpx.Limits(max_connections=20, max_keepalive_connections=10)) as cl:
2. 单个请求独立异常处理与重试
新增封装请求的方法,给每个请求单独添加超时重试逻辑,用指数退避策略避免频繁重试:
async def _safe_request(self, client, url, max_retries=3): for attempt in range(max_retries): try: # 单独设置连接超时和读取超时,精准控制 resp = await client.get(url, timeout=httpx.Timeout(10.0, connect=5.0)) return resp except (httpx.ConnectTimeout, httpx.ReadTimeout): if attempt == max_retries - 1: print(f"重试{max_retries}次后仍失败: {url}") return None # 指数退避,每次重试间隔翻倍 await asyncio.sleep(2 ** attempt) except Exception as e: print(f"请求{url}发生未知错误: {str(e)}") return None
3. 修改批量请求逻辑,避免整批崩溃
将testing方法改为接收text参数,用封装后的安全请求任务替代直接调用cl.get,单个请求失败不会影响整批:
async def testing(self, text): async with httpx.AsyncClient(limits=httpx.Limits(max_connections=20, max_keepalive_connections=10)) as cl: tasks = [self._safe_request(cl, url) for url in text] async_resp = await asyncio.gather(*tasks) for resp, url in zip(async_resp, text): if resp is not None: data = f"{resp} {url}" print(data) self.log('async_tested', data) else: fail_msg = f"请求失败: {url}" print(fail_msg) self.log('async_tested', fail_msg)
4. 优化主循环的调用与异常处理
去掉全局except:,改为捕获具体异常,同时将text传入testing方法:
p = Processor() for i in range(0, 1000, 50): text = [] fr = i to = i + 50 for line in p.open_file('ordered', fr, to): line = line.strip() text.append(line) print(fr, to) try: asyncio.run(p.testing(text)) except Exception as e: print(f"批次{fr}-{to}执行出错: {str(e)}") # 可选:添加整批重试逻辑,或记录错误批次后续处理
5. 优化日志配置
原代码每次调用log都重新初始化logging,会导致重复配置,建议将日志配置移到类初始化阶段:
class Processor: def __init__(self): FORMAT = '%(message)s' logging.basicConfig( handlers=[ logging.FileHandler( filename=f'{THIS_DIR}/async_tested.log', encoding='utf-8' ) ], level=logging.INFO, format=FORMAT, force=True ) def log(self, data): logging.info(data.strip())
额外建议
- 给每个请求添加随机延迟(比如
await asyncio.sleep(random.uniform(0.1, 0.5))),避免请求过于集中触发服务器限流。 - 将失败的URL记录到单独文件,方便后续单独处理。
内容的提问来源于stack exchange,提问作者user407971
相关产品推荐
相关产品推荐

