Bounded work beats unbounded convenience when result sets, batches, and event rates grow.
Pagination, Batch Operations, Schema Versioning, Change Streams, and Backpressure in Applications
Keep application data flow bounded with keyset pagination, chunked bulk writes, schema-version adapters, resumable change streams, and queue-based backpressure.
Learning objectives
Replace deep offset pagination with deterministic range/keyset pagination when the access pattern permits it.
Bound bulk-write size and concurrency while preserving partial-failure evidence.
Read multiple schema versions safely and evolve writers before destructive migrations.
Consume change streams with resume-token checkpoints and bounded queues.
Apply backpressure so downstream latency does not become unbounded memory or database work.
This final chapter pins
MongoDB Community Server 8.3.8 using
mongodb/mongodb-community-server:8.3.8-ubuntu2204-slim, mongosh 2.10.0, and PyMongo 4.17.0.
Labs use AtlasMart synthetic data, loopback-only Docker port
publishing, replica set name atlasmart-rs27 where
topology behavior matters, and explicit operation time budgets.
Authentication/TLS are disabled only for disposable mechanism
labs; the capstone security gate reuses Chapter 22
least-privilege/authentication requirements and treats TLS as
mandatory production acceptance. Default read preference is
primary and majority acknowledgement is used for business writes
unless a failure experiment explicitly states otherwise.
FCV is observed and never changed. Atlas,
Search, Vector Search, Enterprise Advanced, and KMS are optional
and are not required for mandatory work. Runtime load, failover,
pool, latency, retry, and recovery results were not executed in
the generation environment; learners must record their own
evidence instead of copying invented values.
1. Pagination is part of the query contract
Offset pagination is easy to expose as page=5000,
but the server must scan past skipped results before returning
the page. MongoDB documents that skip() becomes
slower as the offset grows. Range pagination instead remembers
the last stable sort key and uses an indexed bound. The
application must choose a deterministic ordering, including a
unique tie-breaker when the primary sort field is not unique.
from pymongo import MongoClient, ASCENDINGclient=MongoClient("mongodb://localhost:27209,localhost:27210,localhost:27211/?replicaSet=atlasmart-rs27")c=client.atlasmart.orders_ch27_l3c.drop()c.insert_many([{"_id":i,"tenantId":"tenant-a","createdSeq":i,"schemaVersion":1,"totalCents":1000+i%500} for i in range(50000)])c.create_index([("tenantId",ASCENDING),("createdSeq",ASCENDING),("_id",ASCENDING)])page_size=50print(c.find({"tenantId":"tenant-a"}).sort([("createdSeq",1),("_id",1)]).skip(40000).limit(page_size).explain()["executionStats"]["totalDocsExamined"])last_seq,last_id=39999,39999q={"tenantId":"tenant-a","$or":[{"createdSeq":{"$gt":last_seq}},{"createdSeq":last_seq,"_id":{"$gt":last_id}}]}print(list(c.find(q).sort([("createdSeq",1),("_id",1)]).limit(page_size))[:2])client.close()
2. Batch work must have a budget and partial-failure contract
Large jobs should chunk input, cap concurrency, record success/error counts, and be restartable. A single giant list increases client memory and can hide which operations completed. Ordered bulk writes stop at the first write error; unordered bulk writes attempt remaining operations and report errors afterward. Neither choice makes the batch atomic.
from pymongo import MongoClient, UpdateOnefrom pymongo.errors import BulkWriteErrorclient=MongoClient("mongodb://localhost:27209,localhost:27210,localhost:27211/?replicaSet=atlasmart-rs27")c=client.atlasmart.orders_ch27_l3for start in range(0,50000,500): ops=[UpdateOne({"_id":i},{"$set":{"schemaVersion":2,"tenantSlug":"tenant-a"}}) for i in range(start,min(start+500,50000))] try: r=c.bulk_write(ops,ordered=False) print(start,"matched",r.matched_count,"modified",r.modified_count) except BulkWriteError as exc: print(start,"writeErrors",len(exc.details.get("writeErrors",[])))client.close()
3. Schema versioning belongs in the application read path
def normalize_order(doc): version=doc.get("schemaVersion",1) if version==1: return {"id":doc["_id"],"tenant":doc["tenantId"],"totalCents":doc["totalCents"]} if version==2: return {"id":doc["_id"],"tenant":doc.get("tenantSlug",doc["tenantId"]),"totalCents":doc["totalCents"]} raise ValueError(f"unsupported schemaVersion={version}")
A safe rollout often follows expand → dual/read-compatible → migrate/backfill → verify → contract. The deliberately wrong approach is to deploy a reader that understands only v2 before the data population and every producer has moved.
4. Change streams need checkpoints and backpressure
A change stream cursor is resumable while the resume token remains locatable in retained history. That does not make the consumer exactly once. Persist a checkpoint after the idempotent sink has completed, and keep downstream work bounded. If a consumer cannot keep up, queue depth, event age, and resume-token age are operational signals.
from queue import Queue, Fullfrom threading import Thread, Eventfrom bson import json_utilfrom pymongo import MongoClientclient=MongoClient("mongodb://localhost:27209,localhost:27210,localhost:27211/?replicaSet=atlasmart-rs27")coll=client.atlasmart.orders_ch27_l3q=Queue(maxsize=100)stop=Event()checkpoint_file="resume-token.json"def producer(): with coll.watch(full_document="updateLookup",batch_size=20,max_await_time_ms=1000) as stream: while not stop.is_set(): event=stream.try_next() if event is None: continue q.put(event,timeout=2) # bounded wait: downstream pressure becomes visibledef consumer(): while not stop.is_set(): event=q.get() try: # Idempotent sink/upsert would run here. with open(checkpoint_file,"w",encoding="utf-8") as f: f.write(json_util.dumps(event["_id"])) finally: q.task_done()Thread(target=producer,daemon=True).start()Thread(target=consumer,daemon=True).start()print("bounded queue capacity",q.maxsize)
For a real service, define behavior when the queue remains full: shed optional work, scale consumers, slow producers where possible, or fail the health check before memory grows without bound. Never drop a change event silently when correctness depends on it.
5. Failure experiment: offset collapse and consumer backlog
Record offset-page totalDocsExamined at increasing
offsets and compare to the range query. Separately, slow the
consumer deliberately (for example, sleep 100 ms per event) and
plot queue depth/event age. The evidence should make the
mechanism visible without claiming one universal page size,
batch size, or queue capacity.
Check your understanding
- Why can deep skip pagination become expensive?
- What does unordered bulk_write guarantee about order?
- Why keep a schemaVersion?
- When should a change-stream checkpoint advance?
- What is backpressure?
Review the answers
1. The server scans past skipped results before returning the requested page, so work grows with the offset.
2. Nothing; it may execute operations in an arbitrary order and attempts remaining operations after individual errors.
3. It lets readers interpret documents during staged schema evolution instead of assuming every document changes atomically.
4. After the event has been handled idempotently by the required sink, so a crash can safely replay uncheckpointed work.
5. A bounded mechanism that makes downstream saturation reduce/admit upstream work instead of accumulating unbounded in-memory or database work.
6. Production judgment
Pagination, bulk APIs, migrations, and streams are all flow-control problems. Choose stable cursor keys, bound every batch/queue, make partial completion observable, preserve backward-compatible readers during schema evolution, and design replay-safe consumers. Lesson 4 moves from these application mechanisms to the architecture that must support them.
Authoritative references
Driver defaults and deployment behavior evolve. Re-check the exact server patch, PyMongo release, topology, and managed-service tier before freezing production assumptions.
- PyMongo driver documentation
- Connect to MongoDB with PyMongo
- PyMongo connection pools
- PyMongo client-side operation timeout
- PyMongo monitoring
- PyMongo CRUD configuration / retries
- PyMongo transactions
- PyMongo bulk writes
- PyMongo release notes
- MongoDB connection strings
- Connection string options
- Retryable writes
- Retryable reads
- Transactions
- Change streams
- cursor.skip() and range pagination
- Explain results
- Read concern
- Write concern
- Read preference
- Replica sets
- Sharding
- Choose a shard key
- Security checklist
- Backup methods
- serverStatus
- MongoDB 8.3 release notes
- mongosh changelog