Python 基础体系 · 第 92/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。

Python 文件上传与后台处理:流式、校验、对象存储和任务状态

文件上传看似只是“接收一个 UploadFile,保存到磁盘”,但一旦文件可能达到数百 MB、上传过程需要校验、处理耗时数分钟,系统就同时面对四个问题:

  1. 如何接收文件而不把完整内容读入内存;
  2. 如何证明收到的内容符合大小、类型、完整性和安全要求;
  3. 如何把原始文件可靠地交给对象存储;
  4. 如何让客户端查询后台处理进度,并正确处理失败、重试和重复执行。

一个可扩展的上传系统通常不是“请求函数直接处理文件”,而是下面这样的数据流:

flowchart LR
    C[客户端] -->|上传字节流| API[FastAPI API]
    API --> V[尺寸/类型/哈希校验]
    V --> T[临时对象或暂存文件]
    T --> DB[(文件元数据与任务状态)]
    API -->|提交任务| B[Broker]
    B --> W[Celery Worker]
    W -->|读取原始对象| S[(对象存储)]
    W --> P[解析/转码/病毒扫描]
    P --> S2[(处理结果对象)]
    W --> DB
    C -->|查询状态| API
    API --> DB

这里有一个重要边界:

  • 上传请求负责接收、校验和持久化原始文件;
  • 后台任务负责耗时处理;
  • 对象存储负责保存大文件;
  • 数据库负责保存文件元数据和任务状态;
  • Broker只负责传递任务消息,不应被当成文件存储。

一、先区分三个“文件接收”模型

1. bytes:简单,但会把完整文件放进内存

FastAPI 可以使用 bytes 声明上传参数:

from typing import Annotated

from fastapi import FastAPI, File

app = FastAPI()


@app.post("/files")
async def upload(file: Annotated[bytes, File()]):
    return {"size": len(file)}

这种方式的语义很直接:框架在调用路径函数前,先取得整个文件内容,然后把它作为一个 bytes 对象传入。文件大小为 NN 字节时,至少需要为内容分配约 NN 字节内存;如果还要复制、解压或计算多个中间结果,峰值内存可能更高。

因此:

contents = await file.read()

并不是“流式处理”。它只是一次性读取。

FastAPI 官方文档也明确区分了 bytesUploadFilebytes 会把完整内容放入内存,适合小文件;UploadFile 使用可在内存和磁盘之间切换的 spooled file,更适合较大文件。(fastapi.tiangolo.com)


2. UploadFile:适合 multipart 表单,但不等于端到端网络流式

浏览器最常见的上传形式是:

Content-Type: multipart/form-data

对应的 FastAPI 代码是:

from fastapi import FastAPI, UploadFile

app = FastAPI()


@app.post("/upload")
async def upload(file: UploadFile):
    return {
        "filename": file.filename,
        "content_type": file.content_type,
    }

接收 multipart 文件还需要安装 python-multipart,因为框架需要解析表单边界、字段名和文件内容。(fastapi.tiangolo.com)

UploadFile 提供:

  • filename:客户端提交的原始文件名;
  • content_type:客户端声明的 MIME 类型;
  • file:底层文件对象;
  • read(size)write(data)seek(offset)close() 等异步接口。

例如分块复制:

from pathlib import Path

from fastapi import FastAPI, HTTPException, UploadFile

app = FastAPI()
UPLOAD_DIR = Path("uploads")
UPLOAD_DIR.mkdir(exist_ok=True)

MAX_SIZE = 100 * 1024 * 1024
CHUNK_SIZE = 1024 * 1024


@app.post("/upload")
async def upload(file: UploadFile):
    destination = UPLOAD_DIR / "received.bin"
    total = 0

    try:
        with destination.open("wb") as output:
            while chunk := await file.read(CHUNK_SIZE):
                total += len(chunk)

                if total > MAX_SIZE:
                    raise HTTPException(
                        status_code=413,
                        detail="file too large",
                    )

                output.write(chunk)
    finally:
        await file.close()

    return {"size": total}

这个循环避免了在业务代码中创建完整的 bytes,但要准确理解它的边界:

UploadFile 的异步读写接口是文件对象接口。FastAPI 会在线程池中执行底层文件操作,再由异步路径函数等待结果;它不意味着应用一定是在收到网络字节的同时处理每一个字节。

FastAPI 文档说明,UploadFile 的异步方法会调用底层 SpooledTemporaryFile 的方法,并在线程池中运行。(fastapi.tiangolo.com)

