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

使用Python asyncio处理百万条记录时触发内存错误如何解决

问题根因

你当前代码的内存溢出和Semaphore配置无关,核心是三个设计错误直接把内存打满:

  • 分批逻辑完全失效:你设置的BATCH_SIZE=1000000是单次读取100万行数据,而且读完之后直接给所有行创建协程任务塞进全局tasks列表。asyncio.Semaphore(500)仅限制同时发起的网络请求数,不会限制任务对象的创建——100万个提前创建的task对象、每个task携带的参数、协程上下文本身就会占用数GB内存,还没等请求发出去内存已经吃了大半。
  • 结果全量驻留内存:你把所有请求的返回值全部攒到responses列表里,等所有请求跑完才一次性返回,百万级返回结果(哪怕是单个小字典)持续堆积,内存占用只会涨不会跌。
  • 异常处理导致内存泄漏:你在两个捕获块里直接return ex返回异常对象,异常对象自带的traceback会持有整个调用栈的引用,关联的上下文对象无法被垃圾回收。当内存不足触发MemoryError时,解释器尝试规范化异常信息时找不到足够内存,就会直接抛出你看到的致命错误中止进程。
  • 额外隐患:500的并发数设置过高,aiohttp默认会把响应体全量读入内存,500个请求同时加载响应体很容易出现内存尖刺,进一步放大OOM概率。
修复方案

按照以下规则调整代码即可彻底解决内存问题:

  • 把批次大小下调到1000~5000的合理区间,每一批次处理完立刻把结果写入磁盘/下游存储,清空当前批次的任务、结果引用,绝对不要全局攒全量任务和结果。
  • 不要提前一次性创建所有协程任务,逐批次创建任务、逐批次等待结果、逐批次落盘。
  • 异常不要直接返回Exception对象,返回可序列化的简单字典/字符串,避免traceback持有上下文导致内存泄漏。
  • 并发数下调到30~100区间(根据接口承载能力调整),过高的并发不仅占内存,还极易触发接口限流。
  • 给aiohttp配置匹配并发数的TCP连接池上限,避免无限制创建连接占用额外内存。

修正后的可运行代码示例:

import asyncio
import json
import aiohttp
from itertools import islice
from aiohttp import ClientSession

# 合理配置参数
BATCH_SIZE = 2000  # 单批次处理2000条,处理完即落盘清空
MAX_CONCURRENT = 50  # 并发数根据接口承载能力调整,普通业务接口50足够
LOGGER = your_logger_instance
API = your_api_prefix
OUTPUT_FILE_PATH = "process_result.jsonl"

async def get_valuation(url, params, api_header, session, semaphore):
    async with semaphore:
        try:
            async with session.get(url, headers=api_header, timeout=30) as response:
                status_code = response.status
                if status_code != 200:
                    return {params: f"not found, {status_code}"}
                asynch_response = await response.json(content_type=None)
                mmr = await get_best_match(params, asynch_response, str(status_code))
                return mmr
        except Exception as ex:
            LOGGER.error(f"Request failed, param: {params}, error: {str(ex)}")
            # 返回简单可序列化结构,不返回异常对象
            return {params: f"request error: {str(ex)}"}


async def wrap_get_fuzzy_match(func, *args, **kwargs):
    try:
        return await func(*args, **kwargs)
    except Exception as err:
        LOGGER.error(f"Task wrap error: {str(err)}")
        return {"error": f"wrap error: {str(err)}"}


async def main(headers, input_file):
    sema = asyncio.Semaphore(MAX_CONCURRENT)
    # 提前打开结果文件,边处理边写入,不攒全量结果
    with open(OUTPUT_FILE_PATH, "w", encoding="utf-8") as out_f:
        # 配置连接池上限和并发数匹配
        connector = aiohttp.TCPConnector(limit=MAX_CONCURRENT + 20)
        async with ClientSession(connector=connector) as session:
            with open(input_file, "r", encoding="utf-8") as f:
                while True:
                    batch = [line.strip('\n') for line in islice(f, BATCH_SIZE)]
                    if not batch:
                        break
                    tasks = []
                    for param in batch:
                        # Python3.7+用create_task替代ensure_future,语义更清晰
                        task = asyncio.create_task(wrap_get_fuzzy_match(
                            get_valuation,
                            url= API + param,
                            params=param,
                            api_header=headers,
                            session=session,
                            semaphore=sema,
                        ))
                        tasks.append(task)
                    # 仅等待当前批次任务完成
                    batch_responses = await asyncio.gather(*tasks)
                    # 结果直接写入文件,不存入全局列表
                    for resp in batch_responses:
                        out_f.write(json.dumps(resp, ensure_ascii=False) + "\n")
                    # 主动清空引用,帮助垃圾回收
                    tasks.clear()
                    batch.clear()
                    batch_responses.clear()
    return "process completed"
额外优化建议
  • 如果内存要求极端严格,可以把asyncio.gather换成asyncio.as_completed,完成一条就写一条,不用等整批次跑完,内存占用会更低。
  • 如果调用的接口返回体体积较大,可以开启aiohttp的流式响应读取,避免全量响应体驻留内存。
  • 处理百万级数据时不要用list/tuple存全量结果,用jsonl、csv等追加写的格式逐行落盘是最稳妥的方案。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.02 05:57:37