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.
Learning objectives
Distinguish collection-, database-, and deployment-scoped change streams and their authorization blast radius.
Filter change events with supported pipeline stages without removing or mutating the resume token.
Use namespace and operationType evidence to prove which events a broader stream sees.
Explain the connection-pool and sharded-router costs of many broad streams.
Diagnose an over-broad deployment watch and replace it with the narrowest useful scope and filter.
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. |
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
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.
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}`);}
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:
docker rm -f atlasmart-ch18-l2 2>/dev/null || truedocker volume rm atlasmart-ch18-l2-data 2>/dev/null || true
Check your understanding
- Which scope should an inventory-only consumer use?
- Does a change-stream pipeline run against current collection documents?
- May the pipeline project away event _id?
- Why can a broad stream increase security risk?
- 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
- MongoDB Change Streams — deployment requirements, majority-committed notification, scopes, sharded behavior, resume tokens, and pre/post images.
- Change Stream Events — event fields, operation types, resume token, update/replace behavior, and expanded events.
- db.collection.watch() — pipeline stages, options, resumability, and mongosh versus driver behavior.
- db.watch() — database-scoped streams.
- Mongo.watch() — deployment-scoped streams.
- update Event — updateDescription, documentKey, full document and pre-image behavior.
- delete Event — delete event and pre-image behavior.
- invalidate Event — stream invalidation and startAfter boundary.
- Change Streams Production Recommendations — sharded total ordering and latency considerations.
- Privilege Actions — changeStream/find authorization requirements.
- Replica Set Oplog — retained history and oplog window.
- PyMongo Driver — official Python driver baseline and change-stream cursor APIs.
- PyMongo 4.17 Release Notes — current driver line used by the course.
- MongoDB 8.3 Release Notes — current stable minor and patch status.
- MongoDB 8.3 Compatibility Changes — current expanded-event field behavior inherited from 8.2.x.
- mongosh Release Notes — mongosh 2.10.0 baseline.