Type something to search...
Building an Event-Driven Audit Log with MongoDB Change Streams

Building an Event-Driven Audit Log with MongoDB Change Streams

Every team eventually gets the question: "Who changed this order, and what did it look like before?" The usual first answer is to add audit logging inside the application. Every endpoint that writes to the database also writes an audit entry. That works until someone adds a new endpoint and forgets, a migration script bypasses the API, or a support engineer fixes a record directly in mongosh. Now your audit log has holes, and you only find them when you need it most.

Change streams flip the model around. Instead of asking every writer to remember to log, you let MongoDB tell you about every change that happens to a collection, no matter where it came from. A single worker process subscribes to the stream, turns each change event into an audit record, and stores it. Because change streams are built on the replication oplog, they're ordered, resumable, and only report changes that have been durably committed.

This guide covers how change events are structured, how to capture before and after images, how to design the audit collection, how to build a worker that survives restarts without losing or duplicating entries, how to attribute changes to users, and the operational pitfalls you'll hit in production.

What a Change Stream Gives You

A change stream is a cursor that never ends. You open it with watch() on a collection, a database, or the whole deployment, and it yields a change event for every write. Change streams require a replica set or sharded cluster, because they read from the oplog. A single-node replica set is fine for local development.

In mongosh, you can watch a collection and make a change from another shell:

const cursor = db.orders.watch();
// In another shell:
// db.orders.updateOne({ _id: 1001 }, { $set: { status: "shipped" } })
cursor.tryNext();

The event looks like this:

{
  _id: { _data: "8266F0B4A2000000012B042C0100296E5A1004..." },
  operationType: "update",
  clusterTime: Timestamp({ t: 1727072418, i: 1 }),
  wallTime: ISODate("2026-09-22T06:20:18.412Z"),
  ns: { db: "shop", coll: "orders" },
  documentKey: { _id: 1001 },
  updateDescription: {
    updatedFields: { status: "shipped" },
    removedFields: [],
    truncatedArrays: []
  }
}

The important fields for auditing:

  • _id is the resume token. It uniquely identifies this event's position in the stream, and you can hand it back to MongoDB later to continue from exactly this point.
  • operationType tells you what happened: insert, update, replace, delete, plus collection-level events like drop and invalidate.
  • clusterTime and wallTime tell you when the change was committed.
  • documentKey identifies the affected document (the _id, plus the shard key on sharded collections).
  • updateDescription gives you the delta for updates, but not the old values.

That last point matters. An update event says status became "shipped", but not what it was before. For an audit log, you usually want both.

Capturing Before and After Images

MongoDB can attach full documents to change events, in two directions.

Post-images (fullDocument) show the document after the change. Set fullDocument: "updateLookup" and MongoDB fetches the current version of the document when it delivers an update event. The catch is that "current" means at lookup time, not at the moment of the change, so if two updates happen in quick succession, both events may show the result of the second one.

Pre-images (fullDocumentBeforeChange) show the document before the change. These require you to enable pre- and post-image recording on the collection first:

db.runCommand({
  collMod: "orders",
  changeStreamPreAndPostImages: { enabled: true },
});

You can also enable it at creation time with db.createCollection("orders", { changeStreamPreAndPostImages: { enabled: true } }).

Once enabled, MongoDB stores the pre-image and post-image of each change in an internal system collection at write time, so the images are exact, not looked up later. You then request them per stream:

const stream = db.orders.watch([], {
  fullDocument: "whenAvailable",
  fullDocumentBeforeChange: "whenAvailable",
});

With "whenAvailable", the field is included if the image exists and omitted otherwise. Use "required" if you'd rather the stream fail than deliver an event without the image, which is a reasonable choice for a compliance audit log.

Pre-images cost storage, so set a retention period. On self-managed deployments, this is a cluster parameter:

db.adminCommand({
  setClusterParameter: {
    changeStreamOptions: {
      preAndPostImages: { expireAfterSeconds: 86400 },
    },
  },
});

Pre-images are also removed as the oplog rolls past them. Make the retention comfortably longer than any outage your audit worker might plausibly have, since a worker that falls behind the retention window will find images missing.

