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
相关产品推荐
相关产品推荐

