如何在FastAPI/Starlette中生成tar.gz流并通过StreamingResponse返回?
如何在FastAPI/Starlette中生成tar.gz流并通过StreamingResponse返回?
我完全懂你遇到的麻烦——tarfile是纯同步的库,它只认带同步write方法的文件对象,但Starlette的StreamingResponse需要异步的方式来发送响应块,直接在同步write里硬塞异步send肯定会出各种协程上下文问题。你之前用asyncio.run_coroutine_threadsafe没成功,问题就出在线程和事件循环的上下文没对齐上。
这里给你一个可行的解决方案,核心思路是用异步队列做同步与异步的桥梁,把tarfile的同步输出和StreamingResponse的异步发送解耦开,同时把耗时的同步打包操作放到线程池里执行,避免阻塞FastAPI的事件循环。
完整实现代码
import asyncio import os import tarfile from pathlib import Path from typing import Mapping, Tuple from concurrent.futures import ThreadPoolExecutor from fastapi import FastAPI from fastapi.responses import StreamingResponse from fastapi.testclient import TestClient from starlette.background import BackgroundTask from starlette.types import Send TAR_FILE_PATH = Path("archive.tar.gz") CHUNK_SIZE = 1024 * 64 # 用大一点的chunk size,效率更高 app = FastAPI() class QueueWriter: """同步的Writer,把tarfile的输出丢进异步队列""" def __init__(self, queue: asyncio.Queue): self.queue = queue def write(self, buffer: bytes) -> int: # 把数据放进队列,这里用put_nowait避免阻塞线程(队列满时会等待) self.queue.put_nowait(buffer) return len(buffer) def close(self): # 发送结束信号 self.queue.put_nowait(None) class TarStreamingResponse(StreamingResponse): def __init__( self, files_to_tar: list[Tuple[str, Path]], # (归档内的文件名, 本地文件路径) status_code: int = 200, headers: Mapping[str, str] | None = None, media_type: str | None = "application/tar+gzip", background: BackgroundTask | None = None, ) -> None: self.files_to_tar = files_to_tar self.status_code = status_code self.media_type = media_type self.background = background self.init_headers(headers) async def stream_response(self, send: Send) -> None: # 发送响应起始帧 await send({ "type": "http.response.start", "status": self.status_code, "headers": self.raw_headers, }) # 创建异步队列,作为同步write和异步send的桥梁 queue: asyncio.Queue[bytes | None] = asyncio.Queue(maxsize=10) writer = QueueWriter(queue) loop = asyncio.get_running_loop() async def send_from_queue(): """从队列取数据并异步发送""" try: while True: buffer = await queue.get() if buffer is None: # 收到结束信号,退出循环 break await send({ "type": "http.response.body", "body": buffer, "more_body": True, }) queue.task_done() finally: # 发送响应结束帧 await send({ "type": "http.response.body", "body": b"", "more_body": False, }) # 启动发送队列数据的异步任务 send_task = asyncio.create_task(send_from_queue()) def pack_tar_files(): """同步打包文件的函数,将在线程池执行""" try: with tarfile.open( mode="w|gz", fileobj=writer, bufsize=CHUNK_SIZE ) as tar: for arcname, file_path in self.files_to_tar: # 添加文件到tar包,用arcname指定归档内的文件名 tar.add(file_path, arcname=arcname) finally: # 不管成功失败,都发送结束信号 writer.close() try: # 把同步打包操作放到线程池执行,避免阻塞事件循环 await loop.run_in_executor(ThreadPoolExecutor(), pack_tar_files) # 等待队列所有数据都处理完 await queue.join() except Exception as e: # 如果打包出错,取消发送任务并抛出异常 send_task.cancel() raise e finally: # 等待发送任务结束 await send_task # 测试用的文件列表 FILES_TO_TAR = [ ("tarfile.py", Path(tarfile.__file__)), ("os.py", Path(os.__file__)) ] @app.get("/") def send_tar() -> StreamingResponse: """返回流式tar.gz文件""" return TarStreamingResponse( files_to_tar=FILES_TO_TAR, headers={ "Content-Disposition": f'attachment; filename="{TAR_FILE_PATH.name}"' }, ) # 测试代码 client = TestClient(app) def test_tar_stream(): test_tar_path = Path("test_archive.tar.gz") if test_tar_path.exists(): os.remove(test_tar_path) # 发起请求并保存响应 response = client.get("/") response.raise_for_status() with open(test_tar_path, "wb") as f: for chunk in response.iter_bytes(CHUNK_SIZE): f.write(chunk) # 验证tar包有效性 with tarfile.open(test_tar_path, "r:gz") as tar: archived_files = {member.name for member in tar.getmembers() if member.isfile()} expected_files = {name for name, _ in FILES_TO_TAR} assert archived_files == expected_files, "归档文件不匹配" print("测试通过!生成的tar.gz有效") os.remove(test_tar_path) if __name__ == "__main__": test_tar_stream()
关键部分解释
QueueWriter类:
- 它是一个同步的文件对象,
write方法把tarfile输出的字节块放进异步队列 - 打包完成后调用
close方法,向队列发送None作为结束信号
- 它是一个同步的文件对象,
线程池执行同步任务:
- 用
loop.run_in_executor把pack_tar_files放到线程池里跑,这样同步的tar打包操作不会阻塞FastAPI的事件循环,保证服务还能处理其他请求
- 用
异步发送任务:
send_from_queue是一个异步任务,不断从队列取数据并调用send发送响应块- 收到
None信号后,发送响应结束帧,完成整个响应流程
异常与资源清理:
- 不管打包成功还是失败,都会发送结束信号,避免队列阻塞
- 如果打包出错,会取消发送任务并抛出异常,保证客户端能收到错误响应
注意事项
- 调整
CHUNK_SIZE可以平衡内存占用和传输效率,大文件建议用64KB或更大的块 - 如果是从S3取文件,只需要把
tar.add换成从S3流式读取并写入tarfile的方式(比如用tarfile.addfile配合S3的get_object的Body流) - 生产环境建议加上超时和异常捕获,避免队列无限阻塞
内容来源于stack exchange
相关产品推荐
相关产品推荐

