Prompt 18 · Lesson 02 · Scope and filtering

Watch a Collection, Database, or Deployment and Filter Change Events with Pipelines

Watch only the namespaces a consumer owns and filter change-event metadata on the server without damaging the resume token.

Intermediate–Advanced120–190 minutesCDC/resumability engineering labMongoDB 8.3.8 · mongosh 2.10.0 · PyMongo 4.17.0Last reviewed: September 2026

Learning objectives

01

Distinguish collection-, database-, and deployment-scoped change streams and their authorization blast radius.

02

Filter change events with supported pipeline stages without removing or mutating the resume token.

03

Use namespace and operationType evidence to prove which events a broader stream sees.

04

Explain the connection-pool and sharded-router costs of many broad streams.

05

Diagnose an over-broad deployment watch and replace it with the narrowest useful scope and filter.

Reproducible lab baseline

This lesson pins MongoDB Community Server 8.3.8 with mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0, and PyMongo 4.17.0. The mandatory lab uses a disposable single-member replica set on loopback port 27156; change streams require a replica set or sharded cluster, so a standalone is intentionally not used. A single-member replica set is enough to learn event/resume mechanics but does not demonstrate high availability, multi-node majority durability, or failover. Authentication and TLS are disabled only for this isolated lab. Feature Compatibility Version (FCV) is inspected. FCV is observed and never changed. Default read/write concern and primary read preference apply unless a command says otherwise. The lab uses collection and database scopes locally. A deployment-scoped stream is also demonstrated through PyMongo MongoClient.watch(); on a sharded cluster the corresponding stream must be opened through mongos. Atlas, Search, Vector Search, KMS, and Enterprise Advanced are not mandatory. Product commands were not executed in this generation environment because Docker, mongod, mongosh, and PyMongo are unavailable here; runtime timings and token values must be measured on the learner machine rather than copied as invented output.

1. Scope is part of the consumer contract

AtlasMart has three consumers: inventory cache invalidation needs only inventory, an audit pipeline needs selected changes across the atlasmart database, and a central governance service might need all non-system collections in an entire deployment. Opening the broadest stream “just in case” increases authorization scope, event volume, filtering work, and connection/resource exposure.

Watch scope Typical API Sees
Collection collection.watch() One non-system collection.
Database database.watch() Non-system collections in one database.
Deployment MongoClient.watch() Non-system collections across databases except admin/local/config.
start disposable replica set (l2)
docker rm -f atlasmart-ch18-l2 2>/dev/null || truedocker volume rm atlasmart-ch18-l2-data 2>/dev/null || truedocker run -d --name atlasmart-ch18-l2 \  -p 127.0.0.1:27156:27017 \  -v atlasmart-ch18-l2-data:/data/db \  mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim --replSet atlasmart-rs18-l2 --oplogSize 128 --bind_ip_alluntil mongosh "mongodb://127.0.0.1:27156/admin?directConnection=true" --quiet --eval 'quit(db.runCommand({ping:1}).ok===1?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27156/admin?directConnection=true" --quiet --eval 'rs.initiate({_id:"atlasmart-rs18-l2",members:[{_id:0,host:"atlasmart-ch18-l2:27017"}]})'until mongosh "mongodb://127.0.0.1:27156/admin?directConnection=true" --quiet --eval 'quit(db.hello().isWritablePrimary?0:1)'; do sleep 1; donemongosh "mongodb://127.0.0.1:27156/admin?replicaSet=atlasmart-rs18-l2" --quiet --eval 'printjson(db.version());printjson(db.runCommand({getParameter:1,featureCompatibilityVersion:1}).featureCompatibilityVersion);printjson(rs.status().members.map(m=>({name:m.name,stateStr:m.stateStr})));' 

2. Generate distinguishable AtlasMart changes

collection, database, and deployment streams
from pprint import pprintfrom pymongo import MongoClientclient = MongoClient("mongodb://127.0.0.1:27156/?replicaSet=atlasmart-rs18-l2")db = client.atlasmartinventory = db.inventory_ch18_l2orders = db.orders_ch18_l2inventory.drop(); orders.drop()inventory.insert_one({"_id":"sku-1","available":5})orders.insert_one({"_id":"o-1","status":"new"})pipeline = [{"$match": {    "operationType": {"$in":["update","replace"]},    "ns.db":"atlasmart"}}]with inventory.watch(pipeline, max_await_time_ms=1000) as collection_stream,      db.watch(pipeline, max_await_time_ms=1000) as database_stream,      client.watch(pipeline, max_await_time_ms=1000) as deployment_stream:    inventory.update_one({"_id":"sku-1"},{"$inc":{"available":-1}})    orders.update_one({"_id":"o-1"},{"$set":{"status":"paid"}})    print("collection:"); pprint(collection_stream.next())    print("database #1:"); pprint(database_stream.next())    print("database #2:"); pprint(database_stream.next())    print("deployment #1:"); pprint(deployment_stream.next())    print("deployment #2:"); pprint(deployment_stream.next())client.close()