Designing the Audit Collection

Keep audit records in their own collection, ideally in a separate database so you can apply different access rules and backup policies. A good audit document is self-contained: someone reading it years later shouldn't need to join against anything to understand what happened.

{
  _id: "8266F0B4A2000000012B042C0100296E5A1004...", // resume token
  at: ISODate("2026-09-22T06:20:18.412Z"),
  op: "update",
  ns: "shop.orders",
  docId: 1001,
  actor: { userId: "u_207", via: "admin-ui", requestId: "req_9f3a" },
  changed: { status: { from: "paid", to: "shipped" } },
  before: { _id: 1001, status: "paid", total: 84.5 /* ... */ },
  after: { _id: 1001, status: "shipped", total: 84.5 /* ... */ }
}

Using the resume token's _data string as the _id is the key design decision. It gives every event a natural unique key, which makes the worker idempotent: if it processes the same event twice after a crash, the second insert fails with a duplicate key error and you simply skip it.

Add indexes for the questions people actually ask:

use audit
db.events.createIndex({ ns: 1, docId: 1, at: -1 }) // history of one record
db.events.createIndex({ "actor.userId": 1, at: -1 }) // what did this user do
db.events.createIndex({ at: 1 }, { expireAfterSeconds: 60 * 60 * 24 * 400 })

The last index is a TTL index that expires entries after about 13 months. Adjust or drop it depending on your retention requirements; some regulations require years, in which case you'd archive to cold storage instead. See TTL Indexes in MongoDB for how expiry works.

Building the Audit Worker

Here's a complete worker with the Node.js driver (6.x). It watches several collections in one database, converts events to audit records, and checkpoints its position so it can resume after a restart.

import { MongoClient } from "mongodb";

const client = new MongoClient(process.env.MONGODB_URI);
await client.connect();

const source = client.db("shop");
const audit = client.db("audit");
const events = audit.collection("events", { writeConcern: { w: "majority" } });
const checkpoints = audit.collection("checkpoints");

const WATCHED = ["orders", "customers", "refunds"];
const WORKER_ID = "shop-audit";

const pipeline = [
  {
    $match: {
      "ns.coll": { $in: WATCHED },
      operationType: { $in: ["insert", "update", "replace", "delete"] },
    },
  },
];

async function loadResumeToken() {
  const cp = await checkpoints.findOne({ _id: WORKER_ID });
  return cp?.token;
}

async function saveResumeToken(token) {
  await checkpoints.updateOne(
    { _id: WORKER_ID },
    { $set: { token, savedAt: new Date() } },
    { upsert: true },
  );
}

function diff(before = {}, after = {}, updateDescription) {
  const keys = updateDescription
    ? [
        ...Object.keys(updateDescription.updatedFields),
        ...updateDescription.removedFields,
      ]
    : [...new Set([...Object.keys(before), ...Object.keys(after)])];

  const changed = {};
  for (const key of keys) {
    const top = key.split(".")[0];
    if (top === "_id" || top === "audit") continue;
    const from = before?.[top];
    const to = after?.[top];
    if (JSON.stringify(from) !== JSON.stringify(to)) {
      changed[top] = { from, to };
    }
  }
  return changed;
}

function toAuditRecord(change) {
  const before = change.fullDocumentBeforeChange;
  const after = change.fullDocument;
  return {
    _id: change._id._data,
    at: change.wallTime,
    op: change.operationType,
    ns: `${change.ns.db}.${change.ns.coll}`,
    docId: change.documentKey._id,
    actor: (after ?? before)?.audit ?? { userId: "unknown" },
    changed: diff(before, after, change.updateDescription),
    before: before ?? null,
    after: after ?? null,
  };
}

async function run() {
  const resumeAfter = await loadResumeToken();
  const stream = source.watch(pipeline, {
    fullDocument: "whenAvailable",
    fullDocumentBeforeChange: "whenAvailable",
    ...(resumeAfter && { resumeAfter }),
  });

  console.log(resumeAfter ? "Resuming audit stream" : "Starting fresh");

  for await (const change of stream) {
    try {
      await events.insertOne(toAuditRecord(change));
    } catch (err) {
      if (err.code !== 11000) throw err; // duplicate = already recorded
    }
    await saveResumeToken(change._id);
  }
}

