使用聚合管道筛选事件
您将向 watch() 传入管道,只接收应用所关注的那部分事件。
使用聚合管道筛选事件 是 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 反馈 — 无需本地设置。