FastAPI流式上传大文件:仅获取指定大小的文件二进制数据块
实现FastAPI大文件流式上传并提取指定大小的纯二进制数据块
我需要创建一个FastAPI端点,实现大文件的流式上传,并且要获取指定大小(CHUNK_SIZE = 1024 * 1024 * 20)的文件二进制数据块,而非包含表单边界及头信息的内容。
当前使用async for chunk in request.stream()获取的结果包含表单边界和头信息,示例如下:
b'--8af8cd2412f5572588d978d2a7c51016\r\nContent-Disposition: form-data; name="file"; filename="toto.txt"\r\nContent-Type: text/plain\r\n\r\n' b'this is my file data! (Can be huge data)\r\n--8af8cd2412f5572588d978d2a7c51016--\r\n'
而我期望只得到文件的二进制数据,示例如下:
b'this is my file data! (Can be huge data)' # 大文件时为 b'...' * 999
我尝试过streaming_form_data但没有成功(不需要将文件写入硬盘),也尝试过自行编写表单数据解析器,均无效。
解决方案:使用python-multipart实现纯二进制块提取
通过python-multipart解析表单数据,配合线程池优化回调逻辑,实现指定大小的文件二进制数据块处理,无需写入硬盘。
步骤1:安装依赖
pip install python-multipart
步骤2:完整代码实现
from concurrent.futures import ThreadPoolExecutor from multipart import MultipartParser from fastapi import Request from fastapi.responses import JSONResponse CHUNK_SIZE = 1024 * 1024 * 20 # 20MB class UploadFileByStream: def __init__(self): self.buffer = None self.bytes_received = None self.executor = ThreadPoolExecutor() self.results = [] self.futures = [] def on_part_begin(self): # 开始处理文件部分时初始化缓冲区 self.buffer = b"" self.bytes_received = 0 def on_part_data(self, data, start, end): # 累加数据到缓冲区 self.buffer += data[start:end] self.bytes_received += end - start # 当缓冲区达到指定大小,提交处理任务 if self.bytes_received >= CHUNK_SIZE: self.send_buffer() def on_part_end(self): # 文件部分结束时,处理剩余的缓冲区数据 if self.buffer: self.send_buffer() def send_buffer(self): # 提交缓冲区数据到线程池处理(替换YOUR_FUNCTION为实际处理逻辑) future = self.executor.submit(YOUR_FUNCTION, self.buffer) self.futures.append(future) # 重置缓冲区 self.buffer = b"" self.bytes_received = 0 def collect_result(self): # 收集所有处理任务的结果 for future in self.futures: try: result = future.result() self.results.extend(result) except Exception as e: print("处理数据块时出错:", e) @app.post("/test/upload") async def upload_test(request: Request): filename = request.headers.get("file") content_type = request.headers.get("Content-Type") # 验证Content-Type是否包含boundary if not content_type or "boundary=" not in content_type: return JSONResponse(content="无效的Content-Type请求头", status_code=400) # 提取边界值 _, boundary = content_type.split("boundary=") file_processor = UploadFileByStream() # 配置解析回调函数 callbacks = { 'on_part_begin': file_processor.on_part_begin, 'on_part_data': file_processor.on_part_data, 'on_part_end': file_processor.on_part_end, } parser = MultipartParser(boundary, callbacks) # 流式读取请求数据并解析 async for chunk in request.stream(): parser.write(chunk) # 收集处理结果 file_processor.collect_result() return {"filename": filename}
代码说明
UploadFileByStream类负责缓冲区管理和数据块处理:on_part_begin:初始化缓冲区,准备接收文件数据on_part_data:累加数据到缓冲区,达到指定大小后提交线程池处理on_part_end:处理文件结束后剩余的缓冲区数据send_buffer:将缓冲区数据提交到线程池执行自定义处理逻辑(需替换YOUR_FUNCTION为实际业务代码)collect_result:收集所有异步任务的处理结果
端点逻辑:解析请求头中的boundary,流式读取请求体并交给
MultipartParser解析,自动提取纯文件二进制数据。
内容的提问来源于stack exchange,提问作者YoYoDev
相关产品推荐
相关产品推荐

