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.

Advanced120–240 minutesDriver/capstone labMongoDB 8.3.8 · mongosh 2.10.0 · PyMongo 4.17.0Last reviewed: September 2026

Learning objectives

01

Replace deep offset pagination with deterministic range/keyset pagination when the access pattern permits it.

02

Bound bulk-write size and concurrency while preserving partial-failure evidence.

03

Read multiple schema versions safely and evolve writers before destructive migrations.

04

Consume change streams with resume-token checkpoints and bounded queues.

05

Apply backpressure so downstream latency does not become unbounded memory or database work.

Reproducible lab baseline

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.

seed and compare offset versus keyset shapes
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.

bounded chunked bulk update
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

read v1 and v2 during migration
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.

bounded producer/consumer skeleton
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

  1. Why can deep skip pagination become expensive?
  2. What does unordered bulk_write guarantee about order?
  3. Why keep a schemaVersion?
  4. When should a change-stream checkpoint advance?
  5. 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.

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.