You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Asyncio+Aiohttp执行超50万异步请求时内存泄漏问题求助

解决Aiohttp+Asyncio批量请求的内存泄漏问题

你的代码存在多个导致内存持续增长的问题,以下是具体问题分析和修复方案:

问题分析

  • 重复创建ClientSession:fetch函数明明已经接收了外部传入的session,却每次请求都新建一个ClientSession,且没有通过async with正确管理生命周期,导致大量未关闭的会话对象堆积在内存中。
  • 重复调用asyncio.gather:fetch_all中先await asyncio.gather(*tasks)完成所有任务,之后又再次调用await asyncio.gather(*tasks),这不仅是冗余操作,还会让tasks列表一直持有已完成任务的引用,无法被垃圾回收器回收,持续占用内存。
  • Semaphore使用方式错误:在循环中对每个URL都执行async with sem,这会导致任务串行创建,完全没利用并发优势,还可能引发任务对象堆积。
  • 重复读取响应内容:同时调用resp.text()和resp.read()会把响应内容在内存中存储两份(二进制字节流和字符串),大量请求下会额外消耗大量内存。

修复后的代码

import asyncio
from asyncio import Semaphore
import aiohttp
from bs4 import BeautifulSoup

async def fetch(session, sem, url):
    try:
        headers = {}  # 替换为你的实际headers
        params = {}   # 替换为你的实际params
        async with sem:
            async with session.get(url, headers=headers, params=params, verify_ssl=False) as resp:
                # 只读取文本内容,不需要同时存二进制和字符串
                text = await resp.text()
                return text
    except Exception as e:
        print(f"url: {url} error happened: {str(e)}")
        # 异常时返回None,避免后续处理报错
        return None

async def fetch_all(urls):
    # 限制并发数
    sem = Semaphore(100)
    # 使用TCPConnector控制连接池,避免连接泄漏
    connector = aiohttp.TCPConnector(limit=100, force_close=True)
    async with aiohttp.ClientSession(connector=connector, cookie_jar=aiohttp.DummyCookieJar()) as session:
        tasks = []
        for url in urls:
            task = asyncio.create_task(fetch(session, sem, url))
            tasks.append(task)
        # 只执行一次gather,获取所有结果
        datas = await asyncio.gather(*tasks, return_exceptions=True)
        return datas

def get_result(data, chupindict):
    if not data:
        return
    try:
        soup = BeautifulSoup(data, 'html.parser')
        name = soup.find('div', class_='name').text.strip()
        chupin = soup.find('div', class_="panel-wrapper", id="出品").text.strip()
        chupindict[name] = chupin
    except Exception as e:
        print(f"解析数据出错: {str(e)}")

if __name__ == "__main__":
    urls = []
    chupindict = {}
    # 生成测试URL(替换为你的实际URL生成逻辑)
    for i in range(0, 500000):
        url = f"http://example.com/test/{i}"
        urls.append(url)
    
    loop = asyncio.new_event_loop()
    asyncio.set_event_loop(loop)
    try:
        datas = loop.run_until_complete(fetch_all(urls))
        # 逐个处理结果,避免一次性持有所有数据
        for data in datas:
            get_result(data, chupindict)
    finally:
        loop.close()

关键修改说明

  1. 复用ClientSession:全局只创建一个ClientSession,所有请求复用这个会话,通过TCPConnector控制连接池大小,避免连接泄漏。
  2. 正确使用Semaphore:将信号量传入fetch函数,在请求内部获取信号量,确保并发数被正确控制,同时任务可以批量创建。
  3. 单次调用gather:只执行一次asyncio.gather获取所有结果,避免任务引用堆积。
  4. 避免重复读取响应:只保留需要的响应格式(这里用文本),减少内存占用。
  5. 结果分批处理:逐个处理返回的结果,而不是一次性持有所有响应数据,降低峰值内存占用。
  6. 异常处理优化:在fetch和get_result中都增加了异常捕获,避免单个请求/解析失败导致程序崩溃,同时减少无效内存占用。

内容的提问来源于stack exchange,提问作者Max Chen

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.11 16:03:22