FastAPI + aioboto3:如何实现iter_chunks流式下载S3文件?
问题描述
我需要开发一个FastAPI接口,从S3下载文件并直接返回给客户端,为此写了如下代码:
from io import BytesIO import aioboto3 import uvicorn from fastapi import FastAPI from starlette.responses import StreamingResponse app = FastAPI() async def download_file(key): session = aioboto3.Session(aws_access_key_id='xxx', aws_secret_access_key='yyy') async with session.client("s3", endpoint_url='http://localhost:9000') as s3: return await s3.get_object(Bucket='my-backups-12345', Key=key) @app.get('/works/') async def works(): # 这个接口可以正常下载文件 key = 'facd0b48-312b-4cf8-8513-7ad15a8f7bae/2-2.xlsx' session = aioboto3.Session(aws_access_key_id='xxx', aws_secret_access_key='yyy') async with session.client("s3", endpoint_url='http://localhost:9000') as s3: resp = await s3.get_object(Bucket='my-backups-12345', Key=key) content = await resp['Body'].read() return StreamingResponse(BytesIO(content)) @app.get('/does-not-work/') async def does_not_work(): # 这个接口无法正常下载文件,想确认iter_chunks的用法是否正确 # 目标是拿到S3的文件块就立即返回 key = 'facd0b48-312b-4cf8-8513-7ad15a8f7bae/2-2.xlsx' session = aioboto3.Session(aws_access_key_id='tlBne2gPKidWDkKx', aws_secret_access_key='p8ClBh5h0nsjb98MeDs6k5l5gI0OZ1bB') async with session.client("s3", use_ssl=False, endpoint_url='http://localhost:9000', verify=False) as s3: resp = await s3.get_object(Bucket='my-backups-12345', Key=key) return StreamingResponse(resp['Body'].iter_chunks()) if __name__ == '__main__': uvicorn.run(app, host="0.0.0.0")
其中/does-not-work接口无法正常返回文件,我想确认自己使用iter_chunks的方式是否正确?我的需求是实现流式返回,拿到文件块就立刻发给客户端。
更新1:我怀疑问题出在最后一个文件块的处理上。
更新2:看起来aioboto3的会话/客户端存在问题,因为我自定义的S3FileStreamingResponse可以正常工作:
import aioboto3 import time from starlette.responses import StreamingResponse class S3FileStreamingResponse(StreamingResponse): """自定义S3文件流式响应""" chunk_size = 4096 async def __call__(self, scope, receive, send) -> None: """处理请求并发送响应""" key = '0d8bba25-4fae-4e5f-8804-4c9282a193cb/2-2.xlsx' session = aioboto3.Session(aws_access_key_id='tlBne2gPKidWDkKx', aws_secret_access_key='p8ClBh5h0nsjb98MeDs6k5l5gI0OZ1bB') async with session.client("s3", use_ssl=False, endpoint_url='http://localhost:9000', verify=False) as s3: resp = await s3.get_object(Bucket='my-backups-12345', Key=key) time.sleep(10) total_size = resp["ContentLength"] sent_size = 0 await send( { "type": "http.response.start", "status": self.status_code, "headers": self.raw_headers, } ) async for chunk in resp["Body"].iter_chunks(): sent_size += len(chunk) await send( { "type": "http.response.body", "body": chunk, "more_body": sent_size < total_size, } ) if self.background is not None: await self.background()
问题原因与解决方案
为什么原/does-not-work接口失效?
核心问题在于aioboto3的客户端上下文管理器在函数返回时就被关闭了。当你在async with块里返回StreamingResponse(resp['Body'].iter_chunks())时,async with会立刻退出,关闭S3客户端连接,导致后续iter_chunks()尝试读取数据时连接已经断开,无法获取文件块。
而自定义的S3FileStreamingResponse能工作,是因为整个流式发送的逻辑都包裹在async with客户端的上下文里,直到所有文件块发送完成,客户端才会被关闭。
正确的实现方式
方式1:用异步生成器包裹读取逻辑
把S3文件的读取逻辑放到异步生成器中,确保客户端在生成器运行期间保持打开状态:
@app.get('/stream/') async def stream_file(): key = 'facd0b48-312b-4cf8-8513-7ad15a8f7bae/2-2.xlsx' session = aioboto3.Session(aws_access_key_id='tlBne2gPKidWDkKx', aws_secret_access_key='p8ClBh5h0nsjb98MeDs6k5l5gI0OZ1bB') async def stream_generator(): async with session.client("s3", use_ssl=False, endpoint_url='http://localhost:9000', verify=False) as s3: resp = await s3.get_object(Bucket='my-backups-12345', Key=key) async for chunk in resp['Body'].iter_chunks(): yield chunk return StreamingResponse(stream_generator())
方式2:优化自定义响应类(复用性更强)
将自定义响应类改成可配置的,支持传入key、bucket等参数,避免硬编码:
import aioboto3 from starlette.responses import StreamingResponse class S3FileStreamingResponse(StreamingResponse): chunk_size = 4096 def __init__(self, key: str, bucket: str, aws_access_key_id: str, aws_secret_access_key: str, endpoint_url: str, **kwargs): super().__init__(content=None, **kwargs) self.key = key self.bucket = bucket self.aws_access_key_id = aws_access_key_id self.aws_secret_access_key = aws_secret_access_key self.endpoint_url = endpoint_url async def __call__(self, scope, receive, send) -> None: session = aioboto3.Session( aws_access_key_id=self.aws_access_key_id, aws_secret_access_key=self.aws_secret_access_key ) async with session.client("s3", use_ssl=False, endpoint_url=self.endpoint_url, verify=False) as s3: resp = await s3.get_object(Bucket=self.bucket, Key=self.key) total_size = resp["ContentLength"] sent_size = 0 await send({ "type": "http.response.start", "status": self.status_code, "headers": self.raw_headers, }) async for chunk in resp["Body"].iter_chunks(chunk_size=self.chunk_size): sent_size += len(chunk) await send({ "type": "http.response.body", "body": chunk, "more_body": sent_size < total_size, }) if self.background is not None: await self.background() # 使用示例 @app.get('/custom-stream/') async def custom_stream(): return S3FileStreamingResponse( key='facd0b48-312b-4cf8-8513-7ad15a8f7bae/2-2.xlsx', bucket='my-backups-12345', aws_access_key_id='tlBne2gPKidWDkKx', aws_secret_access_key='p8ClBh5h0nsjb98MeDs6k5l5gI0OZ1bB', endpoint_url='http://localhost:9000' )
补充说明
iter_chunks()本身的用法是正确的,问题出在客户端生命周期管理上,而非迭代器本身。- 流式返回的优势是不需要把整个文件加载到内存,适合大文件场景,你的思路是对的。
内容的提问来源于stack exchange,提问作者Альберт Александров
相关产品推荐
相关产品推荐

