直接用boto3+gridfsbucket会压垮mongodb,因高并发调用open_upload_stream()且复用同一mongoclient,导致连接池打满、chunk写入乱序及元数据与chunks不一致;s3流式读取未设超时或未分块迭代,易卡死线程或触发oom。

为什么直接用 boto3 + GridFSBucket 会压垮 MongoDB
常见错误是开几十个线程并发调用 bucket.open_upload_stream(),每个线程还复用同一个 MongoClient。结果不是上传快,而是 MongoDB 连接池打满、chunk 写入乱序、部分文件元数据写入成功但 chunks 全丢——因为驱动在高并发下对同一 socket 复用 chunk 请求,触发状态竞争。
更隐蔽的问题是 S3 的 get_object() 默认不设流式读取超时,遇到网络抖动就卡住整个线程,后续所有上传都排队等待;而 boto3 的 StreamingBody 若没显式调用 iter_chunks() 或配 chunk_size,底层可能缓存整块响应体到内存,100MB 文件直接 OOM。
- 必须为每个线程创建独立
MongoClient,且maxPoolSize=1 - S3 客户端要设
config=Config(read_timeout=60, retries={"max_attempts": 2}) - 禁止用
response["Body"].read(),改用response["Body"].iter_chunks(chunk_size=256*1024) - GridFSBucket 必须显式传
chunk_size_bytes=1048576(1MB),避免默认 255KB 导致 GB 文件生成上万 chunk,拖慢索引更新
如何实现带进度感知的失败重试
单纯捕获异常后重试整个文件,会导致重复上传、磁盘浪费、元数据冲突。真正可控的做法是把 S3 对象的 ETag(即 MD5)和文件大小拼成唯一 key,作为 metadata.upload_id 写入 GridFS,并在上传前先查是否已存在同名且完整(length == expected_length)的记录。
- 重试逻辑只针对
NotPrimaryError、MongoNetworkError和超时,其他错误(如权限拒绝、S3 404)直接跳过 - 失败后立即执行
bucket.find({"metadata.upload_id": upload_id, "length": {"$lt": expected_length}}),若命中则用bucket.open_uploadStreamWithId()续传 - 续传时从 S3 拉取需跳过的字节数:用
response = s3.get_object(Bucket=bucket_name, Key=key, Range=f"bytes={uploaded_so_far}-") - 不要依赖 S3 的
LastModified做幂等判断——它精度只有秒级,高频上传可能撞时间戳
怎么安全限流又不拖慢整体吞吐
全局限速(比如每秒最多 5 个文件)看似简单,但实际会让小文件“饿死”大文件——10KB 日志和 200MB 视频共用一个令牌桶,后者一占就是几十秒。更合理的是按字节速率限流 + 并发数软限制。
- 用
aiostream.stream.rate_limit()或手写 token bucket,目标速率设为50 * 1024 * 1024(50MB/s),比 MongoDB 单节点写入瓶颈略低即可 - 并发线程数控制在
min(32, CPU核心数 * 4),再多只会增加上下文切换开销 - 每个线程内上传前检查
client.admin.command("serverStatus")["metrics"]["commands"]["insert"]["failed"],连续 3 次 > 0 则自动降并发 25% - 禁用
retryWrites=true全局配置——GridFS 分块写入本身不可重试,业务层自己控更稳
最容易被忽略的清理盲区
迁移中断后,你看到的是“文件没传完”,但真正危险的是残留:fs.files 里有文档,fs.chunks 里有一半 chunk,而 S3 那边早已返回成功响应。这些孤儿 chunk 不会自动清理,也不参与任何查询,但持续占用磁盘空间,几周后可能吃光 RAID。
- 上传前必须生成 UUID 写入
metadata.upload_id,且该字段建索引:db.fs.files.createIndex({"metadata.upload_id": 1}) - 失败后第一件事是
db.fs.chunks.deleteMany({"files_id": {"$in": list_of_partial_ids}}),再删 fs.files - 不要等“最后统一清理”——迁移脚本退出前必须调用
client.close(),否则未 flush 的 chunk 缓冲还在内存里 - 每天凌晨跑一次
db.runCommand({cleanUpOrphaned: "fs.chunks"}),MongoDB 4.2+ 原生支持,别手写脚本去关联查询