run().catch((err) => {
  console.error("Audit worker stopped:", err);
  process.exit(1);
});

The ordering in the loop is what makes this safe. The audit record is written first, and only then is the checkpoint saved. If the process crashes between the two, the next run resumes from the previous token, sees the same event again, and the duplicate key error tells it the record already exists. You get at-least-once delivery with idempotent writes, which in practice means exactly-once audit records.

Writing the audit record with w: "majority" matters too. Without it, a failover could roll back an audit entry that you've already checkpointed past, and that entry would be lost for good.

Reducing Checkpoint Overhead

Saving the token after every event doubles your writes. For busy collections, checkpoint periodically instead:

let pending = 0;
let lastCheckpoint = Date.now();

for await (const change of stream) {
  await recordEvent(change);
  pending++;
  if (pending >= 100 || Date.now() - lastCheckpoint > 5000) {
    await saveResumeToken(change._id);
    pending = 0;
    lastCheckpoint = Date.now();
  }
}

After a crash you'll reprocess up to 100 events or five seconds of history, and the idempotent _id handles that cleanly. You can also batch the audit inserts with insertMany(docs, { ordered: false }), which continues past duplicate key errors and reports them together.

Handling Quiet Periods

If the watched collections are idle for a long time, your saved token gets old while the oplog keeps moving (other collections are still writing). Eventually the token can fall off the end of the oplog. The driver exposes a post-batch resume token that advances even when no matching events arrive, available as stream.resumeToken. Saving it on a timer keeps your checkpoint fresh:

setInterval(() => {
  if (stream.resumeToken) saveResumeToken(stream.resumeToken);
}, 30_000);

Attributing Changes to Users

Here's the part change streams can't do for you: change events don't say who made the change. MongoDB knows the authenticated database user, but in most applications everything connects as the same service account, so that wouldn't help anyway.

The practical solution is to make the actor part of the data. Every write from your application sets an audit subdocument alongside the real change:

function auditStamp(ctx) {
  return {
    "audit.userId": ctx.user.id,
    "audit.via": ctx.client,
    "audit.requestId": ctx.requestId,
    "audit.at": new Date(),
  };
}

await db
  .collection("orders")
  .updateOne(
    { _id: orderId },
    { $set: { status: "shipped", ...auditStamp(ctx) } },
  );

The worker reads fullDocument.audit to fill in the actor, and the diff function skips the audit field so it doesn't show up as a change. A write that arrives without an audit stamp (a manual fix in the shell, a forgotten script) is still recorded, just with userId: "unknown", which is exactly the kind of gap you want surfaced rather than hidden.

Deletes are trickier, because the pre-image carries the stamp of the last person who edited the document, not the person deleting it. Two approaches work well:

  1. Soft deletes. Instead of removing the document, set deletedAt along with the audit stamp. The audit log records an update with a clear actor. See Implementing Soft Deletes in MongoDB for the full pattern.
  2. Stamp, then delete. Update the audit fields and delete the document inside a transaction. The worker sees both events and can attribute the delete using the stamp from the preceding update.

Watching the Right Scope

You can open a stream at three levels:

ScopeCallUse it when
Collectioncollection.watch()You audit one or two collections
Databasedb.watch()You audit several collections in the same database
Deploymentclient.watch()You need one audit trail across databases

A single database-level stream with a $match on ns.coll is usually more efficient than opening one stream per collection, since each stream is its own cursor reading the oplog. Keep the $match as the first stage so MongoDB can filter events on the server rather than shipping them to the worker.

If you also want schema-level events like index creation or collection renames, pass showExpandedEvents: true. Those are worth recording for security audits, since a dropped index or a renamed collection is often more interesting than an individual document change.

Common Pitfalls