所以 UploadFile 解决的是:

  • 避免业务代码一次性读取完整文件;
  • 利用临时文件承接 multipart 解析结果;
  • 方便把文件传给需要文件对象的库。

它没有自动解决:

  • 代理层是否已经缓冲完整请求;
  • multipart 解析过程是否受框架限制;
  • 上传过程是否具有断点续传能力;
  • 文件是否已经写入对象存储;
  • 任务是否已经可靠提交。

3. Request.stream():直接消费 ASGI 请求体

如果接口本身就是原始二进制上传,例如:

PUT /objects/abc
Content-Type: application/octet-stream

可以直接读取请求流:

from hashlib import sha256
from pathlib import Path

from fastapi import FastAPI, HTTPException, Request

app = FastAPI()

STAGING_DIR = Path("staging")
STAGING_DIR.mkdir(exist_ok=True)

MAX_SIZE = 100 * 1024 * 1024
CHUNK_SIZE = 1024 * 1024


@app.put("/raw-upload/{upload_id}")
async def raw_upload(upload_id: str, request: Request):
    path = STAGING_DIR / f"{upload_id}.part"
    digest = sha256()
    total = 0

    try:
        with path.open("wb") as output:
            async for chunk in request.stream():
                if not chunk:
                    continue

                total += len(chunk)

                if total > MAX_SIZE:
                    raise HTTPException(413, "file too large")

                digest.update(chunk)
                output.write(chunk)

    except Exception:
        path.unlink(missing_ok=True)
        raise

    return {
        "upload_id": upload_id,
        "size": total,
        "sha256": digest.hexdigest(),
    }

ASGI 的 HTTP 规范把请求体表示为一系列 http.request 事件。每个事件包含 body,并通过 more_body 表示后面是否还有数据;这正是应用逐块消费请求体的基础。(asgi.readthedocs.io)

Request.stream() 的核心语义是:

async for chunk in request.stream():
    ...

应用逐次得到字节块,而不是先调用:

body = await request.body()

后者会把整个请求体聚合为一个对象。

但是,流式读取不等于无限制读取。上传链路通常还包括:

客户端
  ↓
CDN / 反向代理
  ↓
ASGI 服务器
  ↓
FastAPI
  ↓
对象存储

任意一层都可能有:

  • 请求体大小限制;
  • 超时;
  • 请求缓冲;
  • 空闲连接超时;
  • 并发连接上限。

因此应用层的 MAX_SIZE 是最后一道防线,而不是唯一的大小控制。应用层通过 Request.stream() 发现超限时,可能已经接收了一部分数据;代理层和 ASGI 服务器仍应配置更早的请求限制。


二、流式上传的正确内存模型

设:

  • NN:文件总大小;
  • CC:单个块大小;
  • MM:固定缓冲和库内部额外开销。

一次性读取的主要内存复杂度接近:

O(N)O(N)

分块读取并立即写入目标的主要内存复杂度接近:

O(C+M)O(C + M)

CNC \ll N 时,峰值内存不再随文件总大小线性增长。

但这只成立于“每个块处理完就释放”的实现。例如下面的代码破坏了流式性质:

chunks = []

async for chunk in request.stream():
    chunks.append(chunk)

data = b"".join(chunks)

虽然网络层是逐块读取的,应用最终仍然构造了完整文件,峰值内存重新接近 O(N)O(N)

下面的代码也可能产生不必要的复制:

async for chunk in request.stream():
    all_data += chunk

因为不可变 bytes 的反复拼接可能反复分配和复制已有内容。上传路径应优先采用:

output.write(chunk)
digest.update(chunk)

而不是累积整个文件。


三、校验必须分为“声明校验”和“内容校验”

客户端通常会发送:

Content-Length: 1048576
Content-Type: image/jpeg

这些字段有用,但都不能单独证明文件可信。

1. Content-Length 只能作为早期拒绝条件

如果请求头给出:

Content-Length: 200000000

而服务端限制为 100 MB,可以立刻拒绝:

from fastapi import Header

@app.put("/raw-upload/{upload_id}")
async def raw_upload(
    upload_id: str,
    request: Request,
    content_length: int | None = Header(default=None),
):
    if content_length is not None and content_length > MAX_SIZE:
        raise HTTPException(413, "file too large")

    ...

但是不能假设:

content_length == 实际读取字节数

