From 1313934289eb43733fea742062c17f5791b0a546 Mon Sep 17 00:00:00 2001 From: jerryliang Date: Fri, 16 Jan 2026 15:45:12 +0800 Subject: [PATCH] =?UTF-8?q?feat:=20=E4=B8=BA=E6=89=80=E6=9C=89=E5=AD=98?= =?UTF-8?q?=E5=82=A8=E5=90=8E=E7=AB=AF=E5=AE=9E=E7=8E=B0=E6=B5=81=E5=BC=8F?= =?UTF-8?q?=E6=96=87=E4=BB=B6=E4=B8=8B=E8=BD=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将文件下载方式从全量读取改为流式传输,减少内存占用并提高大文件下载性能 支持S3、OneDrive、OpenDAL和WebDAV存储后端的流式下载 添加错误处理和HTTP异常捕获 --- core/storage.py | 154 +++++++++++++++++++++++++++++++++--------------- 1 file changed, 105 insertions(+), 49 deletions(-) diff --git a/core/storage.py b/core/storage.py index 0e820bb..a459797 100644 --- a/core/storage.py +++ b/core/storage.py @@ -23,7 +23,7 @@ from core.response import APIResponse from core.settings import data_root, settings from apps.base.models import FileCodes, UploadChunk from core.utils import get_file_url, sanitize_filename -from fastapi.responses import FileResponse +from fastapi.responses import FileResponse, StreamingResponse class FileStorageInterface: @@ -310,20 +310,34 @@ class S3FileStorage(FileStorageInterface): }, ExpiresIn=3600, ) - tmp = io.BytesIO() - async with aiohttp.ClientSession() as session: - async with session.get(link) as resp: - tmp.write(await resp.read()) - tmp.seek(0) - content = tmp.read() - tmp.close() - return Response( - content, + + async def stream_generator(): + async with aiohttp.ClientSession() as session: + async with session.get(link) as resp: + if resp.status != 200: + raise HTTPException( + status_code=resp.status, + detail=f"从S3获取文件失败: {resp.status}" + ) + # 设置块大小(例如64KB) + chunk_size = 65536 + while True: + chunk = await resp.content.read(chunk_size) + if not chunk: + break + yield chunk + + from fastapi.responses import StreamingResponse + headers = { + "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode("latin-1")}"' + } + return StreamingResponse( + stream_generator(), media_type="application/octet-stream", - headers={ - "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode("latin-1")}"' - }, + headers=headers ) + except HTTPException: + raise except Exception: raise HTTPException(status_code=503, detail="服务代理下载异常,请稍后再试") @@ -602,20 +616,32 @@ class OneDriveFileStorage(FileStorageInterface): link = await asyncio.to_thread( self._get_file_url, await file_code.get_file_path(), filename ) - tmp = io.BytesIO() - async with aiohttp.ClientSession() as session: - async with session.get(link) as resp: - tmp.write(await resp.read()) - tmp.seek(0) - content = tmp.read() - tmp.close() - return Response( - content, + + async def stream_generator(): + async with aiohttp.ClientSession() as session: + async with session.get(link) as resp: + if resp.status != 200: + raise HTTPException( + status_code=resp.status, + detail=f"从OneDrive获取文件失败: {resp.status}" + ) + chunk_size = 65536 + while True: + chunk = await resp.content.read(chunk_size) + if not chunk: + break + yield chunk + + headers = { + "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode("latin-1")}"' + } + return StreamingResponse( + stream_generator(), media_type="application/octet-stream", - headers={ - "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode("latin-1")}"' - }, + headers=headers ) + except HTTPException: + raise except Exception: raise HTTPException(status_code=503, detail="服务代理下载异常,请稍后再试") @@ -776,11 +802,35 @@ class OpenDALFileStorage(FileStorageInterface): async def get_file_response(self, file_code: FileCodes): try: filename = file_code.prefix + file_code.suffix - content = await self.operator.read(await file_code.get_file_path()) + # 尝试使用流式读取器 + try: + # OpenDAL 可能提供 reader 方法返回一个异步读取器 + reader = await self.operator.reader(await file_code.get_file_path()) + except AttributeError: + # 如果 reader 方法不存在,回退到全量读取(兼容旧版本) + content = await self.operator.read(await file_code.get_file_path()) + headers = { + "Content-Disposition": f'attachment; filename="{filename}"' + } + return Response( + content, headers=headers, media_type="application/octet-stream" + ) + + async def stream_generator(): + chunk_size = 65536 + while True: + chunk = await reader.read(chunk_size) + if not chunk: + break + yield chunk + headers = { - "Content-Disposition": f'attachment; filename="{filename}"'} - return Response( - content, headers=headers, media_type="application/octet-stream" + "Content-Disposition": f'attachment; filename="{filename}"' + } + return StreamingResponse( + stream_generator(), + media_type="application/octet-stream", + headers=headers ) except Exception as e: logger.info(e) @@ -969,26 +1019,32 @@ class WebDAVFileStorage(FileStorageInterface): try: filename = file_code.prefix + file_code.suffix url = self._build_url(await file_code.get_file_path()) - async with aiohttp.ClientSession(headers={ - "Authorization": f"Basic {base64.b64encode(f'{settings.webdav_username}:{settings.webdav_password}'.encode()).decode()}" - }) as session: - async with session.get(url) as resp: - if resp.status != 200: - raise HTTPException( - status_code=resp.status, - detail=f"文件获取失败{resp.status}: {await resp.text()}", - ) - # 读取内容到内存 - content = await resp.read() - return Response( - content=content, - media_type=resp.headers.get( - "Content-Type", "application/octet-stream" - ), - headers={ - "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode()}"' - }, - ) + + async def stream_generator(): + async with aiohttp.ClientSession(headers={ + "Authorization": f"Basic {base64.b64encode(f'{settings.webdav_username}:{settings.webdav_password}'.encode()).decode()}" + }) as session: + async with session.get(url) as resp: + if resp.status != 200: + raise HTTPException( + status_code=resp.status, + detail=f"文件获取失败{resp.status}: {await resp.text()}", + ) + chunk_size = 65536 + while True: + chunk = await resp.content.read(chunk_size) + if not chunk: + break + yield chunk + + headers = { + "Content-Disposition": f'attachment; filename="{filename.encode("utf-8").decode()}"' + } + return StreamingResponse( + stream_generator(), + media_type="application/octet-stream", + headers=headers + ) except aiohttp.ClientError as e: raise HTTPException( status_code=503, detail=f"WebDAV连接异常: {str(e)}")