如何在MongoDB中通过Aggregate关联表并统计消息状态数量
如何在MongoDB聚合中添加对Messages表的Lookup以获取按Campaign分组的状态统计?
我正在尝试用MongoDB的lookup、pipeline和match关联多张表,需要提取Message表中按Campaign分组的特定状态统计数量。
Message表示例数据
[ { "createdAt":"2023-03-05T10:47:09.646Z", "_id":"6404735c55bfc04cf96ecbbe", "message":"Testing", "userId":"62f13fd054f851ecc2c12e2a", "status":"submitted", "__v":0, "campaignId":1 }, { "createdAt":"2023-03-05T10:47:09.652Z", "_id":"640473a916577b4cfa9c9e29", "message":"Testing one", "userId":"62f13fd054f851ecc2c12e2a", "status":"delivered", "__v":0, "campaignId":2 }, { "createdAt":"2023-03-05T10:47:09.695Z", "_id":"640473ae7052674cfed04602", "message":"Testing three", "userId":"62f13fd054f851ecc2c12e2a", "status":"rejected", "__v":0, "campaignId":3 }, { "createdAt":"2023-03-05T11:05:41.217Z", "_id":"6404779830ca9a5d3603c51a", "paymentMethod":"message 6", "userId":"62f13fd054f851ecc2c12e2a", "status":"submitted", "__v":0, "campaignId":4 } ]
现有查询
Users.aggregate([ { $match: { _id: user._id } }, { $lookup: { as: "userBilling", from: "sms_transactions", let: { user_id: "$_id" }, pipeline: [ { $match: { $expr: { $eq: [ "$userId", "$$user_id" ] } } }, ], }}, ])
预期结果
注:修正了原格式中不符合JSON规范的部分,调整为合法的数组对象结构
{ "_id":"62f13fd054f851ecc2c12e2a", "userBilling":[], "userMessagesCount":[ { "campaignId": "1", "userId":"62f13fd054f851ecc2c12e2a", "submitted":"1", "delivered":"0", "rejected":"0" }, { "campaignId": "2", "userId":"62f13fd054f851ecc2c12e2a", "submitted":"0", "delivered":"1", "rejected":"0" } ] }
各集合Schema
User集合
var UserSchema = new mongoose.Schema({ firstName: {type:String,required: true}, lastName: {type:String,required: true}, email:{ type:String, unique: true}, password: String, phoneNumber:{ type:String, unique: true}, role:String, status:String, confirmed_phone:Boolean, confirmed_email:Boolean, country:{type:String,required: true}, token: String, created_at:{type: Date,default: Date.now()}, },{ autoCreate: true}); mongoose.model('Users', UserSchema); module.exports = mongoose.model('Users');
Transaction集合
var SMSTransactions = new mongoose.Schema( { userId: { type: mongoose.Schema.ObjectId, ref: 'users' }, amount: { type: Number, required: true }, balanceBefore: { type: String }, balanceAfter: { type: String }, invoiceId: { type: String, default: () => MUUID.v4().toString("D"), }, createdAt: { type: Date, default: Date.now() }, updatedAt: mongoose.Date, paymentMethod: { type: String, required: true }, callbackUrl: String, status: String, country : String }, { autoCreate: true } ); mongoose.model("sms_transactions", SMSTransactions); module.exports = mongoose.model("sms_transactions");
Messages集合
var MessagesSchema = new mongoose.Schema({ message: {type:String,required: true}, recipient: String, sender:String, status:String, bal_before:mongoose.Types.Decimal128, balance_after:mongoose.Types.Decimal128, created_at:{type: Date,default: Date.now()}, updated_at: mongoose.Date, type:String, userId: { type: mongoose.Schema.ObjectId, ref: 'users' }, campaignId:String },{ autoCreate: true}); mongoose.model('messages', MessagesSchema); module.exports = mongoose.model('messages');
解决方案
在现有聚合管道中添加一个新的$lookup阶段,针对messages集合,通过内部管道实现按campaignId分组并统计各状态数量:
Users.aggregate([ { $match: { _id: user._id } }, // 原有的userBilling关联 { $lookup: { as: "userBilling", from: "sms_transactions", let: { user_id: "$_id" }, pipeline: [ { $match: { $expr: { $eq: [ "$userId", "$$user_id" ] } } }, ], }}, // 新增的消息统计关联 { $lookup: { as: "userMessagesCount", from: "messages", let: { user_id: "$_id" }, pipeline: [ // 筛选当前用户的所有消息 { $match: { $expr: { $eq: [ "$userId", "$$user_id" ] } } }, // 按campaignId分组,统计各状态数量 { $group: { _id: "$campaignId", userId: { $first: "$userId" }, submitted: { $sum: { $cond: [{ $eq: ["$status", "submitted"] }, 1, 0] } }, delivered: { $sum: { $cond: [{ $eq: ["$status", "delivered"] }, 1, 0] } }, rejected: { $sum: { $cond: [{ $eq: ["$status", "rejected"] }, 1, 0] } } }}, // 调整字段格式,匹配预期结果 { $project: { _id: 0, campaignId: "$_id", userId: 1, submitted: { $toString: "$submitted" }, delivered: { $toString: "$delivered" }, rejected: { $toString: "$rejected" } }} ] }} ])
关键逻辑说明
- 匹配用户消息:在内部管道中先用
$match筛选出当前用户的所有消息; - 分组统计:通过
$group按campaignId聚合,用$cond判断状态并累加计数; - 格式调整:用
$project重命名字段,同时将数字统计结果转为字符串,完全匹配预期输出格式。
内容的提问来源于stack exchange,提问作者lutakyn
相关产品推荐
相关产品推荐