原因包括:

  • 使用 chunked 传输时可能没有 Content-Length
  • 中间层可能重新组织请求;
  • 客户端可能提供错误值;
  • 多部分表单的总长度包含边界和字段元数据,不等于文件大小。

因此必须在实际读取过程中累计:

total += len(chunk)

if total > MAX_SIZE:
    raise HTTPException(413, "file too large")

最终的实际大小是:

S=i=1kciS = \sum_{i=1}^{k} |c_i|

其中 cic_i 是第 ii 个收到的字节块,SS 才是服务端真正接收的文件大小。


2. Content-Type 是声明,不是鉴定结果

攻击者可以把任意文件声明为:

Content-Type: image/jpeg

所以至少要分三层判断:

  1. 扩展名检查:便于用户理解和下载;
  2. MIME 声明检查:便于接口层快速拒绝;
  3. 文件签名或解析检查:判断内容是否真的符合格式。

例如 JPEG 文件通常以 FF D8 FF 开头,PNG 文件以固定的 PNG 签名开头:

def detect_type(prefix: bytes) -> str | None:
    if prefix.startswith(b"\x89PNG\r\n\x1a\n"):
        return "image/png"

    if prefix.startswith(b"\xff\xd8\xff"):
        return "image/jpeg"

    return None

流式校验时只需保留前若干字节:

prefix = bytearray()

async for chunk in request.stream():
    if len(prefix) < 16:
        need = 16 - len(prefix)
        prefix.extend(chunk[:need])

    ...

注意,文件签名只能判断“开头看起来像某种格式”。它不能证明:

  • 文件结构完整;
  • 解码器不会崩溃;
  • 压缩炸弹风险可接受;
  • 文件不含恶意载荷。

对图片、压缩包、Office 文档等复杂格式,还需要在隔离进程或专用扫描服务中解析。


3. 哈希值用于完整性和幂等关联

流式计算 SHA-256:

from hashlib import sha256

digest = sha256()

async for chunk in request.stream():
    digest.update(chunk)
    output.write(chunk)

sha256_hex = digest.hexdigest()

其计算过程可以表示为:

H0=IVH_0 = IV

Hi=f(Hi1,ci)H_i = f(H_{i-1}, c_i)

最终:

H=HkH = H_k

其中:

  • IVIV 是哈希算法的初始状态;
  • cic_i 是第 ii 个数据块;
  • ff 是哈希压缩函数;
  • HH 是整个文件的摘要。

哈希值可以用于:

  • 上传完成后的完整性校验;
  • 内容寻址;
  • 去重;
  • 幂等键;
  • 后台任务输入版本标识。

但哈希值不是授权凭证。客户端提交一个 SHA-256 字符串,并不能证明它有权上传该文件;授权仍由用户身份、租户、对象键和业务权限决定。


四、临时文件必须有明确的提交语义

直接写入最终路径有一个常见故障:

写入 final.bin
进程崩溃
final.bin 只包含前 40%
后台任务读取到半个文件

因此应使用“暂存路径 + 原子提交”:

from pathlib import Path
from uuid import uuid4

async def save_stream(request: Request) -> tuple[Path, int, str]:
    temp_path = Path("staging") / f"{uuid4()}.part"
    final_path = Path("staging") / f"{uuid4()}.bin"

    total = 0
    digest = sha256()

    try:
        with temp_path.open("wb") as output:
            async for chunk in request.stream():
                if not chunk:
                    continue

                total += len(chunk)
                if total > MAX_SIZE:
                    raise HTTPException(413, "file too large")

                digest.update(chunk)
                output.write(chunk)

            output.flush()

        temp_path.replace(final_path)
        return final_path, total, digest.hexdigest()

    except Exception:
        temp_path.unlink(missing_ok=True)
        final_path.unlink(missing_ok=True)
        raise

这里的关键条件是:

只有完整写入并关闭暂存文件后,才把它暴露为“已提交文件”。

replace() 的原子性依赖源文件和目标文件位于同一文件系统;如果跨文件系统移动,可能退化为复制加删除,不能直接假设具有同样语义。

在对象存储中也应保留类似状态:

UPLOADING → UPLOADED → PROCESSING → SUCCEEDED
                         ↓
                       FAILED

UPLOADING 状态表示上传尚未提交;后台任务只能读取 UPLOADED 状态的对象。


五、对象存储解决的不是“上传校验”,而是“持久化和扩展”

对象存储把数据组织成:

