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

MongoDB 3.4副本集并发findAndModify更新同一文档的竞态问题问询

解决MongoDB 3.4副本集下任务工作器的竞态条件问题

你遇到的这个场景太典型了——多线程同时完成任务时,都误以为自己是最后一个,导致重复触发下一阶段。结合MongoDB 3.4的特性,我给你几个经过生产环境验证的解决方案:

方案一:维护任务集全局状态文档(首推)

这个方案利用MongoDB单文档操作的原子性,从根源上避免竞态。核心思路是给每个任务集单独维护一个状态文档,记录总任务数和已完成任务数,每次任务完成时原子递增已完成数,同时判断是否达到总任务数。

实现步骤:

  1. 任务集创建时,插入状态文档示例:
    {
      "_id": "task_set_123",
      "totalTasks": 5,
      "completedTasks": 0,
      "status": "IN_PROGRESS"
    }
    
  2. 单个任务完成时的Java代码实现:
    // 假设任务集合为taskCollection,任务集状态集合为taskSetCollection
    // 1. 标记当前任务为完成
    UpdateResult taskUpdateResult = taskCollection.updateOne(
        Filters.eq("_id", taskId),
        Updates.set("status", "COMPLETED")
    );
    
    if (taskUpdateResult.getModifiedCount() == 1) {
        // 2. 原子递增已完成数,仅当任务集未完成时执行
        FindOneAndUpdateOptions options = new FindOneAndUpdateOptions()
            .returnDocument(ReturnDocument.AFTER)
            .upsert(false);
    
        Document updatedTaskSet = taskSetCollection.findOneAndUpdate(
            Filters.and(
                Filters.eq("_id", taskSetId),
                Filters.eq("status", "IN_PROGRESS"),
                Filters.lt("completedTasks", "totalTasks")
            ),
            Updates.combine(
                Updates.inc("completedTasks", 1),
                Updates.set("status", Filters.cond(
                    Filters.eq("completedTasks", Filters.subtract("totalTasks", 1)),
                    "COMPLETED",
                    "IN_PROGRESS"
                ))
            ),
            options
        );
    
        // 3. 若更新后任务集状态为COMPLETED,说明当前是最后一个完成的任务
        if (updatedTaskSet != null && "COMPLETED".equals(updatedTaskSet.getString("status"))) {
            triggerNextProcessingPhase(taskSetId);
        }
    }
    

优点:

  • 完全依赖MongoDB原子操作,无竞态风险
  • 性能优异,无需额外锁机制
  • 状态清晰,便于监控任务集进度

方案二:分布式锁+全局检查

如果任务集的总任务数无法提前确定(比如动态添加任务),可以用分布式锁保证只有一个线程能执行「检查所有任务是否完成」的操作。

实现步骤:

  1. 用MongoDB文档作为分布式锁,每个任务集对应一个锁文档
  2. 任务完成后尝试获取锁,成功后再检查任务状态:
    // 锁集合:lockCollection
    String lockId = taskSetId + "_completion_lock";
    FindOneAndUpdateOptions lockOptions = new FindOneAndUpdateOptions()
        .upsert(true)
        .returnDocument(ReturnDocument.AFTER);
    
    // 尝试获取锁(设置10秒过期,避免死锁)
    Document lock = lockCollection.findOneAndUpdate(
        Filters.and(
            Filters.eq("_id", lockId),
            Filters.or(
                Filters.eq("locked", false),
                Filters.lt("expireAt", new Date())
            )
        ),
        Updates.combine(
            Updates.set("locked", true),
            Updates.set("expireAt", new Date(System.currentTimeMillis() + 10000))
        ),
        lockOptions
    );
    
    if (lock != null && lock.getBoolean("locked")) {
        try {
            // 检查当前任务集是否所有任务都已完成
            long incompleteTasks = taskCollection.countDocuments(
                Filters.and(
                    Filters.eq("taskSetId", taskSetId),
                    Filters.ne("status", "COMPLETED")
                )
            );
    
            if (incompleteTasks == 0) {
                triggerNextProcessingPhase(taskSetId);
                // 标记任务集为完成,避免后续重复检查
                taskSetCollection.updateOne(
                    Filters.eq("_id", taskSetId),
                    Updates.set("status", "COMPLETED")
                );
            }
        } finally {
            // 释放锁
            lockCollection.updateOne(
                Filters.eq("_id", lockId),
                Updates.set("locked", false)
            );
        }
    }
    

注意事项:

  • 必须给锁设置过期时间,防止线程崩溃导致锁永久占用
  • 锁的粒度要合适,避免影响其他任务集的执行

方案三:原子更新+乐观锁(不推荐)

如果不想维护额外状态文档,可以在任务集文档添加version字段做乐观锁,但高并发下更新失败率高,需要重试逻辑,性能不如前两个方案。

核心代码示例:

boolean triggerNextPhase = false;
int retryCount = 3;

while (retryCount-- > 0 && !triggerNextPhase) {
    Document taskSet = taskSetCollection.find(Filters.eq("_id", taskSetId)).first();
    if (taskSet == null) break;

    long completedCount = taskCollection.countDocuments(
        Filters.and(
            Filters.eq("taskSetId", taskSetId),
            Filters.eq("status", "COMPLETED")
        )
    );

    if (completedCount == taskSet.getInteger("totalTasks")) {
        UpdateResult result = taskSetCollection.updateOne(
            Filters.and(
                Filters.eq("_id", taskSetId),
                Filters.eq("version", taskSet.getInteger("version")),
                Filters.eq("status", "IN_PROGRESS")
            ),
            Updates.combine(
                Updates.set("status", "COMPLETED"),
                Updates.inc("version", 1)
            )
        );

        if (result.getModifiedCount() == 1) {
            triggerNextPhase = true;
        }
    } else {
        break;
    }
}

if (triggerNextPhase) {
    triggerNextProcessingPhase(taskSetId);
}

总结

优先选择方案一,它利用MongoDB原生原子性,高效且彻底解决竞态问题。如果任务集总任务数动态变化,再考虑方案二。方案三只适合并发度较低的场景,否则重试会影响性能。

内容的提问来源于stack exchange,提问作者Eddie

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 10:19:17