AsyncIO中Timeout异常处理:实现重试或跳过失败任务
解决Asyncio对接爬虫服务的Timeout异常:重试与跳过方案
核心问题拆解
你的代码存在三个关键问题导致Timeout错误无法处理:
- 未给API请求设置显式超时限制,超时错误会直接抛出
asyncio.gather默认会因单个任务异常终止所有任务,无法实现"跳过失败任务"- 缺少重试逻辑,失败后直接终止流程
修改后的完整代码
import asyncio import aiohttp from time import perf_counter import csv path = "*******************" domains = [] total_count = 0 # 域名读取逻辑保留 with open(path, 'r') as file: csvreader = csv.reader(file) for row in csvreader: try: website = row[4].split("//")[-1].split("www.")[-1].split('/')[0] if website == "": continue domains.append(website) except: continue sample = domains[0:50] async def fetch(s, body, max_retries=3, retry_delay=2): """带超时、重试的API请求处理""" for attempt in range(max_retries + 1): try: # 设置10秒总超时,避免无限等待 async with s.post( 'https://****************', json=body, timeout=aiohttp.ClientTimeout(total=10) ) as r: if r.status != 200: print(f"请求失败,状态码: {r.status} | 域名: {body['domain']}") if attempt < max_retries: await asyncio.sleep(retry_delay) continue return None enrich_response = await r.json() employees = enrich_response.get('employees', []) global total_count for employee in employees: job_title = employee.get('job_title', '') if job_title in ("Owner", "CEO"): print(employee) print("*" * 50) total_count += 1 print(f"Total Count: {total_count}") return enrich_response except asyncio.TimeoutError: print(f"超时重试 {attempt+1}/{max_retries+1} | 域名: {body['domain']}") if attempt < max_retries: await asyncio.sleep(retry_delay) continue print(f"重试耗尽,跳过域名: {body['domain']}") return None except Exception as e: print(f"未知错误: {str(e)} | 域名: {body['domain']}") if attempt < max_retries: await asyncio.sleep(retry_delay) continue print(f"重试耗尽,跳过域名: {body['domain']}") return None async def fetch_all(s, bodies): tasks = [] for body in bodies: tasks.append(asyncio.create_task(fetch(s, body))) # return_exceptions=True 让单个任务失败不影响全局,异常会被作为结果返回 res = await asyncio.gather(*tasks, return_exceptions=True) # 过滤掉异常和失败结果,只保留成功数据 return [item for item in res if not isinstance(item, Exception) and item is not None] async def main(): bodies = [] for domain in sample: bodies.append({ "api_key": "********************************", "domain": domain }) async with aiohttp.ClientSession() as session: data = await fetch_all(session, bodies) print(f"成功完成请求数: {len(data)}") if __name__ == '__main__': start = perf_counter() try: asyncio.run(main()) except Exception as e: print(f"主程序异常: {str(e)}") stop = perf_counter() print(f"耗时: {stop - start:.2f} 秒")
关键修改说明
1. 超时与重试机制
- 给
aiohttp.post添加ClientTimeout,强制限制请求总时长 - 用循环实现最多3次重试,每次失败后等待2秒再尝试
- 针对性捕获
TimeoutError和其他异常,输出明确日志便于排查
2. 跳过失败任务
- 调用
asyncio.gather时添加return_exceptions=True,单个任务抛出异常不会终止整个批量任务,异常会被作为结果返回 - 最后过滤掉异常和
None结果,只保留成功的请求数据
3. 其他优化
- 用
dict.get()替代直接索引,避免因返回数据结构异常导致的KeyError - 简化职位判断逻辑,代码更简洁
- 增加详细日志,明确显示哪个域名请求失败
可选调整
- 不需要重试直接跳过:将
max_retries设为0即可 - 可根据服务商限制调整
retry_delay(重试间隔)和max_retries(最大重试次数) - 避免全局变量:可以用
asyncio.Lock保护计数器,或者让fetch返回找到的目标职位数量,最后在main中汇总统计
内容的提问来源于stack exchange,提问作者Dominic Jay
相关产品推荐
相关产品推荐

