使用asyncio+aiohttp批量请求30000个URL耗时过长如何优化
问题描述
现有包含30000个URL的列表,需要提取每个URL对应页面的文本内容,最终汇总得到存储所有页面文本的列表。当前基于asyncio实现的版本运行速度优于requests同步实现,但处理完全部30000个URL仍需10分钟,未达效率预期,需调整代码进一步提升处理效率。
现有实现代码
核心异步处理逻辑
async def main(session, url): async with session.get(url, timeout=False) as resp: text = await resp.text() return text async def obtain_text_from_url(urls, list_text): my_conn = aiohttp.TCPConnector(limit=200) async with aiohttp.ClientSession(connector=my_conn) as session: tasks = [] for url in source_boamp: tasks.append(asyncio.ensure_future(main(session, url))) tasks2 = await asyncio.gather(*tasks) for task in tasks2: list_text.append(task) return list_text
注:原代码存在缩进错误、遍历变量未使用入参urls、引用未定义变量source_boamp的问题,会直接导致运行报错
任务启动函数
def launch(urls): list_text = [] asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) asyncio.run(obtain_text_from_url(urls, liste_text)) return liste_text
注:原代码存在变量名拼写错误,入参传的liste_text和之前定义的list_text不一致,会直接触发变量不存在报错
效率优化方案
- 不要一次性创建全部30000个任务直接丢给
asyncio.gather:瞬间创建数万协程会占用大量不必要的内存,同时瞬时发起的海量请求极易触发目标站点限流、丢包、连接重置,反而拖慢整体处理速度。用asyncio.Semaphore控制全局并发上限,根据本地带宽、目标站点承载能力在300-1000区间测试最优并发值,不要固定使用200的连接数限制。 - 优化TCPConnector配置:
- 新增
enable_cleanup_closed=True参数,自动清理已关闭的失效连接,减少无效连接占用 - 配置
limit_per_host参数和总连接数limit匹配,避免单域名并发被默认规则限制 - 移除
timeout=False配置,设置合理的全局超时(比如10-15秒),直接丢弃长时间无响应的卡死请求,避免单个慢请求长期占用连接资源阻塞队列
- 新增
- 调整结果收集逻辑:用
asyncio.as_completed替代asyncio.gather,拿到一个响应结果就立刻存入结果列表,不需要等所有请求全部完成再统一处理,减少等待开销 - 优化响应解析逻辑:如果提前知道目标页面编码,直接用
await resp.read()拿到二进制内容后手动指定编码解码,不要直接调用resp.text()——resp.text()会自动做编码检测,这个过程在页面量大的时候性能开销非常高。如果需要提取页面特定文本而非全量源码,换用selectolax、lxml这类C实现的解析库,解析速度比纯Python实现快数倍。 - 开启TCP连接复用:创建ClientSession时统一配置
Connection: keep-alive请求头,复用已建立的TCP连接,减少重复TCP握手的开销。 - 增加轻量失败处理:对超时、5xx类服务端错误的请求做1-2次重试,遇到4xx类客户端错误(比如404、403)直接跳过,不要反复重试无效请求。
- 替换默认事件循环:Windows环境下可安装
uvloop替换asyncio默认事件循环,调度效率比默认实现高30%左右,注意0.17以上版本的uvloop已原生支持Windows。
注意:如果目标站点配置了反爬策略,过高并发会直接触发IP封禁、验证码拦截,这种情况需要配合代理池使用,否则仅调整代码无法突破速度上限。
优化后的核心逻辑参考
import asyncio import aiohttp # 根据测试调整并发上限 CONCURRENCY = 500 TIMEOUT = aiohttp.ClientTimeout(total=15) HEADERS = {"Connection": "keep-alive"} async def fetch(session, sem, url): async with sem: try: async with session.get(url, timeout=TIMEOUT) as resp: if resp.status != 200: return "" content = await resp.read() # 已知站点编码可直接指定,跳过自动编码检测 return content.decode("utf-8", errors="ignore") except Exception: return "" async def obtain_text_from_url(urls): sem = asyncio.Semaphore(CONCURRENCY) connector = aiohttp.TCPConnector( limit=CONCURRENCY, limit_per_host=CONCURRENCY, enable_cleanup_closed=True ) list_text = [] async with aiohttp.ClientSession(connector=connector, headers=HEADERS) as session: tasks = [asyncio.create_task(fetch(session, sem, url)) for url in urls] for task in asyncio.as_completed(tasks): res = await task if res: list_text.append(res) return list_text def launch(urls): asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy()) # 安装uvloop后可打开下面注释替换事件循环 # import uvloop # asyncio.set_event_loop_policy(uvloop.EventLoopPolicy()) return asyncio.run(obtain_text_from_url(urls))
内容的提问来源于stack exchange,提问作者Fitz
相关产品推荐
相关产品推荐

