Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Use the official MongoDB Go driver’s Watch() method to open a change stream on a collection, database, or deployment. Iterate with Next(), decode each event, and durably save the event’s _id resume token together with the work it represents. On restart, pass that token back with resumeAfter (or use startAfter when starting after an invalidate event), while keeping the original pipeline and options unchanged.

Choose the stream scope

The object on which you call Watch() determines which changes can arrive:

Scope Go call What it observes
Collection coll.Watch(ctx, pipeline, options...) Changes for one collection.
Database db.Watch(ctx, pipeline, options...) Eligible collection changes in that database. System collections and the admin, local, and config databases are excluded.
Deployment/client client.Watch(ctx, pipeline, options...) Eligible changes across databases.

These scopes and their eligibility rules are described in MongoDB’s change-stream documentation. A broad scope can produce substantially more events, so filter deliberately rather than opening a deployment-wide stream by default.

Open and consume a stream in Go

An empty pipeline returns all changes available at the selected scope. Add aggregation stages, commonly a $match, to reduce the event types or namespaces your consumer receives.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.
cs, err := coll.Watch(ctx, mongo.Pipeline{})
if err != nil {
    return err
}
defer cs.Close(ctx)

for cs.Next(ctx) {
    var event bson.M
    if err := cs.Decode(&event); err != nil {
        return err
    }

    // Process the event, then durably record event["_id"]
    // as the resume token for the work just completed.
}

return cs.Err()

This is the consumption pattern shown in the MongoDB Go Driver change-stream guide. Adapt cancellation, shutdown, decoding, retry policy, and token storage to your application. Check the error from Watch(), close the stream when finished, and inspect Err() after iteration ends; a loop ending is not by itself proof that processing completed normally.

Filter by operation type

For example, a pipeline can pass only inserts, updates, and replacements:

pipeline := mongo.Pipeline{
    {{Key: "$match", Value: bson.D{
        {Key: "operationType", Value: bson.D{
            {Key: "$in", Value: bson.A{"insert", "update", "replace"}},
        }},
    }}},
}
cs, err := coll.Watch(ctx, pipeline)

Keep the exact pipeline available: MongoDB warns that changing the pipeline or options when resuming can produce unpredictable behavior, affect consistency, or prevent resumption.

Understand what an event contains

Resume token

Every change event’s _id value is its resume token. Persist it reliably with the side effects represented by that event. If the process crashes, restart the stream with the last token that is known to have been processed.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Update delta versus a full document

Update notifications normally describe a delta, such as changed and removed fields, rather than including the complete document. Requesting UpdateLookup asks MongoDB to look up a full post-update document:

opts := options.ChangeStream().SetFullDocument(options.UpdateLookup)
cs, err := coll.Watch(ctx, mongo.Pipeline{}, opts)

The looked-up document is the most current majority-committed version available when MongoDB performs the lookup. It can therefore include writes that occurred after the update that generated the event; it is not an event-time snapshot. The update delta remains the description of the original update.

Pre-images and post-images

For exact before-and-after data, configure the collection with changeStreamPreAndPostImages, then select the appropriate full-document options. Availability depends on the event and configuration: inserts have no pre-image, and deletes have no post-image. WhenAvailable permits an event when an image is unavailable, while Required causes the operation to fail if the requested image cannot be supplied. Pre-image options use FullDocumentBeforeChange. The driver guide documents these settings and their prerequisites.

Resume safely after disconnects

  1. Process an event and complete the side effect your application considers successful.
  2. Persist that event’s _id token atomically with the processing result, or use an equivalent idempotent design.
  3. On reconnect, recreate the stream with the same pipeline and options.
  4. Set resumeAfter to the last durably processed token. MongoDB continues after that event.
  5. Use startAfter when you need to start a new stream after an invalidate event.

MongoDB must still retain the operation history referenced by the token. If the oplog no longer contains that history, the token cannot be used to resume; choose an application-specific recovery path rather than silently skipping events. The MongoDB change-stream manual describes these resume semantics.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

Make processing and token storage agree

Saving a token before its event’s side effect can lose work after a crash. Saving it only after the side effect can replay the event, so handlers should be idempotent or otherwise tolerate duplicates. The durable checkpoint should represent the last event your application is prepared to consider complete.

Independent reader supportYour contribution helps us test, update, and keep practical guides available for everyone.Support on Ko-Fi

Handle update lookups and filtered streams carefully

MongoDB documents a specific failure mode when fullDocument: "updateLookup" is combined with a $match filter. During rapid deletions or traffic spikes, a deleted document can produce a null fullDocument, and the stream may fail with a Resume Token Not Found error. The production recommendations identify pre-/post-images with whenAvailable as an alternative to evaluate when this combination is important to your workload.

This is not a reason to avoid filtering or lookups universally. It is a reason to test the exact pipeline, event volume, delete behavior, and recovery path you will run in production.

Control event size

Change-stream response documents are limited to BSON’s 16 MB document limit. Large source documents and full-document lookups can make an event too large even when the original write succeeded. MongoDB documents the $changeStreamSplitLargeEvent stage beginning in Server 6.0.9. Confirm your deployment version and test the resulting event shape before depending on that stage. See the production recommendations for the operational guidance.

Free tools Windows power users keep installed

One-click scans. No signup required.

Special offer. See more information about Outbyte and uninstall instructions. Please review EULA and Privacy policy.

A practical option decision guide

Need Prefer Important limitation
Small, precise change notifications Default update deltas You must apply the delta or fetch a document yourself.
A current full document after updates UpdateLookup The result can include later majority-committed writes and is not an event-time snapshot.
Exact before/after images Configured pre-/post-images with suitable availability options The collection must enable changeStreamPreAndPostImages; an image may not exist for a particular event.
Restart after an ordinary interruption resumeAfter plus the saved _id The oplog must still contain the referenced history.
Start after invalidation startAfter Use the original stream definition and handle invalidation explicitly.

Production checklist

  • Select the narrowest scope that satisfies the requirement.
  • Use a pipeline to exclude irrelevant operation types or namespaces.
  • Decode into a type that preserves the fields your handler needs, including _id.
  • Close streams and cancel contexts during shutdown.
  • Persist resume tokens with completed work and make handlers duplicate-safe.
  • Recreate a stream with the same pipeline and options when resuming.
  • Monitor for oplog-history loss, Resume Token Not Found errors, and event-size failures.
  • Test updates, replacements, deletes, invalidation, reconnects, traffic spikes, and large documents against the MongoDB Server version you deploy.

Product prices and availability are accurate as of the date/time indicated and are subject to change. Any price and availability information displayed on Amazon at the time of purchase will apply.