Django+MongoDB集群环境下原子更新遇Retryable write禁止报错
问题解决:MongoDB集群+Celery环境下的重试写入报错与原子更新实现
报错原因
这个错误是MongoDB重试写入机制触发的:同一个会话中,旧的重试写任务(txnNumber 4)还未完成,新的重试写任务(txnNumber 6)已启动,MongoDB禁止这种重叠操作,避免破坏数据一致性。
放到你的代码场景里,问题根源在这几点:
- Celery是多进程运行环境,你用的
threading.Lock仅在单个进程的线程内有效,跨进程完全无法控制并发 FileField.replace操作和using_large_results、results等字段更新是分开执行的,属于同一个会话内的多个独立写操作,直接触发了重试写冲突- 直接修改EmbeddedDocument属性后,未通过父文档的原子更新提交,导致会话复用出现txnNumber重叠问题
解决方案
1. 替换无效线程锁为分布式乐观锁
Celery多进程/分布式环境下,线程锁等于摆设,改用基于数据库的乐观锁最靠谱——给父文档新增版本号字段,每次更新前校验版本,更新成功则自增版本号,并发冲突时自动重试。
2. 实现全流程原子更新
因为MyWrapper是嵌入式文档,所有修改必须通过父文档的原子操作完成,不能零散地修改字段、写入GridFS文件。要把GridFS文件读写和文档字段更新绑定成一个原子流程。
3. 可选:禁用重试写入
如果仍存在会话冲突,可直接在MongoDB连接配置中关闭重试写入功能。
修改后的代码示例
首先给父文档添加版本号字段(假设父文档名为TaskResult):
from mongoengine import Document, EmbeddedDocument, DictField, BooleanField, FileField, IntField import json from gridfs import GridFS import logging logger = logging.getLogger(__name__) class InvalidResultError(Exception): pass class MyWrapper(EmbeddedDocument): MAXIMUM_MONGO_DOCUMENT_SIZE = 12582912 results = DictField(default={}) _large_results = FileField(required=True) using_large_results = BooleanField(default=False) def large_results(self): try: self._large_results.seek(0) return json.load(self._large_results) except Exception: return {} def __get_true_result(self): if self.using_large_results: self._large_results.seek(0) try: return json.loads(self._large_results.read() or '{}') except Exception as e: logger.exception("解析大结果JSON失败") raise InvalidResultError from e else: return self.results.copy() # 父文档必须添加version字段用于乐观锁 class TaskResult(Document): wrapper = EmbeddedDocumentField(MyWrapper) version = IntField(default=0) # 其他业务字段...
然后编写原子更新的核心逻辑:
def update_task_result(task_result_id, result, result_class, update=False): class_name = result_class.__name__ db = TaskResult._get_db() fs = GridFS(db) while True: # 1. 原子性读取当前文档与版本号 task_result = TaskResult.objects(id=task_result_id).select_related('wrapper').first() if not task_result: raise ValueError("找不到对应的任务结果") current_version = task_result.version wrapper = task_result.wrapper valid_result = wrapper.__get_true_result() # 2. 计算待更新的结果数据 try: current = valid_result[class_name] if update else {} except KeyError: current = {} if update: current.update(result) else: current = result valid_result[class_name] = current json_result = json.dumps(valid_result) using_large = len(json_result) >= MyWrapper.MAXIMUM_MONGO_DOCUMENT_SIZE # 3. 准备原子更新操作 update_ops = {} new_file_id = None if using_large: # 处理大结果:删除旧文件,写入新文件 if wrapper._large_results.grid_id: fs.delete(wrapper._large_results.grid_id) new_file = fs.put(json_result.encode('utf-8'), content_type='application/json') new_file_id = new_file # 设置嵌入式文档字段 update_ops["wrapper.results"] = {} update_ops["wrapper.using_large_results"] = True update_ops["wrapper._large_results"] = new_file_id else: # 处理小结果:更新results字段,清空大文件占位 if wrapper._large_results.grid_id: fs.delete(wrapper._large_results.grid_id) new_file = fs.put(b'{}', content_type='application/json') new_file_id = new_file update_ops["wrapper.results"] = valid_result update_ops["wrapper.using_large_results"] = False update_ops["wrapper._large_results"] = new_file_id # 4. 原子更新父文档,通过版本号校验实现乐观锁 updated_count = TaskResult.objects( id=task_result_id, version=current_version ).update( **update_ops, inc__version=1 # 版本号自增,并发冲突时自动重试 ) if updated_count == 1: # 更新成功,退出循环 break # 更新失败说明存在并发修改,重新执行流程
如果仍存在重试写入冲突,可在MongoDB连接时禁用重试:
from mongoengine import connect connect( db='你的数据库名', host='mongodb://集群地址', retryWrites=False # 关闭重试写入功能 )
关键改动说明
- 乐观锁控并发:通过父文档的
version字段确保同一时间只有一个更新操作成功,并发冲突时自动重试 - 原子操作绑定:将GridFS文件操作与文档字段更新绑定为一个流程,避免多个独立写操作导致的会话冲突
- 弃用无效锁:替换
threading.Lock为数据库级乐观锁,适配Celery多进程/分布式环境 - 会话隔离:每次更新使用独立会话,避免出现txnNumber重叠问题
内容的提问来源于stack exchange,提问作者Shashank Pathak
相关产品推荐
相关产品推荐