Expected invariant: the collection stream sees only the inventory update; the database and deployment streams can see both AtlasMart updates because both satisfy the pipeline. Exact event tokens and timestamps differ every run.

3. Filtering happens on change-event fields

The pipeline runs over change event documents, not directly over current collection documents. Therefore a predicate such as {"fullDocument.status":"paid"} only works when the event shape actually contains fullDocument. Filtering on stable metadata such as operationType, ns, and documentKey is often cheaper and less surprising.

legal server-side filter and illegal token removal
const c = db.getSiblingDB("atlasmart").orders_ch18_l2;const good = c.watch([  {$match:{operationType:{$in:["insert","update","replace","delete"]}}},  {$set:{consumerKind:"orders-cdc"}}]);print("The _id resume token is preserved by the pipeline.");good.close();// Deliberately wrong: MongoDB rejects a change-stream pipeline that removes _id.try {  c.watch([{$project:{_id:0,operationType:1}}]).next();} catch (e) {  print(`expected token-preservation failure: ${e.codeName || e.message}`);}
Expanded events and Stable API

DDL-style events such as createIndexes require showExpandedEvents:true. In MongoDB 8.3, fields such as collectionUUID and updateDescription.disambiguatedPaths should not be assumed present unless expanded events are requested. showExpandedEvents itself is not Stable API V1, so API-strict clients must treat that as an explicit compatibility choice.

4. Wrong approach: one deployment-wide firehose per worker

Opening a deployment-wide stream in every worker and filtering most events in application code can consume unnecessary connections and bandwidth, enlarge RBAC scope, and create duplicate processing responsibilities. Prefer one ownership model per consumer group, the narrowest useful watch scope, a selective server-side pipeline, and explicit fan-out downstream if multiple workers need the same logical event feed.

Signal Healthy interpretation Failure smell
Open streams vs pool size Pool has headroom beyond persistent streams. Notification latency rises as streams compete for too few connections.
Filtered/input event ratio Most transported events are relevant. Consumer discards nearly everything after receiving it.
RBAC scope Only required collection/database privileges. Consumer can watch unrelated tenant/business data.
Processing lag Checkpoint stays comfortably inside retained history. Checkpoint age approaches oplog window.

5. Production judgment

Collection streams are the default when ownership is narrow. Database/deployment streams are useful for cross-namespace integrations, but document why that breadth is required. On sharded clusters, mongos creates streams on each shard even when the consumer's eventual filter is selective, so broad streams are not free. Security teams should review find + changeStream privileges at the exact scope. Tests should include empty batches, connection interruption, filter evolution, and unknown/forward-compatible fields.

Bridge. Lesson 3 focuses on event payload semantics—especially the difference between an update delta, a current lookup, and stored pre/post images.

Cleanup/reset

Everything in this lesson is disposable. Remove only the chapter-specific container and volume:

cleanup
docker rm -f atlasmart-ch18-l2 2>/dev/null || truedocker volume rm atlasmart-ch18-l2-data 2>/dev/null || true

Check your understanding

  1. Which scope should an inventory-only consumer use?
  2. Does a change-stream pipeline run against current collection documents?
  3. May the pipeline project away event _id?
  4. Why can a broad stream increase security risk?
  5. Is showExpandedEvents Stable API V1?
Review the answers

1. A collection-scoped stream unless there is a concrete need for broader namespaces.

2. No. It runs against change-event documents.

3. No. _id is the resume token and MongoDB rejects pipelines that modify/remove it.

4. It requires wider find/changeStream privileges and exposes more namespaces to the consumer.

5. No. Treat it as an explicit compatibility/version choice.

Authoritative references

Keep knowledge open

Help the academy stay free and grow.

If these tutorials save you time, a small donation supports new lessons, technical review, diagrams, examples, and long-term maintenance.

ETHEthereum / ERC-20 only
0x716c4Ab160C4B66F31a28AE2448BfF68fc3a2ef0

Send only Ethereum or ERC-20 compatible assets to this address.