bucket + object key + object bytes + metadata

以 S3 模型为例,对象键唯一标识存储桶中的对象;所谓“目录”通常只是键中的前缀,并不是传统文件系统目录。(docs.aws.amazon.com)

不要直接使用用户文件名作为对象键:

key = file.filename

原因包括:

  • 文件名可能重复;
  • 可能包含路径片段;
  • 可能导致覆盖;
  • 可能包含特殊字符;
  • 文件名不应承担租户隔离职责。

更安全的键可以由服务端生成:

tenant/{tenant_id}/uploads/{upload_id}/original
tenant/{tenant_id}/uploads/{upload_id}/derived/thumbnail.jpg

用户原始文件名只作为元数据保存:

{
  "upload_id": "01J...",
  "original_filename": "report.pdf",
  "content_type": "application/pdf",
  "size": 7340032,
  "sha256": "...",
  "object_key": "tenant/acme/uploads/01J.../original"
}

对象键的生成原则是:

业务身份由服务端生成;
展示名称保存为元数据;
对象键不信任客户端输入;
原始对象和处理结果使用不同键。

1. 服务端中转上传

最简单的数据流是:

客户端 → FastAPI → 对象存储

FastAPI 在读取每个块时,将其写入本地暂存文件,上传完成后再调用对象存储 SDK。

这种模式便于:

  • 统一鉴权;
  • 统一大小限制;
  • 在服务端计算哈希;
  • 统一记录数据库状态。

代价是应用服务器承受了完整的网络流量和磁盘暂存压力。


2. 预签名 URL 上传

大文件常用:

客户端 → FastAPI:申请上传
FastAPI → 客户端:返回预签名 URL
客户端 → 对象存储:直接上传
客户端 → FastAPI:确认上传

预签名 URL 是带有有限权限和有效期的授权 URL。AWS 文档说明,使用预签名 URL 可以让调用方在没有云平台长期凭证的情况下上传对象,但权限受签发者权限限制;如果相同对象键已存在,上传可能替换原对象。(docs.aws.amazon.com)

因此,申请接口应由服务端生成不可猜测、不会与其他文件冲突的对象键:

@app.post("/uploads")
async def create_upload(user: CurrentUser):
    upload_id = generate_id()
    key = f"tenant/{user.tenant_id}/uploads/{upload_id}/original"

    # 实际项目中由对象存储 SDK 生成预签名 URL
    return {
        "upload_id": upload_id,
        "object_key": key,
        "upload_url": "...",
    }

客户端上传成功后,不能仅凭客户端说“上传完成”就把状态改成 UPLOADED。服务端应通过对象存储的 HEAD 或元数据接口确认:

  • 对象存在;
  • 实际大小符合预期;
  • 内容类型符合允许范围;
  • 哈希或校验值符合要求;
  • 对象键属于当前租户。

3. 分段上传

对象存储的 multipart upload 将一个对象拆成多个连续分段,每个分段可以独立上传;某一段失败时只需要重传该段,而不必重传整个对象。AWS 文档还说明,完成分段上传时需要提交分段编号和对应的 ETag;使用额外校验和时,分段编号还必须从 1 开始连续排列。(docs.aws.amazon.com)

分段上传的状态可以这样保存:

INITIATED
  ↓
PART 1 uploaded
PART 2 uploaded
PART 3 failed → retry PART 3
  ↓
COMPLETING
  ↓
COMPLETED

不要把 ETag 简单当成整个对象的 MD5。对于 multipart upload,完成后的 ETag 不一定是对象内容的 MD5。对象完整性应使用对象存储支持的校验和机制,或由服务端重新计算哈希。(docs.aws.amazon.com)


六、后台处理:BackgroundTasks 和 Celery 不是同一种保证

1. FastAPI BackgroundTasks

FastAPI 的 BackgroundTasks 用于在响应返回后执行任务:

from fastapi import BackgroundTasks, FastAPI

app = FastAPI()


def write_audit_log(upload_id: str) -> None:
    with open("audit.log", "a", encoding="utf-8") as output:
        output.write(f"{upload_id}\n")


@app.post("/uploads/{upload_id}/finalize")
async def finalize(upload_id: str, background_tasks: BackgroundTasks):
    background_tasks.add_task(write_audit_log, upload_id)
    return {"upload_id": upload_id, "status": "accepted"}

