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

如何通过原子操作实现MongoDB通知的合并更新逻辑?

用MongoDB原子操作实现连续通知合并逻辑

你的现有代码存在并发风险:如果两个相同通知同时请求,或者在查询最后一条通知和更新之间插入了其他通知,会导致错误更新旧的通知文档,不符合"仅合并连续相同通知"的需求。以下是两种原子化的解决方案:

方案一:单命令原子操作(无需修改Schema)

利用MongoDB的findOneAndUpdate配合聚合表达式,直接匹配用户最新的且message相同的通知,匹配到则更新count,否则插入新文档:

async function createNotification(user_id, message) {
  const notifications = db.collection('notifications');
  await notifications.findOneAndUpdate(
    // 查询条件:匹配用户、当前message,且是该用户最新的通知
    {
      user_id,
      message,
      $expr: {
        $eq: [
          "$timestamp",
          {
            $max: {
              $lookup: {
                from: "notifications",
                localField: "user_id",
                foreignField: "user_id",
                as: "user_timestamps",
                pipeline: [{ $project: { timestamp: 1, _id: 0 } }]
              }
            }
          }
        ]
      }
    },
    // 更新操作:count加1,timestamp更新为当前时间
    {
      $inc: { count: 1 },
      $currentDate: { timestamp: true }
    },
    // 配置:无匹配则插入新文档,返回更新后的文档
    {
      upsert: true,
      returnDocument: 'after'
    }
  );
}

说明

  • 查询条件里的$lookup会获取该用户所有通知的时间戳,$max取最大值后和当前文档的timestamp对比,确保只匹配用户最新的通知
  • $inc会自动处理新插入的文档:count字段不存在时,直接设为1
  • $currentDate会将timestamp设为当前服务器时间,不管是更新还是插入

方案二:事务实现(高效推荐)

如果你的MongoDB版本支持事务(4.0+),推荐用事务包裹查询和更新/插入操作,既保证原子性,又避免了$lookup带来的性能损耗:

async function createNotification(user_id, message) {
  const session = await db.startSession();
  session.startTransaction();
  try {
    const notifications = db.collection('notifications');
    // 在事务内查询最新通知
    const lastNotification = await notifications.findOne(
      { user_id },
      { sort: { timestamp: -1 }, session }
    );

    if (lastNotification?.message === message) {
      // 更新最新通知的count和timestamp
      await notifications.updateOne(
        { _id: lastNotification._id },
        { $inc: { count: 1 }, $currentDate: { timestamp: true } },
        { session }
      );
    } else {
      // 插入新通知
      await notifications.insertOne(
        { user_id, message, count: 1, timestamp: new Date() },
        { session }
      );
    }
    await session.commitTransaction();
  } catch (error) {
    await session.abortTransaction();
    throw error;
  } finally {
    session.endSession();
  }
}

说明

  • 事务内的所有操作要么全部成功,要么全部回滚,彻底避免了并发场景下的间隙问题
  • 相比方案一,无需执行$lookup,性能更优,尤其是用户通知数量较多时

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 00:42:06