内存有限时,如何用aiohttp分批异步处理大文件中的URL?
如何分批读取大文件并异步处理URL?
问题描述
我有一个每行存一个URL的超大文件,目前用aiohttp做异步批量请求处理。但文件体积太大、内存有限,而且我不清楚文件的总行数,逐行计数又特别耗时。我想实现这样的流程:
- 读取100,000行存入列表
- 处理该列表的URL请求
- 暂停文件读取,等处理完再继续
- 重复上述步骤直到所有行处理完毕
我原本想搞异步方案,但可能对相关库理解有偏差,试了下面的代码,感觉不太对:
counter = 0 inputs=[] async with open("test.txt") as f: for line in f: counter=counter+1 if counter%100000 != 0: inputs.append(line.strip()) else: await get_req_fn(inputs) inputs=[]
解决方案
先给你梳理下现有代码的几个小问题,再给出优化后的实现:
async with open完全没必要——本地文件读取是同步操作,用普通的with open就够了,异步上下文管理器在这里是多余的。- 最后一批不足100000行的URL会被漏掉,循环结束后没有处理剩余的
inputs。 - 批量处理时的会话管理可以更规范,
aiohttp的ClientSession最好复用或者按批次合理创建。
推荐实现(同步读文件+异步处理)
本地文件读取的速度远快于网络请求,同步读取不会拖慢整体效率,反而代码更简洁:
import aiohttp import asyncio async def process_batch(session, urls): """异步批量处理URL请求""" # 创建所有请求任务,用return_exceptions避免单个请求失败导致整批挂掉 tasks = [session.get(url) for url in urls] responses = await asyncio.gather(*tasks, return_exceptions=True) # 这里写你的响应处理逻辑,比如解析、存储结果等 for url, resp in zip(urls, responses): if isinstance(resp, Exception): print(f"请求 {url} 失败: {str(resp)}") else: # 示例:打印响应状态码 print(f"请求 {url} 成功,状态码: {resp.status}") # 可以在这里解析resp.text()或者resp.json() async def main(): batch_size = 100000 inputs = [] # 同步读取文件,本地IO速度足够快,不会成为瓶颈 with open("test.txt", "r", encoding="utf-8") as f: for line in f: url = line.strip() if url: # 跳过空行 inputs.append(url) # 攒够一批就处理 if len(inputs) == batch_size: async with aiohttp.ClientSession() as session: await process_batch(session, inputs) inputs = [] # 处理最后一批不足batch_size的URL if inputs: async with aiohttp.ClientSession() as session: await process_batch(session, inputs) if __name__ == "__main__": asyncio.run(main())
可选:异步读取文件(适合极端大文件场景)
如果你的文件大到同步读取都会阻塞事件循环(这种情况很少见),可以把文件读取放到线程池里,用asyncio.to_thread实现异步读取:
import aiohttp import asyncio async def read_file_batches(file_path, batch_size): """异步生成每一批URL列表""" inputs = [] # 把同步读取操作放到线程池,不阻塞事件循环 def read_sync(): with open(file_path, "r", encoding="utf-8") as f: yield from f async for line in asyncio.to_thread(read_sync): url = line.strip() if url: inputs.append(url) if len(inputs) == batch_size: yield inputs inputs = [] if inputs: yield inputs async def process_batch(session, urls): # 和上面的process_batch逻辑一致 tasks = [session.get(url) for url in urls] responses = await asyncio.gather(*tasks, return_exceptions=True) for url, resp in zip(urls, responses): if isinstance(resp, Exception): print(f"请求 {url} 失败: {str(resp)}") else: print(f"请求 {url} 成功,状态码: {resp.status}") async def main(): batch_size = 100000 async with aiohttp.ClientSession() as session: # 异步迭代每一批URL async for batch in read_file_batches("test.txt", batch_size): await process_batch(session, batch) if __name__ == "__main__": asyncio.run(main())
关键要点
- 批量处理的合理性:100000个并发请求可能会触发目标服务器的限流,你可以根据实际情况调整
batch_size,或者在process_batch里用asyncio.Semaphore控制并发数。 - 会话复用:
aiohttp.ClientSession最好复用(比如整个main里只用一个),频繁创建销毁会话会影响效率。 - 异常处理:一定要用
return_exceptions=True或者捕获单个任务的异常,避免一个请求失败导致整批处理中断。
内容的提问来源于stack exchange,提问作者user9749763
相关产品推荐
相关产品推荐