它适合:

  • 写一条通知;
  • 发送轻量邮件;
  • 做短时间的本地后处理;
  • 不需要独立进程、持久队列和跨机器调度的任务。

它不适合承载关键的大型处理任务,因为任务与当前 Web 进程生命周期相关。进程崩溃、容器重启或部署替换时,尚未执行或正在执行的任务可能丢失。

FastAPI 官方文档也把 BackgroundTasks 定位为小型后台操作,并建议对于需要多进程、多服务器或重计算的场景使用 Celery 等独立任务系统。(fastapi.tiangolo.com)

特别要避免:

background_tasks.add_task(process_file, file.file)

请求结束后,上传文件对象可能已经被关闭或清理。后台任务应接收稳定的引用:

background_tasks.add_task(process_file, object_key)

任务自己根据 object_key 重新打开对象。


2. Celery 的四个角色

Celery 的基本数据流是:

FastAPI
  ↓ publish message
Broker
  ↓ deliver message
Worker
  ↓ execute task
Result backend / database
  • Broker:传递任务消息,例如 Redis 或 RabbitMQ;
  • Worker:从 Broker 取消息并执行任务;
  • ACK:Worker 向 Broker 确认消息已处理或已接管;
  • Result backend:保存任务状态和结果。

Celery 任务消息应只包含小型、可序列化的标识:

process_upload.delay(upload_id)

而不是:

process_upload.delay(file_bytes)

把大文件放进 Broker 会导致消息体膨胀、复制增多、Broker 内存压力上升,也使任务重试变得昂贵。文件内容应放在对象存储,消息只携带:

{
  "upload_id": "...",
  "object_key": "...",
  "expected_sha256": "...",
  "schema_version": 1
}

七、ACK、重试和幂等必须一起设计

1. ACK 不是“任务成功”的同义词

ACK 是消息消费协议中的确认。它回答的是:

这个消息是否已经被 Worker 接管并从 Broker 的投递流程中确认?

它不自动回答:

文件是否处理成功?

Celery 默认会在任务执行前确认消息;如果 Worker 在执行中崩溃,消息可能不会再次执行。对于幂等任务,可以配置 acks_late=True,让确认延后到任务执行完成之后,但这样 Worker 在执行中崩溃时,任务可能被再次执行。Celery 文档明确要求:使用延迟 ACK 时,任务必须具有幂等性。(docs.celeryq.dev)


2. 幂等性的形式化条件

设任务函数为:

F(x)F(x)

如果重复执行同一个输入不会造成额外副作用,则要求:

F(F(x))F(x)F(F(x)) \equiv F(x)

这里的等价不是指 Python 返回值对象相同,而是指外部可观察结果相同。

非幂等示例:

@app.task
def add_balance(user_id: str, amount: int):
    db.execute(
        "UPDATE account SET balance = balance + ? WHERE user_id = ?",
        (amount, user_id),
    )

任务执行两次,余额增加两次。

幂等改写可以使用任务唯一键:

@app.task
def add_balance_once(operation_id: str, user_id: str, amount: int):
    with db.transaction():
        inserted = db.execute(
            """
            INSERT INTO applied_operations(operation_id)
            VALUES (?)
            ON CONFLICT DO NOTHING
            """,
            (operation_id,),
        )

        if inserted.rowcount == 0:
            return "already-applied"

        db.execute(
            "UPDATE account SET balance = balance + ? WHERE user_id = ?",
            (amount, user_id),
        )

    return "applied"

第一次执行:

插入 operation_id 成功
更新余额
提交事务

第二次执行:

插入 operation_id 冲突
跳过余额更新
返回 already-applied

文件处理通常可以设计为“确定性输出键”:

derived/{upload_id}/thumbnail.jpg

Worker 先检查输出对象是否已经存在;如果存在,再验证它对应的输入版本或哈希,确认后直接返回。这样即使任务重复执行,也不会不断产生新的结果对象。


3. 只对可恢复错误重试

网络超时、对象存储临时不可用、数据库暂时连接失败,通常具有重试价值:

from celery import Celery

celery_app = Celery(
    "worker",
    broker="redis://localhost:6379/0",
    backend="redis://localhost:6379/1",
)


@celery_app.task(
    bind=True,
    autoretry_for=(TimeoutError, ConnectionError),
    retry_backoff=True,
    retry_jitter=True,
    max_retries=5,
)
def process_upload(self, upload_id: str):
    ...

