直接用celery.task更新MongoDB易卡住,因worker无连接池管理,频繁新建MongoClient致连接数暴涨、TIME_WAIT堆积;并发更新同一集合若无写关注或事务控制,可能丢数据或重复写入。

为什么直接用 celery.task 更新 MongoDB 容易卡住?
因为 Celery worker 默认不带 MongoDB 连接池管理,每次任务里新建 pymongo.MongoClient 会导致连接数暴涨、TCP TIME_WAIT 堆积,甚至触发系统级连接限制。更隐蔽的问题是:如果多个任务并发更新同一集合,没加写关注(w)或事务控制,可能丢数据或触发重复写入。
- 别在任务函数里反复调用
MongoClient(),应复用连接实例 - 批量更新必须显式设置
acknowledged=True,否则update_many()可能静默失败 - 避免用
find().batch_size(N)配合循环更新——这会把游标留在服务端,worker 退出后游标自动销毁,导致漏更新
如何安全地在 Celery 中复用 MongoDB 连接?
Celery 的 on_worker_process_init 钩子是初始化单例连接的正确位置,而不是在任务里 lazy 初始化。MongoDB 连接对象本身是线程安全的,但不能跨进程共享,所以每个 worker 子进程都要有自己的 client 实例。
from celery import Celery
from pymongo import MongoClient
<p>app = Celery('tasks', broker='redis://localhost')</p><h1>全局占位,实际由钩子填充</h1><p>mongo_client = None</p><p>@app.on_worker_process_init.connect
def init_mongo_client(**kwargs):
global mongo_client
mongo_client = MongoClient('mongodb://localhost:27017/', maxPoolSize=100)</p><p>@app.task
def batch_update_users(user_ids, update_data):
db = mongo_client['myapp']
result = db.users.update_many(
{'_id': {'$in': user_ids}},
{'$set': update_data},</p><h1>关键:确保写操作被确认</h1><pre class="brush:php;toolbar:false;"> upsert=False
)
return result.modified_count
update_many() 和 bulk_write() 该怎么选?
当你要对不同文档施加不同更新逻辑(比如有的要 $inc,有的要 $set,有的还要 $unset),必须用 bulk_write();如果只是统一字段覆盖,update_many() 更简洁、网络开销更低。
-
update_many():适合同质化更新,一次发一个命令,支持collation和hint,但无法混合操作类型 -
bulk_write():支持UpdateOne/ReplaceOne/DeleteOne混合,可设ordered=False让错误项跳过,但要注意:未指定upsert=True时匹配不到文档不会报错 - 两者都默认使用
w=1,生产环境建议显式传write_concern={'w': 'majority'}
怎么防止批量更新任务重试导致数据重复?
Celery 默认重试机制和 MongoDB 的“非幂等更新”一结合,就容易出问题。比如任务执行到一半 worker 挂了,重试时又跑一遍 update_many() —— 如果条件只靠 user_ids,那第二次什么都不会改,看似安全;但如果更新逻辑含 $inc 或时间戳,就会翻车。
- 给更新条件加版本号或时间窗,例如
{'_id': {'$in': ids}, 'updated_at': {'$lt': datetime.utcnow() - timedelta(hours=1)}} - 用
find_one_and_update()做单文档幂等控制,再配合 Redis 记录已处理 ID(适合中小批量) - 最稳妥的是在 MongoDB 层加唯一索引约束,让重复更新直接抛
DuplicateKeyError,然后在任务里捕获并忽略
真正麻烦的不是怎么写,而是怎么验证——批量任务跑完后,务必比对 matched_count 和 modified_count,前者是命中数,后者才是真改了的数。差太多说明有文档被条件过滤掉了,或者写关注没生效。
Python免费学习笔记(深入):立即使用
在学习笔记中,你将探索 Python 的核心概念和高级技巧!











