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

如何在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" }
        }}
      ]
    }}
  ])

关键逻辑说明

  1. 匹配用户消息:在内部管道中先用$match筛选出当前用户的所有消息;
  2. 分组统计:通过$group按campaignId聚合,用$cond判断状态并累加计数;
  3. 格式调整:用$project重命名字段,同时将数字统计结果转为字符串,完全匹配预期输出格式。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 07:57:52