Celery 支持 autoretry_for、指数退避、抖动和最大重试次数。指数退避的直觉是:第一次失败后短暂等待,连续失败时逐渐延长等待;抖动则避免大量任务在相同时间同时重试。(docs.celeryq.dev)

不应重试:

  • 文件格式非法;
  • 用户无权访问;
  • 文件超过业务限制;
  • 病毒扫描失败且结论明确;
  • 代码 bug 导致的确定性异常。

否则会形成无意义的重试循环。


八、任务状态应该由业务状态和队列状态共同组成

Celery 内置状态包括 PENDINGSTARTEDSUCCESSFAILURERETRY。其中 PENDING 既可能表示“等待执行”,也可能表示“这个任务 ID 对结果后端而言未知”;它不是严格意义上的“排队中”。(docs.celeryq.dev)

因此面向客户端的状态不应直接暴露 Celery 原始状态,而应建立业务状态:

CREATED
  ↓
UPLOADING
  ↓
UPLOADED
  ↓
QUEUED
  ↓
PROCESSING
  ├──→ SUCCEEDED
  ├──→ FAILED
  └──→ CANCELLED

推荐保存类似记录:

CREATE TABLE uploads (
    upload_id TEXT PRIMARY KEY,
    tenant_id TEXT NOT NULL,
    object_key TEXT NOT NULL,
    original_filename TEXT NOT NULL,
    content_type TEXT,
    size_bytes INTEGER,
    sha256 TEXT,
    status TEXT NOT NULL,
    task_id TEXT,
    error_code TEXT,
    error_message TEXT,
    created_at TEXT NOT NULL,
    updated_at TEXT NOT NULL
);

客户端查询接口:

from pydantic import BaseModel

class UploadStatus(BaseModel):
    upload_id: str
    status: str
    size_bytes: int | None = None
    sha256: str | None = None
    error_code: str | None = None


@app.get("/uploads/{upload_id}", response_model=UploadStatus)
async def get_upload(upload_id: str, user: CurrentUser):
    row = find_upload_for_tenant(upload_id, user.tenant_id)

    if row is None:
        raise HTTPException(404, "upload not found")

    return row

这里的租户条件不能省略:

WHERE upload_id = :upload_id
  AND tenant_id = :tenant_id

否则用户只要猜中或获得另一个 upload_id,就可能读取其他租户的文件状态。


九、一个完整的提交接口

下面是一个使用本地暂存文件的简化端到端示例。它展示的是控制流和状态边界,生产环境可以把暂存文件替换为对象存储上传。

from contextlib import asynccontextmanager
from hashlib import sha256
from pathlib import Path
from uuid import uuid4

from fastapi import FastAPI, HTTPException, Request, status

STAGING = Path("staging")
STAGING.mkdir(exist_ok=True)

MAX_SIZE = 100 * 1024 * 1024
CHUNK_SIZE = 1024 * 1024


def create_upload_record(upload_id: str, object_key: str) -> None:
    # 实际项目中写入数据库:
    # status = UPLOADING
    pass


def mark_uploaded(
    upload_id: str,
    *,
    size_bytes: int,
    sha256_hex: str,
) -> None:
    # 实际项目中使用事务:
    # 只有当前状态仍为 UPLOADING 时,才更新为 UPLOADED
    pass


def enqueue_processing(upload_id: str) -> str:
    # 实际项目中:
    # task = process_upload.delay(upload_id)
    # return task.id
    return f"local-task-{upload_id}"


@asynccontextmanager
async def lifespan(app: FastAPI):
    # 初始化数据库连接池、对象存储客户端等共享资源
    yield
    # 关闭连接池和客户端


app = FastAPI(lifespan=lifespan)


@app.put("/uploads/{upload_id}")
async def upload(upload_id: str, request: Request):
    object_key = f"original/{upload_id}"
    create_upload_record(upload_id, object_key)

    temp_path = STAGING / f"{uuid4()}.part"
    digest = sha256()
    total = 0

    try:
        with temp_path.open("wb") as output:
            async for chunk in request.stream():
                if not chunk:
                    continue

                total += len(chunk)
                if total > MAX_SIZE:
                    raise HTTPException(
                        status_code=status.HTTP_413_REQUEST_ENTITY_TOO_LARGE,
                        detail="file too large",
                    )

                digest.update(chunk)
                output.write(chunk)

            output.flush()

        # 此处应先将 temp_path 提交到对象存储,
        # 并确认对象存在、大小和校验值符合预期。
        mark_uploaded(
            upload_id,
            size_bytes=total,
            sha256_hex=digest.hexdigest(),
        )

        task_id = enqueue_processing(upload_id)

        return {
            "upload_id": upload_id,
            "status": "QUEUED",
            "task_id": task_id,
        }

    except HTTPException:
        temp_path.unlink(missing_ok=True)
        raise
    except Exception:
        temp_path.unlink(missing_ok=True)
        # 实际项目中将数据库状态更新为 FAILED
        raise HTTPException(500, "upload failed")

