aiohttp流式get请求必须用response.content.iter_chunks()或iter_any()逐块读取,避免read()或text()全量加载内存;需配合适当超时、连接复用及异常处理确保大文件传输健壮性。

用 aiohttp 发起流式 GET 请求并分块读取
直接调用 response.text() 或 response.read() 会把整个响应体加载进内存,大文件(比如几百 MB 的日志或导出数据)极易触发 MemoryError。必须用流式接口逐段读取。
关键点是:不等响应结束,拿到 aiohttp.ClientResponse 后立即调用 content.iter_chunks() 或 content.iter_any(),配合 async for 循环处理:
import aiohttp
import asyncio
<p>async def stream_large_file(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as response:</p><h1>确保状态码正常,避免把 404/500 当作数据流</h1><pre class="brush:python;toolbar:false;"> response.raise_for_status()
async for chunk, _ in response.content.iter_chunks():
# 这里处理每个 bytes 块,例如写入文件、解析 JSON 行、校验 CRC
process_chunk(chunk)
-
iter_chunks()返回(bytes, bool)元组,第二个值为True表示这是最后一块 - 如果只关心原始字节且不需要判断结尾,用
iter_any()更轻量 - 务必在
async with session.get(...)内完成整个流读取,否则连接可能提前关闭
边下载边解压或解析 JSONL 的典型场景
真实业务中,你往往不是“存文件”,而是“用数据”——比如从 S3 预签名 URL 流式拉取一个 .gz 压缩包,实时解压并逐行处理 JSONL;或直接消费 API 返回的未压缩 JSONL 流。
注意:Python 标准库的 gzip.decompress() 是全量解压函数,不能用于流式。必须用 gzip.GzipFile 包装一个可读的异步流:
import gzip
from io import BytesIO
<p>async def stream_and_gunzip(url):
async with aiohttp.ClientSession() as session:
async with session.get(url) as response:</p><h1>把 response.content 包装成类文件对象</h1><pre class="brush:python;toolbar:false;"> # 注意:aiohttp 的 content 不是 file-like,需用 BytesIO + 分块喂入
buffer = BytesIO()
async for chunk in response.content.iter_any():
buffer.write(chunk)
# 这里不能直接传 buffer 给 GzipFile —— 它需要 seekable,而流式 buffer 不满足
# 正确做法:用 streaming decompressor,如 https://pypi.org/project/aiofiles/ + gzip 搭配,或改用 zlib.decompressobj()
更稳妥的做法是用 zlib.decompressobj() 手动流式解压:
- 初始化
decompressor = zlib.decompressobj(16 + zlib.MAX_WBITS)(支持 gzip header) - 每次收到
chunk就调用decompressor.decompress(chunk),剩余未完成部分用decompressor.flush()收尾 - JSONL 场景下,把解压出的 bytes 按
\n切分,逐行json.loads(),避免累积整块文本
aiohttp 默认超时和连接池对长流的影响
默认 aiohttp.ClientTimeout(total=5*60),但大文件传输可能耗时远超 5 分钟,尤其网络波动时。单纯调大 total 不够——TCP 连接空闲超时(如 Nginx 的 keepalive_timeout)、代理中断、SSL 心跳缺失都可能导致流意外断开。
- 显式设置
timeout=aiohttp.ClientTimeout(total=None, sock_read=600):禁用总超时,但保留单次读操作最长等待(防卡死) - 启用连接复用:
connector = aiohttp.TCPConnector(keepalive_timeout=3600),避免频繁建连开销 - 加
raise_for_status=False并手动检查response.status,因为某些服务在流中途出错时可能不发完整 HTTP 错误码,而是静默截断 - 务必捕获
aiohttp.ClientPayloadError和asyncio.TimeoutError,它们常出现在流被中断时
为什么不用 requests + stream=True?
requests 的 stream=True 是同步阻塞流,底层仍是 urllib3 的 blocking socket read。在 asyncio 协程里直接调用它,会阻塞整个事件循环,失去异步意义——哪怕你只在一个协程里用,也会拖慢其他并发任务。
常见误用:
# ❌ 错误:在 async def 里调用 requests.get(..., stream=True)
async def bad_example():
resp = requests.get(url, stream=True) # 这里就卡住 event loop 了
for chunk in resp.iter_content(8192):
...
结论很直接:只要主流程是 async,HTTP 客户端就必须选原生异步库,aiohttp 是目前最成熟的选择;httpx 也支持 async,但要注意其默认后端(asyncio vs trio)和连接池行为略有差异。
真正难的从来不是“怎么读”,而是“怎么确保每一块都可靠到达、不丢不重、能恢复”。流式处理的健壮性,藏在超时策略、重试逻辑、断点续传标记、以及 chunk 边界是否对齐业务语义(比如一个 JSON 对象不能被切在中间)这些细节里。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











