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

内存有限时,如何用aiohttp分批异步处理大文件中的URL?

如何分批读取大文件并异步处理URL?

问题描述

我有一个每行存一个URL的超大文件,目前用aiohttp做异步批量请求处理。但文件体积太大、内存有限,而且我不清楚文件的总行数,逐行计数又特别耗时。我想实现这样的流程:

  1. 读取100,000行存入列表
  2. 处理该列表的URL请求
  3. 暂停文件读取,等处理完再继续
  4. 重复上述步骤直到所有行处理完毕

我原本想搞异步方案,但可能对相关库理解有偏差,试了下面的代码,感觉不太对:

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=[]

解决方案

先给你梳理下现有代码的几个小问题,再给出优化后的实现:

  1. async with open完全没必要——本地文件读取是同步操作,用普通的with open就够了,异步上下文管理器在这里是多余的。
  2. 最后一批不足100000行的URL会被漏掉,循环结束后没有处理剩余的inputs。
  3. 批量处理时的会话管理可以更规范,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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 09:20:10