MongoDB Academy · 课时

使用聚合管道筛选事件

您将向 watch() 传入管道,只接收应用所关注的那部分事件。

第 3 / 4 课13 个步骤

使用聚合管道筛选事件 是 CoddyKit 上的免费 MongoDB Academy 课时。 这是第 3 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 MongoDB Academy 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 MongoDB Academy 课程共包含 4 节课。

为什么要筛选变更流事件?

如果不进行筛选,变更流会传递集合中的每一个变更事件。在繁忙的生产集合中,这可能意味着每秒数千个事件,其中大多数与您的应用无关。使用聚合管道在服务器端进行筛选,可以减少网络流量,降低应用的 CPU 使用率,并确保事件处理程序只处理相关事件。筛选发生在事件离开 MongoDB 之前——只有匹配的事件会传输到客户端。

向 watch() 传递管道

watch() 的第一个参数是聚合管道数组。MongoDB 会先将此管道应用于每个变更事件文档,然后决定是否将其传递给应用。并非所有聚合阶段都允许用于变更流管道——只有特定的子集可以使用,主要包括 $match、$project、$addFields、$replaceRoot 和 $redact。不允许使用 $group 和 $lookup 阶段。

// Only receive insert events — filter everything else
const changeStream = db.collection('orders').watch([
  {
    $match: {
      operationType: 'insert'
    }
  }
]);

for await (const change of changeStream) {
  // Only insert events arrive here
  console.log('New order:', change.fullDocument._id);
}

按操作类型筛选

按 operationType 筛选是最常见的管道筛选方式。您可以将 $match 与单个操作类型字符串配合使用,也可以使用 $in 匹配多个类型。当应用关心插入和更新但不关心删除,或者不同微服务需要订阅同一集合中的不同操作类型时,这种方式非常有用。

// React only to new orders and status updates
const stream = db.collection('orders').watch([
  {
    $match: {
      operationType: { $in: ['insert', 'update'] }
    }
  }
]);

// Or match a single type:
const deletedStream = db.collection('orders').watch([
  { $match: { operationType: 'delete' } }
]);

按已更改字段筛选更新事件

您可以根据哪些字段被修改来筛选更新事件,方法是在 $match 阶段查询 updateDescription.updatedFields 对象。这样,您就可以只订阅特定字段的变化,例如仅在文档的 status 字段变为特定值时接收事件。这比接收所有更新后再在应用代码中筛选更加高效。

// Only receive updates where status changed to 'shipped'
const shippedStream = db.collection('orders').watch([
  {
    $match: {
      operationType: 'update',
      'updateDescription.updatedFields.status': 'shipped'
    }
  }
], { fullDocument: 'updateLookup' });

for await (const change of shippedStream) {
  const order = change.fullDocument;
  await sendShippingEmail(order.customerId, order.trackingNumber);
}

按文档字段值筛选

对于插入事件,您可以根据 fullDocument 子文档中的字段进行筛选。例如,只接收 fullDocument.priority 为 'high',或 fullDocument.region 等于 'US-WEST' 的插入事件。这种服务器端筛选在多租户架构中特别有用,因为不同的应用实例可能只需要接收不同数据子集的事件。

// Only receive inserts for high-priority orders in the US-WEST region
const priorityStream = db.collection('orders').watch([
  {
    $match: {
      operationType: 'insert',
      'fullDocument.priority': 'high',
      'fullDocument.region': 'US-WEST'
    }
  }
]);

for await (const change of priorityStream) {
  await escalateOrder(change.fullDocument);
}

使用 $project 重塑事件

变更流管道中的 $project 阶段会在事件传递给应用之前重塑事件文档。您可以只保留处理程序所需的字段、重命名字段,或计算派生字段。这样可以减小通过网络传输的数据量,并通过只呈现所需数据来简化事件处理程序代码。

// Project only the fields the handler needs
const stream = db.collection('users').watch([
  { $match: { operationType: { $in: ['insert', 'update'] } } },
  {
    $project: {
      operationType: 1,
      'documentKey._id': 1,
      'updateDescription.updatedFields.email': 1,
      'fullDocument.email': 1,
      'fullDocument.name': 1
    }
  }
]);

// Handler receives trimmed events with only email and name

使用 $addFields 丰富事件

$addFields 阶段允许您向变更事件文档中添加计算字段。您可以添加事件处理时间戳、根据操作类型推导类别,或计算路由键。这些丰富后的字段会包含在应用接收到的事件中,使后续代码能够直接使用预先计算的值,而无需重新计算。

const stream = db.collection('payments').watch([
  {
    $addFields: {
      processedAt: '$$NOW', // current timestamp as event enrichment
      eventCategory: {
        $switch: {
          branches: [
            { case: { $eq: ['$operationType', 'insert'] }, then: 'NEW_PAYMENT' },
            { case: { $eq: ['$operationType', 'update'] }, then: 'PAYMENT_UPDATE' }
          ],
          default: 'OTHER'
        }
      }
    }
  }
]);

串联多个阶段

