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

Mongoid聚合:解决投影输出中重复消息去重问题

问题:消除Mongoid聚合中未读最后消息的重复项

我有带标签的会话数据,会话包含带created_at和read_at(可为nil)的消息。需求是获取指定标签下所有会话的未读消息,以及最后一条消息(无论是否已读)。但当前实现的聚合代码中,若最后一条消息未读,会在结果里出现重复项,需要解决这个问题。


现有实现代码

match = { '$match': {
  'tag_ids': {
    '$in': [
      BSON::ObjectId('658db0d00f5b4f4471bdb3b6'), # inbox
    ],
  }
} }

unread_messages = { '$lookup': {
  'from':         'wco_email_message',
  'localField':   '_id',
  'foreignField': 'conversation_id',
  'pipeline': [
    { '$match': { 'read_at': nil } },
    { '$sort': { 'created_at': 1 } },
    { '$project': { 'created_at': 1 } },
  ],
  'as':           'unread_messages',
} }

last_message = { '$lookup': {
  'from':         'wco_email_message',
  'localField':   '_id',
  'foreignField': 'conversation_id',
  'pipeline': [
    { '$sort': { 'created_at': -1 } },
    { '$limit': 1 },
    { '$project': { 'created_at': 1 } },
  ],
  'as':           'last_message',
} }
project = { '$project': {
  "subject": 1,
  "all_messages": { '$concatArrays': [ '$unread_messages', '$last_message' ] },
} }

group = { '$group': {
  '_id': {
    "subject": 1,
    "all_messages": { '$addToSet': '$all_messages' },
  },
} }


outs = WcoEmail::Conversation.collection.aggregate([
  match,
  unread_messages,
  last_message,
  project,
  # group,
]).to_a

当前输出结果(存在重复项)

[
  {
    "_id"=>BSON::ObjectId('65985e070f5b4f7476a2b1c8'),
    "subject"=>"Re: JavaScript Developer _ Austin, TX.",
    "all_messages"=>[
      {"_id"=>BSON::ObjectId('65985ba10f5b4f52505dbe1a'), "created_at"=>2024-01-05 19:42:25.238 UTC},
      {"_id"=>BSON::ObjectId('65985e070f5b4f7476a2b1c9'), "created_at"=>2024-01-05 19:52:39.506 UTC},
      {"_id"=>BSON::ObjectId('65985eda0f5b4f760e1f4265'), "created_at"=>2024-01-05 19:56:10.674 UTC},
      {"_id"=>BSON::ObjectId('65985eda0f5b4f760e1f4265'), "created_at"=>2024-01-05 19:56:10.674 UTC}
    ]
  },
  {
    "_id"=>BSON::ObjectId('659f36550f5b4f114ea648ac'),
    "subject"=>"You've been chosen!",
    "all_messages"=>[
      {"_id"=>BSON::ObjectId('659f3bdc0f5b4f1a31928d7c'), "created_at"=>2024-01-11 00:52:44.109 UTC},
      {"_id"=>BSON::ObjectId('659f3bdc0f5b4f1a31928d7c'), "created_at"=>2024-01-11 00:52:44.109 UTC}
    ]
  },
  {
    "_id"=>BSON::ObjectId('659f3c970f5b4f1b3ee43a78'),
    "subject"=>"Re: Wasya Co Inquiry",
    "all_messages"=>[{"_id"=>BSON::ObjectId('659f3c970f5b4f1b3ee43a79'), "created_at"=>2024-01-11 00:55:51.927 UTC}]
  },
  {
    "_id"=>BSON::ObjectId('65bb350b0f5b4f1e126bb4ee'),
    "subject"=>"a Deploy Script",
    "all_messages"=>[
      {"_id"=>BSON::ObjectId('65bb37ea0f5b4f347256f121'), "created_at"=>2024-02-01 06:19:22.991 UTC},
      {"_id"=>BSON::ObjectId('65bb37ea0f5b4f347256f121'), "created_at"=>2024-02-01 06:19:22.991 UTC}
    ]
  },
  {
    "_id"=>BSON::ObjectId('65d524d8767ccd3c1cf87d2b'),
    "subject"=>"Resend of AWS Message",
    "all_messages"=>[
      {"_id"=>BSON::ObjectId('65d524d8767ccd3c1cf87d2c'), "created_at"=>2024-02-20 22:16:56.068 UTC},
      {"_id"=>BSON::ObjectId('65d5258a767ccd3c1cf87d2d'), "created_at"=>2024-02-20 22:19:54.454 UTC},
      {"_id"=>BSON::ObjectId('65d5258a767ccd3c1cf87d2d'), "created_at"=>2024-02-20 22:19:54.454 UTC}
    ]
  }
]

解决方案

方案1:合并后去重(简单直接)

在合并数组后,用$addToSet自动去重(依赖消息_id的唯一性),再按created_at排序恢复顺序:

outs = WcoEmail::Conversation.collection.aggregate([
  match,
  unread_messages,
  last_message,
  # 先合并数组
  { '$project': {
      "subject": 1,
      "merged_messages": { '$concatArrays': [ '$unread_messages', '$last_message' ] }
    }
  },
  # 去重并排序
  { '$addFields': {
      "all_messages": {
        '$sortArray': {
          'input': { '$reduce': {
              'input': '$merged_messages',
              'initialValue': [],
              'in': { '$cond': [
                  { '$in': [ '$$this._id', '$$value._id' ] },
                  '$$value',
                  { '$concatArrays': [ '$$value', [ '$$this' ] ] }
                ]
              }
            }
          },
          'sortBy': { 'created_at': 1 }
        }
      }
    }
  },
  # 清理临时字段
  { '$project': { "subject": 1, "all_messages": 1 } }
]).to_a

方案2:优化Lookup逻辑(从根源避免重复)

用一次$lookup同时获取未读消息和最后一条消息,再处理合并,减少查询次数:

match = { '$match': {
  'tag_ids': {
    '$in': [ BSON::ObjectId('658db0d00f5b4f4471bdb3b6') ]
  }
} }

messages_lookup = { '$lookup': {
  'from': 'wco_email_message',
  'localField': '_id',
  'foreignField': 'conversation_id',
  'pipeline': [
    { '$sort': { 'created_at': -1 } },
    { '$group': {
        '_id': '$conversation_id',
        'unread_messages': { '$push': { '$cond': [ { '$eq': [ '$read_at', nil ] }, { '_id': '$_id', 'created_at': '$created_at' }, '$$REMOVE' ] } },
        'last_message': { '$first': { '_id': '$_id', 'created_at': '$created_at' } }
      }
    },
    { '$project': {
        'unread_messages': { '$filter': { 'input': '$unread_messages', 'cond': { '$ne': [ '$$this', '$$REMOVE' ] } } },
        'last_message': 1
      }
    }
  ],
  'as': 'message_data'
} }

project = { '$project': {
  'subject': 1,
  'all_messages': {
    '$sortArray': {
      'input': { '$reduce': {
          'input': { '$concatArrays': [ { '$arrayElemAt': [ '$message_data.unread_messages', 0 ] }, [ { '$arrayElemAt': [ '$message_data.last_message', 0 ] } ] ] },
          'initialValue': [],
          'in': { '$cond': [
              { '$in': [ '$$this._id', '$$value._id' ] },
              '$$value',
              { '$concatArrays': [ '$$value', [ '$$this' ] ] }
            ]
          }
        }
      },
      'sortBy': { 'created_at': 1 }
    }
  }
} }

outs = WcoEmail::Conversation.collection.aggregate([
  match,
  messages_lookup,
  project
]).to_a

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.29 08:14:55