You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 13:29:53