Filtrering af hændelser med en aggregeringspipeline
De lærende vil sende en pipeline til watch() for kun at modtage det udvalg af hændelser, som applikationen har brug for.
Filtrering af hændelser med en aggregeringspipeline er en gratis MongoDB Academy-lektion på CoddyKit. Dette er lektion 3 af 4. Du kan læse hele lektionen gratis nedenfor — og derefter øve dig praktisk i browseren med en indbygget kodeeditor og en AI-vejleder, der er tilgængelig døgnet rundt. Den er en del af læringsforløbet i MongoDB Academy, og dine fremskridt synkroniseres på tværs af nettet og CoddyKit-appen. MongoDB Academy-kurset indeholder 4 lektioner i alt.
Hvorfor filtrere ændringsstreamhændelser?
Uden filtrering leverer en ændringsstream alle ændringshændelser i en collection. I en travl produktionscollection kan det betyde tusindvis af hændelser i sekundet, hvoraf de fleste ikke er relevante for din applikation. Filtrering på serveren ved hjælp af en aggregeringspipeline reducerer netværkstrafikken, sænker CPU-forbruget i din applikation og sikrer, at din hændelseshåndtering kun behandler relevante hændelser. Filtreringen sker, før hændelserne forlader MongoDB – kun hændelser, der matcher, sendes til din klient.
Videregivelse af en pipeline til watch()
Det første argument til watch() er en array med en aggregeringspipeline. MongoDB anvender denne pipeline på hvert dokument med en ændringshændelse, før det afgøres, om det skal leveres til din applikation. Ikke alle aggregeringsstadier er tilladt i ændringsstreampipelines – kun en bestemt delmængde, primært $match, $project, $addFields, $replaceRoot og $redact. Stadierne $group og $lookup er ikke tilladt.
// 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);
}Filtrering efter handlingstype
Filtrering efter operationType er det mest almindelige pipelinefilter. Du kan bruge $match med én streng for handlingstypen eller med $in for at matche flere typer. Det er nyttigt, når din applikation er interesseret i indsættelser og opdateringer, men ikke sletninger, eller når forskellige mikrotjenester abonnerer på forskellige handlingstyper i den samme collection.
// 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' } }
]);Filtrering af opdateringshændelser efter ændrede felter
Du kan filtrere opdateringshændelser baseret på hvilke felter der blev ændret ved at forespørge på objektet updateDescription.updatedFields i stadiet $match. Det giver dig mulighed for kun at abonnere på bestemte feltændringer – for eksempel kun når et dokuments status-felt ændres til en bestemt værdi. Det er mere effektivt end at modtage alle opdateringer og filtrere dem i applikationskoden.
// 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);
}Filtrering efter feltværdier i dokumenter
Ved indsættelseshændelser kan du filtrere ud fra felter i underdokumentet fullDocument. Du kan for eksempel kun modtage indsættelser, hvor fullDocument.priority er 'high', eller hvor fullDocument.region er lig med 'US-WEST'. Denne filtrering på serveren er især effektiv i arkitekturer med flere lejere, hvor forskellige applikationsinstanser har brug for hændelser for forskellige delmængder af data.
// 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);
}Brug af $project til at omforme hændelser
Stadiet $project i en ændringsstreampipeline omformer hændelsesdokumentet, før det leveres til din applikation. Du kan medtage kun de felter, din håndtering har brug for, omdøbe felter eller beregne afledte felter. Det reducerer størrelsen på den nyttelast, der sendes over netværket, og forenkler koden til din hændelseshåndtering ved kun at præsentere de nødvendige data.
// 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 nameBrug af $addFields til at berige hændelser
Stadiet $addFields giver dig mulighed for at tilføje beregnede felter til dokumentet med ændringshændelsen. Du kan tilføje et tidsstempel for, hvornår hændelsen blev behandlet, udlede en kategori ud fra handlingstypen eller beregne en routingnøgle. Disse berigede felter medtages i den hændelse, din applikation modtager, så efterfølgende kode kan bruge forudberegnede værdier uden at beregne dem igen.
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'
}
}
}
}
]);Sammenkædning af flere stadier
Du kan sammenkæde flere pipelinestadier i en ændringsstreampipeline for at opnå effektiv sammensætning. Et almindeligt mønster er: $match for at filtrere hændelser → $addFields for at berige → $project for at beskære. Hvert stadie behandler outputtet fra det foregående. Husk, at stadiernes rækkefølge er vigtig – anvend det mest selektive $match først for at minimere de dokumenter, som senere stadier behandler.
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 } }
]);Ydelsespåvirkning ved filtrering på serveren
Pipelinefiltrering på serveren i ændringsstreams er væsentligt mere effektiv end at modtage alle hændelser og filtrere dem i applikationskoden. Uden filtrering på serveren skal hver hændelse serialiseres og sendes over netværket. Med et $match-stadie evaluerer MongoDB filteret internt og sender kun hændelser, der matcher. For collections med høj trafik kan det reducere netværksforbruget og applikationens CPU-forbrug med flere størrelsesordener.
// 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
}Tilladte og forbudte stadier
MongoDB begrænser, hvilke aggregeringsstadier der kan bruges i ændringsstreampipelines. Tilladt: $match, $project, $addFields, $replaceRoot, $replaceWith, $redact. Forbudt: $group, $lookup, $unwind, $geoNear, $out, $merge og flere andre. Hvis du forsøger at bruge et forbudt stadie, opstår der en fejl, når streamen åbnes. Hvis du har brug for komplekse transformationer, skal du udføre dem i applikationskoden efter modtagelsen af de (forfiltrerede) hændelser.
// 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 } }
]);Filtreringsmønster til flere lejere
I applikationer med flere lejere deler flere lejere én collection med et tenantId-felt. I stedet for at køre én ændringsstream pr. lejer (hvilket er dyrt) kan du køre én stream pr. serviceinstans med et $match-filter på fullDocument.tenantId, begrænset til de lejere, som instansen betjener. Det kan skaleres til hundredvis af lejere med langt færre åbne markører på MongoDB-serveren.
// 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' });Hurtigt tjek
Test din forståelse af begreberne MongoDB & NoSQL Databases fra denne lektion.
Opsummering af lektionen
I denne lektion har du lært at videregive en aggregeringspipeline som det første argument til watch() for at filtrere hændelser på serveren, at tilladte stadier omfatter $match, $project, $addFields, $replaceRoot og $redact – men ikke $group eller $lookup, og at filtrering på serveren drastisk reducerer netværkstrafikken og applikationens CPU-forbrug sammenlignet med filtrering i applikationen. Dernæst undersøger vi, hvordan ændringsstreams genoptages efter en afbrydelse ved hjælp af genoptagelsestokens.
Lær JavaScript med en AI-underviser — gratis
Skriv og kør rigtig kode i din browser, få øjeblikkelig hjælp fra en AI-underviser døgnet rundt, og fortsæt, hvor du slap, på web eller i appen.
- Kurser
- 30
- Lektioner
- 120
Ofte stillede spørgsmål
Er lektionen “Filtrering af hændelser med en aggregeringspipeline” gratis?
Ja — hele teksten til “Filtrering af hændelser med en aggregeringspipeline” kan læses gratis her på nettet. Hvis du vil øve dig interaktivt med en indbygget kodeeditor og en AI-vejleder døgnet rundt og få adgang til resten af MongoDB Academy-kurset, skal du opgradere til CoddyKit PRO. MongoDB Academy-kurset indeholder 4 lektioner i alt.
Hvad lærer jeg i “Filtrering af hændelser med en aggregeringspipeline”?
De lærende vil sende en pipeline til watch() for kun at modtage det udvalg af hændelser, som applikationen har brug for. Du øver dig i MongoDB Academy med praktisk kode, som du kører direkte i browseren, og en AI-vejleder døgnet rundt besvarer dine spørgsmål, mens du arbejder dig gennem lektionen.
Skal jeg have erfaring for at begynde på MongoDB Academy?
Der kræves ingen tidligere erfaring. MongoDB Academy på CoddyKit er tilrettelagt for både begyndere og øvede, så du kan starte her eller fra begyndelsen og lære i dit eget tempo. Dette er lektion 3 af 4.
Hvor lang tid tager lektionen “Filtrering af hændelser med en aggregeringspipeline”?
De fleste CoddyKit-lektioner tager cirka 5–10 minutter. Hver lektion er kort og interaktiv, så du gør løbende fremskridt og kan fortsætte, hvor du slap – på både web og app.
Kan jeg skrive og køre kode i denne MongoDB Academy-lektion?
Ja. Alle MongoDB Academy-lektioner har en indbygget kodeeditor, så du kan skrive og køre rigtig kode direkte i din browser og få øjeblikkelig feedback fra AI – uden lokal opsætning.
Alle lektioner i dette kursus
- Åbning af en ændringsstrøm på en samling
- Strukturen af et ændringshændelsesdokument
- Filtrering af hændelser med en aggregeringspipeline
- Genoptagelse af ændringsstrømme efter afbrydelser