FastAPI中如何通过StreamingResponse流式生成完整多文件Zip归档?
问题:流式生成多文件Zip归档时文件损坏的修复方法
我有一段生成Zip归档并进行流式传输的代码。处理大请求时,数据开始流式传输前可能需要数分钟处理时间,导致请求取消后Python进程仍会继续运行。我希望通过生成器函数逐个生成压缩后的文件,同时让用户获得完整的Zip归档,以此提升请求取消的鲁棒性。
我写了一个最小可复现示例:生成单个文件(N=1)时代码正常运行,但生成2个及以上文件(N>=2)时,生成的Zip文件损坏。示例代码如下:
# test endpoint from fastapi import APIRouter, Path from numpy import random import io, zipfile from fastapi.responses import StreamingResponse router = APIRouter(tags=["Make files"]) # make some fake data and zip, then yield def make_data(N): """ Make fake data. """ CHUNK_SIZE = 1024*1024 for n in range(N): content = random.random(100) name = f'{n:02}.txt' # Create new in-memory zip file for each file s = io.BytesIO() with zipfile.ZipFile(s, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=2) as zf: # Add file content to the in-memory zip file zf.writestr(name, content) # Seek to the beginning of the in-memory zip file s.seek(0) # Yield the content of the in-memory zip file for the current file while chunk := s.read(CHUNK_SIZE): yield chunk # streamingresponse def stream_data(N): """ Stream the files. """ return StreamingResponse( make_data(N), media_type="application/zip", headers={"Content-Disposition": f"attachment; filename=download.zip"}) # endpoint @router.get("/{N}") async def yield_files( N: int = Path(..., decription="random files to make")): return stream_data(N)
修复方案
你当前的核心问题是:每次循环都创建独立的ZipFile对象,相当于把多个完整的Zip文件拼接在一起,而非在同一个归档中添加多文件,最终导致Zip损坏。要实现流式生成合法的多文件Zip归档,需要复用同一个ZipFile实例,并增量输出压缩数据。
以下是修改后的完整代码:
# test endpoint from fastapi import APIRouter, Path from numpy import random import io, zipfile from fastapi.responses import StreamingResponse router = APIRouter(tags=["Make files"]) # make some fake data and stream the zip archive def make_data(N): """ Generate fake data and stream the complete zip archive incrementally. """ CHUNK_SIZE = 1024 * 1024 # 创建单个内存流用于整个Zip归档 zip_stream = io.BytesIO() # 手动管理ZipFile生命周期,避免with语句自动关闭 zf = zipfile.ZipFile(zip_stream, "w", compression=zipfile.ZIP_DEFLATED, compresslevel=2) try: last_pos = 0 # 记录上次读取到的流位置,用于增量输出 for n in range(N): # 将numpy数组转为字节串(原代码直接传数组会报错) content = random.random(100).tobytes() name = f'{n:02}.txt' # 向归档中添加当前文件 zf.writestr(name, content) # 获取当前流的末尾位置 zip_stream.seek(0, io.SEEK_END) current_pos = zip_stream.tell() # 读取并yield本次新增的压缩数据块 zip_stream.seek(last_pos) remaining = current_pos - last_pos while remaining > 0: chunk_size = min(CHUNK_SIZE, remaining) chunk = zip_stream.read(chunk_size) yield chunk remaining -= chunk_size # 更新上次位置为当前末尾 last_pos = current_pos # 关闭ZipFile,写入归档结束标记(必须操作,否则Zip文件损坏) zf.close() # 读取并yield归档结束标记的剩余数据 zip_stream.seek(last_pos) while chunk := zip_stream.read(CHUNK_SIZE): yield chunk finally: # 确保资源释放 zf.close() # streamingresponse def stream_data(N): """ Stream the zip archive. """ return StreamingResponse( make_data(N), media_type="application/zip", headers={"Content-Disposition": f"attachment; filename=download.zip"}) # endpoint @router.get("/{N}") async def yield_files( N: int = Path(..., description="number of random files to generate")): # 修正拼写错误 return stream_data(N)
关键修改说明
- 复用单个ZipFile:全程使用同一个ZipFile实例,所有文件都添加到同一个归档中,保证Zip结构合法。
- 增量输出数据:每次添加文件后只输出流中新增的部分,避免重复发送已传输的数据,同时实现流式传输的即时性。
- 修复数据类型问题:将numpy数组转为字节串,避免
writestr方法报错。 - 手动关闭ZipFile:必须关闭ZipFile来写入归档结束标记,这是生成合法Zip文件的必要步骤。
- 修正拼写错误:将
decription改为正确的description。
修改后,生成器会逐个处理文件并即时输出压缩数据,用户取消请求时生成器会立即停止运行,不会浪费进程资源,同时生成的Zip文件完整可用。
内容的提问来源于stack exchange,提问作者Adam
相关产品推荐
相关产品推荐

