Python 基础体系 · 第 92/112 篇。示例统一以 Python 3.14 为语言基线;第三方库使用与其兼容的现代稳定版本,版本敏感行为会单独说明。
Python 文件上传与后台处理:流式、校验、对象存储和任务状态
文件上传看似只是“接收一个 UploadFile,保存到磁盘”,但一旦文件可能达到数百 MB、上传过程需要校验、处理耗时数分钟,系统就同时面对四个问题:
- 如何接收文件而不把完整内容读入内存;
- 如何证明收到的内容符合大小、类型、完整性和安全要求;
- 如何把原始文件可靠地交给对象存储;
- 如何让客户端查询后台处理进度,并正确处理失败、重试和重复执行。
一个可扩展的上传系统通常不是“请求函数直接处理文件”,而是下面这样的数据流:
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 对象传入。文件大小为 字节时,至少需要为内容分配约 字节内存;如果还要复制、解压或计算多个中间结果,峰值内存可能更高。
因此:
contents = await file.read()
并不是“流式处理”。它只是一次性读取。
FastAPI 官方文档也明确区分了 bytes 和 UploadFile:bytes 会把完整内容放入内存,适合小文件;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 服务器仍应配置更早的请求限制。
二、流式上传的正确内存模型
设:
- :文件总大小;
- :单个块大小;
- :固定缓冲和库内部额外开销。
一次性读取的主要内存复杂度接近:
分块读取并立即写入目标的主要内存复杂度接近:
当 时,峰值内存不再随文件总大小线性增长。
但这只成立于“每个块处理完就释放”的实现。例如下面的代码破坏了流式性质:
chunks = []
async for chunk in request.stream():
chunks.append(chunk)
data = b"".join(chunks)
虽然网络层是逐块读取的,应用最终仍然构造了完整文件,峰值内存重新接近 。
下面的代码也可能产生不必要的复制:
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")
最终的实际大小是:
其中 是第 个收到的字节块, 才是服务端真正接收的文件大小。
2. Content-Type 是声明,不是鉴定结果
攻击者可以把任意文件声明为:
Content-Type: image/jpeg
所以至少要分三层判断:
- 扩展名检查:便于用户理解和下载;
- MIME 声明检查:便于接口层快速拒绝;
- 文件签名或解析检查:判断内容是否真的符合格式。
例如 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()
其计算过程可以表示为:
最终:
其中:
- 是哈希算法的初始状态;
- 是第 个数据块;
- 是哈希压缩函数;
- 是整个文件的摘要。
哈希值可以用于:
- 上传完成后的完整性校验;
- 内容寻址;
- 去重;
- 幂等键;
- 后台任务输入版本标识。
但哈希值不是授权凭证。客户端提交一个 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. 幂等性的形式化条件
设任务函数为:
如果重复执行同一个输入不会造成额外副作用,则要求:
这里的等价不是指 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 内置状态包括 PENDING、STARTED、SUCCESS、FAILURE 和 RETRY。其中 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 参数,旧式 startup 和 shutdown 事件处理器不会同时执行;两种方式应择一使用。(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 完整学习路线:从语言模型、并发到 Web、数据、AI 与生产交付
- 上一篇:Python Web 认证与授权:Session、JWT、OAuth 2.0、RBAC 和 CSRF
- 下一篇:Python 网络采集:Requests、HTTPX、BeautifulSoup、限速和合规
- 延伸:FastAPI 完整基础:路由、依赖注入、校验、异步和生命周期
- 延伸:Celery 任务队列:Broker、Worker、ACK、重试、定时和幂等
官方资料
本文依据 Python 官方文档、相关 PEP 与生态项目官方文档重新梳理;正文、示例与工程清单由 WR BLOG 编写。

评论
0 条讨论