You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 16:14:51