Falling off the oplog. If your worker is down longer than the oplog window, resuming fails with a ChangeStreamHistoryLost error. There's no way to recover those events from the stream. Monitor the replication oplog window, alert when the worker's lag grows, and size the oplog generously on audited clusters. On Atlas you can set a minimum oplog window in the cluster's advanced settings.

Relying on updateLookup for history. The looked-up document reflects the state when the event was delivered, not when the change happened. For an audit log, enable pre- and post-images and use whenAvailable or required. Otherwise your "before" and "after" values can silently lie.

Hitting the 16 MB event limit. A change event is itself a BSON document. With both images attached, an event for a document near 8 MB can exceed the limit and break the stream. Use a $changeStreamSplitLargeEvent stage at the end of your pipeline to split oversized events into fragments, or project the images down to the fields you actually need to audit.

Letting the audit log be edited. An audit trail that the application can modify isn't much of an audit trail. Give the worker a database role with only insert and find on audit.events, and give everyone else read-only access. Pair that with Atlas or database-level auditing if you need to prove nobody tampered with the collection.

Running two workers on the same stream. Two instances with the same checkpoint will both write every event. The idempotent _id prevents duplicates, but you're doing double the work and racing on the checkpoint. Run one active worker, supervised by your process manager or orchestrator, and let it restart on failure.

Forgetting invalidate events. If a watched collection is dropped or renamed, a collection-level stream receives an invalidate event and closes. You can't resume after an invalidate with resumeAfter; you need startAfter with that token. Database-level streams are more forgiving, which is another reason to prefer them.

Querying the Audit Log

Once events are flowing, answering "what happened to order 1001?" is a simple query:

db.events
  .find(
    { ns: "shop.orders", docId: 1001 },
    { at: 1, op: 1, "actor.userId": 1, changed: 1 },
  )
  .sort({ at: -1 });
[
  {
    _id: "8266F0B4A2000000012B042C...",
    at: ISODate("2026-09-22T06:20:18.412Z"),
    op: "update",
    actor: { userId: "u_207" },
    changed: { status: { from: "paid", to: "shipped" } },
  },
  {
    _id: "8266F0A91C000000032B042C...",
    at: ISODate("2026-09-21T17:02:44.905Z"),
    op: "insert",
    actor: { userId: "u_1042" },
    changed: {
      status: { from: null, to: "paid" },
      total: { from: null, to: 84.5 },
    },
  },
];

And "what did this support agent change yesterday?" uses the actor index:

db.events
  .find({
    "actor.userId": "u_207",
    at: { $gte: ISODate("2026-09-21"), $lt: ISODate("2026-09-22") },
  })
  .sort({ at: 1 });

If you need per-document version history inside the source collection itself, rather than in a central log, the approaches in Tracking Document History and Versioning Changes in MongoDB complement this setup well.

Conclusion

A change-stream audit log captures every committed write, regardless of which service, script, or shell made it. Enable pre- and post-images on audited collections, use the resume token as the audit record's _id for idempotency, write with majority write concern, checkpoint after writing, and stamp every application write with an actor so the log can say who did what. Then protect the audit collection with narrow roles and watch the oplog window so the worker never falls behind.

Start small: enable changeStreamPreAndPostImages on your most sensitive collection, run the worker above against a staging cluster, and make a few edits through your app and through mongosh. Compare the entries, and you'll see exactly which writes are missing an actor stamp today.

Tags :
Share :

Related Posts

A Complete Guide to MongoDB Query Operators

A Complete Guide to MongoDB Query Operators

Your first MongoDB queries are usually simple equality filters: find the user with this email, find orders with this status. That covers a surprising

Continue Reading
Async MongoDB in Python with Motor and FastAPI

Async MongoDB in Python with Motor and FastAPI

FastAPI runs your endpoints on an event loop. That's what lets a single worker juggle hundreds of concurrent requests: while one request waits on the

Continue Reading
Atlas Online Archive: Tiering Cold Data to Cut Costs

Atlas Online Archive: Tiering Cold Data to Cut Costs

Look at almost any production database and you'll find the same shape. A small slice of recent data gets nearly all the reads and writes: this week's

Continue Reading