使用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
相关产品推荐
相关产品推荐