这个接口的关键顺序是:

1. 创建 UPLOADING 记录
2. 分块接收
3. 累计实际大小
4. 增量计算哈希
5. 写入暂存位置
6. 完成对象持久化
7. 更新为 UPLOADED
8. 提交后台任务
9. 返回 QUEUED

不能在第 5 步后直接把数据库更新为 UPLOADED,因为本地暂存文件可能尚未成功上传到对象存储。

也不能先提交 Celery 任务,再写入对象存储,否则 Worker 可能先读到“不存在”或“不完整”的对象。


十、数据库和消息队列之间存在“双写”问题

下面的代码看似合理:

mark_uploaded(upload_id)
process_upload.delay(upload_id)

但它包含两个独立操作:

数据库提交成功
消息发布失败

此时文件状态是 UPLOADED,却没有后台任务。

反过来也可能发生:

消息发布成功
数据库事务回滚

Worker 收到任务后找不到数据库记录。

常见解决方案是 事务消息表,也称 Outbox 模式

CREATE TABLE outbox (
    event_id TEXT PRIMARY KEY,
    event_type TEXT NOT NULL,
    aggregate_id TEXT NOT NULL,
    payload TEXT NOT NULL,
    published_at TEXT
);

在同一个数据库事务中:

1. 更新 uploads.status = 'UPLOADED'
2. 插入 outbox 事件
3. 提交事务

独立发布器不断读取未发布事件:

读取 outbox
  ↓
发布 Celery 消息
  ↓
发布成功
  ↓
设置 published_at

如果发布器崩溃:

消息已经发布,但 published_at 尚未写入

恢复后可能再次发布同一事件。因此消费者仍必须幂等。Outbox 解决的是“尽量不丢消息”,不是让消息系统和数据库凭空获得全局原子事务。


十一、生命周期资源不能靠模块全局变量随意管理

对象存储客户端、数据库连接池和扫描器客户端通常是跨请求共享的资源。FastAPI 推荐使用 lifespan 管理应用启动和关闭:yield 之前执行初始化,yield 之后执行清理;应用开始接收请求前会完成启动逻辑,关闭时执行清理逻辑。(fastapi.tiangolo.com)

from contextlib import asynccontextmanager

from fastapi import FastAPI


@asynccontextmanager
async def lifespan(app: FastAPI):
    app.state.storage = create_storage_client()
    app.state.db = create_database_pool()

    yield

    await app.state.db.close()
    await app.state.storage.close()


app = FastAPI(lifespan=lifespan)

需要注意:

  • 每个 Web Worker 都可能拥有自己的一份客户端和连接池;
  • Celery Worker 是独立进程,不应直接复用 FastAPI 进程中的对象;
  • 任务函数应在 Worker 自己的生命周期中建立或获取资源;
  • 连接池大小必须结合 Web 并发数、Worker 数和数据库上限计算。

FastAPI 文档还说明,如果使用 lifespan 参数,旧式 startupshutdown 事件处理器不会同时执行;两种方式应择一使用。(fastapi.tiangolo.com)


十二、异步并不自动等于“不阻塞”

async def 只表示函数可以暂停并把控制权交回事件循环。它不能把同步阻塞代码自动变成异步:

@app.post("/bad")
async def bad():
    data = blocking_library.read_all()
    result = cpu_heavy_parse(data)
    return result

如果 blocking_library.read_all()cpu_heavy_parse() 长时间运行,当前事件循环线程会被占住,其他请求无法及时推进。

Python 的 asyncio 适合高层次的 I/O 并发;对于阻塞代码可以使用线程或进程池,但 CPU 密集型 Python 代码在通常的 CPython 构建中仍受 GIL 影响,线程不等于多核并行。Python 3.14 文档同时说明了 free-threaded 构建的支持情况,但这不是默认运行模式,不能把它当成普通部署环境的前提。(docs.python.org)