您可以在变更流管道中串联多个管道阶段,实现强大的组合处理。常见模式是:使用 $match 筛选事件 → 使用 $addFields 丰富事件 → 使用 $project 精简事件。每个阶段都会处理前一个阶段的输出。请记住,阶段顺序很重要——应先应用筛选性最强的 $match,以减少后续阶段需要处理的文档数量。

const stream = db.collection('inventory').watch([
  // Stage 1: filter to updates only
  { $match: { operationType: 'update' } },
  // Stage 2: add computed field
  {
    $addFields: {
      isLowStock: {
        $lt: ['$updateDescription.updatedFields.quantity', 10]
      }
    }
  },
  // Stage 3: only pass through low-stock events
  { $match: { isLowStock: true } },
  // Stage 4: trim to essential fields
  { $project: { 'documentKey._id': 1, operationType: 1 } }
]);

服务器端筛选的性能影响

在变更流中进行服务器端管道筛选,比接收所有事件后再在应用代码中筛选高效得多。如果不进行服务器端筛选,每个事件都必须序列化并通过网络传输。使用 $match 阶段时,MongoDB 会在内部评估筛选条件,只传输匹配的事件。对于高流量集合,这可以将网络使用量和应用 CPU 使用量降低几个数量级。

// Inefficient: receive all events, filter in JS
for await (const change of db.collection('orders').watch()) {
  if (change.operationType === 'insert' && change.fullDocument.total > 1000) {
    // Most events are discarded here — wasted network I/O
  }
}

// Efficient: filter server-side
for await (const change of db.collection('orders').watch([
  { $match: { operationType: 'insert', 'fullDocument.total': { $gt: 1000 } } }
])) {
  // Only matching events arrive here
}

允许使用的阶段与禁止使用的阶段

MongoDB 会限制可用于变更流管道的聚合阶段。允许使用:$match、$project、$addFields、$replaceRoot、$replaceWith、$redact。禁止使用:$group、$lookup、$unwind、$geoNear、$out、$merge 等。尝试使用禁止的阶段会在打开流时导致错误。如果需要复杂转换,请在接收经过预筛选的事件后,于应用代码中完成。

// WRONG — $group is not allowed in change stream pipelines
db.collection('orders').watch([
  { $group: { _id: '$fullDocument.region', count: { $sum: 1 } } } // Error!
]);

// RIGHT — use only permitted stages in the pipeline
db.collection('orders').watch([
  { $match: { operationType: 'insert' } },
  { $project: { 'fullDocument.region': 1, 'fullDocument.total': 1 } }
]);

多租户筛选模式

在多租户应用中,多个租户通过 tenantId 字段共享一个集合。与其为每个租户运行一个变更流(成本较高),不如让每个服务实例运行一个流,并针对该实例负责的租户,对 fullDocument.tenantId 使用 $match 筛选条件。这样,只需在 MongoDB 服务器上打开少得多的游标,就可以扩展到数百个租户。

// Service instance handles tenants T1 and T2 only
const myTenants = ['T1', 'T2'];

const stream = db.collection('events').watch([
  {
    $match: {
      $or: [
        { 'fullDocument.tenantId': { $in: myTenants } },  // for inserts
        { 'updateDescription.updatedFields.tenantId': { $in: myTenants } }  // for updates
      ]
    }
  }
], { fullDocument: 'updateLookup' });

快速检查

请测试您对本课中 MongoDB 和 NoSQL 数据库概念的理解。

课程回顾

在本课中,您学到了:将聚合管道作为 watch() 的第一个参数,以便在服务器端筛选事件;允许使用的阶段包括 $match、$project、$addFields、$replaceRoot 和 $redact,但不包括 $group 或 $lookup;以及与应用端筛选相比,服务器端筛选可以大幅减少网络流量和应用 CPU 使用量。接下来,我们将学习如何使用恢复令牌,在中断后恢复变更流。

免费开始

用 AI 导师学习 JavaScript — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
30
课程
120

常见问题解答

「使用聚合管道筛选事件」课时是免费的吗?

是的 — 「使用聚合管道筛选事件」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 MongoDB Academy 课程的其余内容,请升级到 CoddyKit PRO。 MongoDB Academy 课程共包含 4 节课。

「使用聚合管道筛选事件」这节课中我会学到什么?

您将向 watch() 传入管道,只接收应用所关注的那部分事件。 你通过在浏览器中直接运行的动手代码来练习 MongoDB Academy,全天候 AI 导师会在你学习这节课的过程中回答你的问题。

学习 MongoDB Academy 需要有经验吗?

无需任何先前经验。CoddyKit 上的 MongoDB Academy 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 3 节课,共 4 节。

「使用聚合管道筛选事件」课时需要多长时间?

大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。

我能在这节 MongoDB Academy 课中编写并运行代码吗?

能。每节 MongoDB Academy 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。

此课程中的所有课时

  1. 在集合上打开变更流
  2. 变更事件文档结构
  3. 使用聚合管道筛选事件
  4. 中断后恢复变更流
← 返回 MongoDB Academy