0Pricing
MongoDB Academy · レッスン

集約パイプラインによるイベントのフィルタリング

watch()にパイプラインを渡し、アプリケーションが必要とするイベントだけを受信します。

「集約パイプラインによるイベントのフィルタリング」はCoddyKit上の無料MongoDB Academyレッスンです。 これはレッスン3/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはMongoDB Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 MongoDB Academyコースには全4レッスンが含まれています。

このレッスンの一部はまだ翻訳されておらず、英語で表示されています。

Why Filter Change Stream Events?

Without filtering, a change stream delivers every change event on a collection. In a busy production collection, this can mean thousands of events per second, most of which your application does not care about. Filtering at the server using an aggregation pipeline reduces network traffic, lowers CPU usage in your application, and ensures your event handler only processes relevant events. Filtering happens before events leave MongoDB—only matching events are transmitted to your client.

Passing a Pipeline to watch()

The first argument to watch() is an aggregation pipeline array. MongoDB applies this pipeline to each change event document before deciding whether to deliver it to your application. Not all aggregation stages are permitted in change stream pipelines—only a specific subset is allowed, primarily $match, $project, $addFields, $replaceRoot, and $redact. The $group and $lookup stages are not allowed.

// 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);
}

Filtering by Operation Type

Filtering by operationType is the most common pipeline filter. You can use $match with a single operation type string or with $in to match multiple types. This is useful when your application cares about inserts and updates but not deletes, or when different microservices subscribe to different operation types on the same 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' } }
]);

Filtering Update Events by Changed Fields

You can filter update events based on which fields were modified by querying the updateDescription.updatedFields object in the $match stage. This lets you subscribe only to specific field changes—for example, only when a document's status field transitions to a particular value. This is more efficient than receiving all updates and filtering in application code.

// 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);
}

Filtering by Document Field Values

For insert events, you can filter based on fields in the fullDocument sub-document. For example, receive only inserts where fullDocument.priority is 'high' or fullDocument.region equals 'US-WEST'. This server-side filtering is especially powerful in multi-tenant architectures where different application instances need events for different subsets of 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);
}

Using $project to Reshape Events

The $project stage in a change stream pipeline reshapes the event document before it is delivered to your application. You can include only the fields your handler needs, rename fields, or compute derived fields. This reduces the payload size transmitted over the network and simplifies your event handler code by presenting only the data it needs.

// 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

Using $addFields to Enrich Events

The $addFields stage lets you add computed fields to the change event document. You can add a timestamp when the event was processed, derive a category from the operation type, or compute a routing key. These enriched fields are included in the event your application receives, allowing downstream code to use pre-computed values without recalculating them.

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'
        }
      }
    }
  }
]);

Chaining Multiple Stages

You can chain multiple pipeline stages in a change stream pipeline for powerful composition. A common pattern is: $match to filter events → $addFields to enrich → $project to trim. Each stage processes the output of the previous one. Remember that the order of stages matters—apply the most selective $match first to minimize the documents processed by later stages.

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 } }
]);

Performance Impact of Server-Side Filtering

Server-side pipeline filtering in change streams is significantly more efficient than receiving all events and filtering in application code. Without server-side filtering, every event must be serialized and transmitted over the network. With a $match stage, MongoDB evaluates the filter internally and transmits only matching events. For high-traffic collections, this can reduce network usage and application CPU by orders of magnitude.

// 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
}

Permitted vs Forbidden Stages

MongoDB restricts which aggregation stages can be used in change stream pipelines. Permitted: $match, $project, $addFields, $replaceRoot, $replaceWith, $redact. Forbidden: $group, $lookup, $unwind, $geoNear, $out, $merge, and several others. Attempting to use a forbidden stage causes an error when opening the stream. If you need complex transformations, do them in application code after receiving the (pre-filtered) events.

// 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 } }
]);

Multi-Tenant Filtering Pattern

In multi-tenant applications, multiple tenants share one collection with a tenantId field. Rather than running one change stream per tenant (expensive), run one stream per service instance with a $match filter on fullDocument.tenantId scoped to the tenants that instance serves. This scales to hundreds of tenants with far fewer open cursors on the MongoDB server.

// 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' });

Quick Check

Test your understanding of MongoDB & NoSQL Databases concepts from this lesson.

Lesson Recap

In this lesson you learned: pass an aggregation pipeline as the first argument to watch() to filter events server-side, permitted stages include $match, $project, $addFields, $replaceRoot, and $redact — but not $group or $lookup, and server-side filtering dramatically reduces network traffic and application CPU compared to application-side filtering. Next up we explore resuming change streams after an interruption using resume tokens.

よくある質問

「集約パイプラインによるイベントのフィルタリング」レッスンは無料ですか?

はい。「集約パイプラインによるイベントのフィルタリング」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、MongoDB Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 MongoDB Academyコースには全4レッスンが含まれています。

「集約パイプラインによるイベントのフィルタリング」で何を学びますか?

watch()にパイプラインを渡し、アプリケーションが必要とするイベントだけを受信します。 ブラウザで直接実行するハンズオンコードでMongoDB Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。

MongoDB Academyを始めるのに経験は必要ですか?

事前経験は必要ありません。CoddyKitのMongoDB Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン3/4です。

「集約パイプラインによるイベントのフィルタリング」レッスンにはどのくらい時間がかかりますか?

ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。

このMongoDB Academyレッスンでコードを書いて実行できますか?

はい。すべてのMongoDB Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. コレクションでの変更ストリームの開始
  2. 変更イベントドキュメントの構造
  3. 集約パイプラインによるイベントのフィルタリング
  4. 中断後の変更ストリームの再開
← MongoDB Academyに戻る