上传路径通常是 I/O 密集型,适合分块读取和写入;视频转码、压缩、OCR、病毒扫描等重任务则应放入独立 Worker,必要时再使用独立进程、专用服务或具备资源隔离的执行环境。


十三、失败路径必须有清理动作

上传和处理至少要考虑以下失败点:

阶段 失败表现 应有动作
读取请求体 客户端断开、超时 删除暂存文件,记录失败原因
大小校验 超过上限 立即停止读取,删除暂存文件
类型校验 MIME 或文件签名不符 不提交对象,不创建处理任务
对象存储上传 网络错误 删除或标记未完成的对象,允许重试
数据库提交 事务失败 不返回成功状态,等待客户端重试
消息发布 Broker 不可用 使用 outbox 或返回可重试错误
Worker 处理 临时错误 重试;任务必须幂等
Worker 处理 永久错误 FAILED,保存稳定的错误码
客户端断开 响应发送失败 不一定代表服务端上传失败,查询数据库确认

客户端断开尤其容易误判。ASGI 规范定义了 http.disconnect 事件;如果应用继续向已关闭连接发送数据,服务器可能抛出 OSError 子类异常。(asgi.readthedocs.io)

例如客户端在上传完成后立即断网,服务端可能已经把对象写入成功,但响应没有成功返回。此时客户端不应简单地重新创建一个全新文件,而应使用 upload_id 查询状态,必要时继续或重试同一上传。


十四、诊断时要同时看四条证据链

遇到“上传成功但任务没处理”时,不要只看 HTTP 日志,应按以下顺序检查:

1. API:请求是否读完?实际字节数是多少?
2. 对象存储:对象是否存在?大小和哈希是否一致?
3. 数据库:状态是否从 UPLOADING 进入 UPLOADED?
4. Broker/Worker:消息是否发布、接收、重试或失败?

建议每个日志事件都带上:

upload_id
task_id
tenant_id
object_key
sha256
status

不要只记录用户文件名,因为同名文件可以属于不同上传,也可能被覆盖或重复提交。

一种有用的状态查询结果是:

{
  "upload_id": "u_123",
  "status": "PROCESSING",
  "progress": {
    "current": 3,
    "total": 5
  },
  "object": {
    "key": "tenant/acme/uploads/u_123/original",
    "size_bytes": 7340032,
    "sha256": "..."
  },
  "task_id": "celery-task-id"
}

进度值也必须定义清楚。3/5 可以表示“已处理 3 个阶段”,不应伪装成精确的百分比。若任务无法准确估计剩余工作,使用阶段状态比虚假的 87% 更可靠。


十五、常见误区

误区一:UploadFile 就是网络层零拷贝流

不是。它是适合文件上传的文件对象抽象,底层可能使用内存和磁盘之间的 spool,也可能经过 multipart 解析和线程池文件操作。真正需要按 ASGI 请求事件消费时,应使用 Request.stream()

误区二:限制了 Content-Length 就完成了大小校验

不是。Content-Length 可以缺失或不可信;必须累计实际收到的字节数。

误区三:扩展名是 .jpg,文件就是 JPEG

不是。扩展名和 Content-Type 都是弱证据,至少还应检查文件签名,并在隔离环境中进行真正解析。

误区四:后台任务只要返回 202 就可靠了

不是。HTTP 202 Accepted 只表示请求已被接受,不能证明任务已持久化、消息已发布或处理最终会成功。

误区五:任务重试会自动避免重复副作用

不是。重试、延迟 ACK、Worker 崩溃和网络超时都可能造成重复执行;幂等键、唯一约束和确定性输出必须由业务代码设计。

误区六:对象存储键就是文件系统路径

不是。对象存储通常是键值模型,前缀只是逻辑组织方式;对象键仍需由服务端生成,并进行租户隔离和覆盖控制。(docs.aws.amazon.com)


一个完整的文件上传系统,可以用一句话概括其正确边界:

请求负责接收和校验,
对象存储负责保存字节,
数据库负责记录事实,
Broker 负责传递任务,
Worker 负责处理,
幂等设计负责承受重复,
状态机负责向客户端解释结果。

只要这几个职责没有混在一起,流式上传、对象存储和后台任务就能分别扩展;一旦把文件内容塞进消息、把任务状态藏在进程内存、把客户端文件名当成对象键,系统通常会在大文件、进程重启或网络故障下暴露问题。


系列导航与关联阅读

官方资料

本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。