Estuary

MongoDB to Snowflake: Change Streams, CDC & Document Modelling Guide

Compare MongoDB to Snowflake methods, including change stream CDC, batch loads, Debezium, and Estuary, with oplog sizing, VARIANT vs flattened modeling, and fixes for arrays and deletes.

MongoDB to Snowflake integration guide cover showing the MongoDB and Snowflake logos
Share this article

MongoDB data can be moved to Snowflake with a one-time export and COPY INTO, a scheduled job that queries recently updated documents, or change data capture (CDC) from MongoDB change streams. For collections that change continuously, CDC is usually the better fit: an initial scan loads existing documents, then every insert, update and delete is merged into a Snowflake table keyed on _id. Polling never sees deletes.

Two constraints shape the route. The oplog must retain more history than your longest pipeline outage, or the stream can't resume. And documents must become tables: land them as VARIANT, promote the fields you filter and join on to typed columns, and give each array its own view. Estuary is one managed option for this route: it reads MongoDB change streams into a durable collection and loads Snowflake tables on a configurable schedule.

Choosing a MongoDB to Snowflake method

ApproachBest fitFreshnessOperational burdenExample implementation
One-time exportMigration or rare refreshPoint in timeLow, but manualmongoexport → stage → COPY INTO a VARIANT table
Scheduled incremental batchDeletes don't matter; you run an orchestratorAt the scheduleMedium: DAGs, watermark queries, MERGEAirflow querying updatedAt past a watermark
Self-managed streaming CDCYou already operate Kafka ConnectContinuous to Kafka; Snowflake at your load cadenceHigh: Kafka, Debezium, your own mergeDebezium MongoDB connector + Snowflake Kafka connector
Managed CDC serviceProduction collections, no pipeline codeContinuous capture, scheduled loadsLow: configurationEstuary

Only change streams report deletes, so the rest of this guide assumes CDC.

How MongoDB to Snowflake CDC works

CDC on this route reads change streams, MongoDB's supported API over the oplog. Tailing local.oplog.rs yourself gains nothing and loses:

  • Resume tokens, to restart after the last processed event.
  • Total ordering across shards through mongos, instead of merging per-shard oplogs.
  • Server-side $match and $project before events leave the server.
  • Full documents on update. Update events carry only changed fields (updateDescription) unless you request fullDocument: "updateLookup".

Every approach shares one limit: a stream can only resume from a point the oplog still holds.

A new pipeline records the stream position before scanning collections, so writes during the scan aren't missed. It loads the scan, merges change events on _id from that position, and checkpoints the resume token after each batch lands in Snowflake, as a staged MERGE or as appended rows.

A MongoDB change stream and collection scan merged on _id or appended in Snowflake.
Record the stream position before the scan: a stream can only resume from a point the oplog still holds.

MongoDB requirements

  • A replica set or sharded cluster on WiredTiger. A standalone mongod has no change streams.
  • A user with find and changeStream on the replicated databases:
js
use admin db.createRole({ role: "snowflakeReplicator", roles: [], privileges: [{ resource: { db: "shop", collection: "" }, actions: ["find", "changeStream"] }] }) db.createUser({ user: "replicator", pwd: passwordPrompt(), roles: ["snowflakeReplicator"] })
  • Polling for what streams can't watch:admin, local, config and system.* collections, plus time series collections and views, which emit no change events.
  • Network access: allowlisted IPs or an SSH bastion.
Where MongoDB runsChange streamsOplog retention controlWatch out for
Self-managed replica setYesreplSetResizeOplog with minRetentionHoursRun it on every member; it doesn't replicate
MongoDB AtlasYesMinimum Oplog Window: M10+, storage auto-scaling, 24 h defaultUse the mongodb+srv:// string, not the Atlas SQL interface; ?authSource=admin if the user authenticates against admin
Sharded clusterYes, through mongosPer shardThe effective window is the smallest shard's
Amazon DocumentDBOff by default; enable with modifyChangeStreamsNo oplog; change_stream_log_retention_duration, 3 h default, 7 days maxEvents over 16 MB fail the read

Sizing the oplog window

Size the oplog by time, not bytes. Its window, from oldest to newest entry, must outlast the longest time your pipeline could stay stopped, and must also cover the snapshot if your tool snapshots before opening the stream.

js
rs.printReplicationInfo() // configured size and "log length start to end" // Self-managed MongoDB: run on EVERY member; it does not replicate db.adminCommand({ replSetResizeOplog: 1, minRetentionHours: 48 }) // example value

With minRetentionHours set, the oplog grows past its configured size until entries age out; budget disk for peak writes. The WiredTiger default is 5% of free disk (990 MB to 50 GB). Multi-document updates, in-place updates and heavy deletes fill it fastest. Multiply peak oplog growth per hour by the hours you need.

Snowflake requirements

Create a database, schema, virtual warehouse, and a dedicated role and user with grants on them, authenticated with a key pair (JWT).

  • Set AUTO_SUSPEND = 60 on the warehouse; Snowflake's minimum billing granularity is 60 seconds.
  • Don't enable QUOTED_IDENTIFIERS_IGNORE_CASE for the pipeline user.
  • BSON dates are UTC instants: use TIMESTAMP_NTZ holding UTC, or TIMESTAMP_TZ, and match any account-level TIMESTAMP_TYPE_MAPPING.

Setting up MongoDB to Snowflake replication

  1. Confirm change streams: a replica set or sharded cluster, or DocumentDB with modifyChangeStreams.
  2. Size the oplog window for your worst outage.
  3. Create the find + changeStream user and open network access.
  4. Prepare Snowflake as above.
  5. Choose the landing model per collection: VARIANT plus promoted columns (below), and MERGE for mutable collections or append for immutable ones.
  6. Snapshot, then stream. Record the stream position, load existing documents, then merge each batch of events, keeping the last event per _id:
sql
MERGE INTO shop.orders t USING ( SELECT * FROM shop.orders_changes QUALIFY ROW_NUMBER() OVER (PARTITION BY _id ORDER BY change_seq DESC) = 1 ) s ON t._id = s._id WHEN MATCHED AND s.op = 'delete' THEN DELETE WHEN MATCHED THEN UPDATE SET t.doc = s.doc, t.updated_at = s.event_time WHEN NOT MATCHED AND s.op <> 'delete' THEN INSERT (_id, doc, updated_at) VALUES (s._id, s.doc, s.event_time);

change_seq is the event's stream position. clusterTime alone isn't unique, because events in one transaction share it.

  1. Validate. Compare counts per collection. Update a document, delete one, and $unset a field, then check each in Snowflake after one load interval. Stop the pipeline briefly and confirm it resumes from its token.

What happens after the initial load?

  • Updates merge by _id: the whole document with updateLookup, or updateDescription applied by your own code without it.
  • Deletes arrive as a delete event carrying only _id. The merge removes the row (hard delete) or flags it (soft delete). For the deleted contents, enable changeStreamPreAndPostImages (MongoDB 6.0+) for pre-images.
  • Removed fields.$unset produces an update event, not a delete.
  • New fields land inside VARIANT with no DDL, and become typed columns only when you add them to a view or ALTER TABLE. When a field changes type, add a new column rather than altering the old one.
  • Ordering. Change streams guarantee order, including across shards. Parallel batch loads, or deduplicating on a non-unique timestamp, lose it. To keep every version, append events to a history table.

Modelling MongoDB documents in Snowflake: VARIANT vs flattened columns

Do both. Land the document, or its nested parts, as VARIANT so no field is lost, then promote the fields queries filter, join or aggregate on into typed columns or views. Snowflake advises VARIANT when you're unsure how data will be queried, and extracting dates, numbers stored as strings and arrays. JSON has no date type, so dates inside VARIANT stay strings, slower to query and larger to store.

Snowflake extracts up to 200 VARIANT elements per partition into columnar storage, but skips any element containing even one null, or mixed types. Optional, inconsistently typed MongoDB fields often do, and queries on them scan the whole value.

Flattening nested objects

Use path notation and casts, usually in a view over the landing table. For this order document:

json
{ "_id": "65a1f0c2e4b0a1b2c3d4e5f6", "placedAt": "2026-03-02T14:05:11.000Z", "customer": { "name": "Ada", "address": { "city": "Leeds", "postcode": "LS1" } }, "items": [ { "sku": "A-1", "qty": 2, "price": 9.5 }, { "sku": "B-7", "qty": 1, "price": 30 } ], "total": 49.0 }
sql
CREATE OR REPLACE VIEW shop.orders_flat AS SELECT _id AS order_id, doc:placedAt::TIMESTAMP_NTZ AS placed_at, doc:customer.name::STRING AS customer_name, doc:customer.address.city::STRING AS ship_city, doc:total::NUMBER(12,2) AS order_total FROM shop.orders;
  • Element names are case-sensitive:doc:Customer won't match customer.
  • Double-quote keys with spaces or special characters, for example doc:"shipping info".city.
  • Flatten only the levels people query. MongoDB allows 100 levels of nesting; the rest can stay in VARIANT.
  • Warehouse or pipeline? Views over VARIANT are cheap to change. Pipeline flattening gives real columns but ties table shape to pipeline configuration.

Arrays without double counting

Give each array its own grain: a child view with one row per element, keyed by parent _id plus index.

sql
CREATE OR REPLACE VIEW shop.order_items AS SELECT o._id AS order_id, li.index AS line_no, li.value:sku::STRING AS sku, li.value:qty::NUMBER AS qty, li.value:price::NUMBER(12,2) AS unit_price FROM shop.orders o, LATERAL FLATTEN(input => o.doc:items, outer => TRUE) li;
  • outer => TRUE keeps orders with an empty or missing array as one row of NULLs. Without it, those orders vanish.
  • Double counting comes from summing a parent measure across exploded rows: SUM(order_total) over order_items counts each order once per line. Aggregate each grain separately, then join on order_id.
  • Never explode two sibling arrays (say items and payments) in one query; you get their cartesian product.
  • Arrays change in place. A view follows rewrites automatically; a materialized child table needs a delete-and-insert per changed parent _id.

Fields with mixed types

Keep the field in VARIANT, measure the spread, cast defensively:

sql
SELECT TYPEOF(doc:price) AS json_type, COUNT(*) FROM shop.products GROUP BY 1; SELECT _id, TRY_TO_DECIMAL(doc:price::STRING, 12, 2) AS price FROM shop.products;

Count the NULLs that TRY_ functions return as a data-quality check. The durable fix is at the source: correct the documents, add MongoDB schema validation, reload the collection.

BSON to Snowflake type mapping

BSON typeSnowflake typeNotes
ObjectIdSTRING (24 hex chars)First 4 bytes are a creation timestamp
DateTIMESTAMP_NTZ (UTC) or TIMESTAMP_TZMillisecond precision, UTC
Int32, Int64NUMBER
DoubleFLOATWatch for NaN and Infinity
Decimal128NUMBER(p,s) or STRINGSTRING beyond 38 digits
String, BooleanSTRING, BOOLEANVARCHAR defaults to a 16 MB maximum
Object, ArrayVARIANTFlatten as above
BinaryBINARY or base64 STRING
Null / missingJSON null vs SQL NULLSnowflake distinguishes the two

Performance, latency and Snowflake cost

Capture is continuous. Snowflake freshness depends on how often you merge, and each merge spends warehouse credits.

  • Merge less often. Repeated updates to one _id between merges collapse into one write, which matters for carts, sessions and counters.
  • Auto-suspend at 60 seconds, so an idle warehouse stops billing.
  • Mind micro-partitions. Merges that touch fewer partitions cost less; clustering keys cost extra credits.
  • Append instead of merging for collections whose documents are never updated.
  • Count source reads.updateLookup adds one MongoDB read per update.

Production considerations

  • Retries. A restarted pipeline resumes from its last checkpointed token. MERGE tables absorb a replayed batch; append tables duplicate it unless the read position and Snowflake write commit together.
  • Source impact. Snapshots are full collection scans sorted on an indexed field; schedule large ones off-peak.
  • Backfills. Re-snapshot one collection at a time and merge into existing tables rather than truncating, so they stay queryable.
  • Monitoring. Alert when pipeline lag nears the oplog window, and on time since the last successful Snowflake load.
  • Security. Least-privilege role, Snowflake key-pair auth, and an IP allowlist, SSH bastion or private networking.

Common problems and troubleshooting

Change stream can't resume after the oplog window is exceeded

The stream fails with ChangeStreamHistoryLost. Only a source re-read proves nothing was missed. From safest to cheapest:

  1. Re-snapshot only the affected collections. Enlarge the oplog window first, then start a new stream.
  2. Bounded catch-up, if every write sets an indexed updatedAt: start a new stream at "now", merge documents updated after your last checkpoint, and reconcile deletes by comparing _id sets. This is only as correct as your application's updatedAt discipline.

Change events larger than 16 MB

A change event carrying a document, its pre-image and updateDescription can exceed 16 MB, failing the read with BSONObjectTooLarge. Project away updateDescription, drop pre-images you don't need, or add $changeStreamSplitLargeEvent (MongoDB 7.0, 6.0.9+) as the last stage and reassemble the fragments. MongoDB advises documents under 8 MB when requesting both pre- and post-images.

Lag and missing events on sharded clusters

  • Cold shards with few writes delay the merged stream; lower periodicNoopIntervalSecs.
  • Shard-key changes make updateLookup return null, and orphaned documents from chunk migrations emit no events (MongoDB 5.3+).
  • Snapshots need an explicit sort, because results merge from several shards.

Using Estuary for MongoDB to Snowflake

Estuary pairs the MongoDB capture with the Snowflake materialization. The capture opens database-level change streams and writes each MongoDB collection to an Estuary collection keyed on _id; the materialization stages each transaction in the internal stage flow_v1 and applies it with MERGE, or COPY INTO if all its keys are new.

  • Oplog window: Estuary's connector docs recommend at least 24 hours, and more for reliability.
  • Change events: updates always carry the full document (updateLookup, or stored post-images when every captured collection in the database has them), plus pre-images where enabled; updateDescription is projected out to keep events under 16 MB. No oplog.rs access is needed.
  • Batch collections: views and time series use Batch Snapshot or Batch Incremental mode, polled every 24 hours by default. Only change-stream mode captures deletes.
  • Document shape: one table per collection, _id as a string (ObjectId as hex). Top-level fields become columns (Depth 1 field selection); nested objects, arrays and multi-type fields become VARIANT. Depth 2, Unlimited Depth or selected nested fields add columns: promotion by field selection instead of views.
  • Types: dates become millisecond UTC timestamps. Decimal128, binary (base64), NaN and Infinity are captured as strings; inference marks numeric strings, so Decimal128 lands as FLOAT (INTEGER if every value is whole) unless you set castToString or a DDL override (see Choose how documents land below).
  • Deletes and updates: soft (_meta/op) unless hardDelete: true; merge by default, or delta updates that append.
  • Load interval: the sync schedule defaults to 30 minutes when caught up. Updates to one _id between syncs combine into one write.
  • Limits: schema inference tracks up to 1,000 field locations for MongoDB; fields beyond that are pruned from field selection until added to the read schema. The DocumentDB variant is batch-only and requires an SSH tunnel; Cosmos DB has its own variant.

Implementing MongoDB to Snowflake with Estuary

  1. Confirm the prerequisites.
  2. Create the MongoDB capture. In the Estuary dashboard, open Sources, click New Capture, choose MongoDB and name the capture. Fill in the endpoint properties:
    • Address: host and port, or a mongodb+srv:// string, with ?authSource=admin if the user authenticates against admin.
    • User and Password: the find + changeStream user.
    • Database (optional): a comma-separated list; empty discovers every database.
    • Capture Batch Collections in Addition to Change Stream Collections: check to include views and time series collections.
    • Through a bastion, add the SSH Endpoint and private key in the network tunnel settings.
MongoDB capture form with Address, User, Password (masked) and Database filled in.
  1. Select collections. Click Next. Estuary discovers one binding per collection, all selected by default.
    • Remove collections you don't need; the filter box above the list narrows a long one.
    • Check each binding's Capture Mode. Change Stream Incremental is the only mode that captures deletes. Views and time series use Batch Snapshot or Batch Incremental, with a Cursor Field and optional Polling Schedule (capture modes).
    • Click Save and publish. The backfill and change stream capture start automatically.
MongoDB capture bindings with a collection whose Capture Mode is Change Stream Incremental.
  1. Create the Snowflake materialization. Click Materialize on the published capture, which pre-fills its collections, and choose Snowflake. Fill in the endpoint properties:
    • Host (Account URL): orgname-accountname.snowflakecomputing.com, without the protocol.
    • Database, Schema, Warehouse and Role from the setup script.
    • User and Private Key: key-pair (JWT); the User Password option is deprecated.
    • Snowflake Timestamp Type: TIMESTAMP_NTZ (normalize to UTC) is recommended for new tasks (timestamp mapping).
    • Hard Delete: on deletes rows in Snowflake; off (the default) keeps them, flagged in _meta/op.
Snowflake endpoint form with connection fields, Timestamp Type and a redacted private key.
  1. Choose how documents land. Per binding, open the Config tab's Field Selection table (field selection). The tradeoffs are covered in VARIANT vs flattened columns.
    • Set Field Depth: Depth 1 (the default), Depth 2 or Unlimited Depth.
    • Require or exclude individual fields.
    • Enable delta updates only on append-only collections.
Field Selection table at Depth 1, with the nested customer field as one column.
  • Keep Decimal128 amounts exact: in the specification editor, add price: { castToString: true } under the binding's fields.include to keep the digits as TEXT, or price: { DDL: "NUMBER(38,2)" } for a numeric column, which holds values up to about 15 significant digits exactly (custom column types). A DDL override applies when the table is created.
  1. Set the sync schedule. Under Sync Schedule (reference):
    • Sync Frequency: 30 minutes by default once caught up. Shorter intervals spend more warehouse credits (cost).
    • Timezone, Fast Sync Start Time, Fast Sync Stop Time, Fast Sync Enabled Days (optional): use that frequency inside a window, every 4 hours outside it.
Snowflake Sync Schedule with Sync Frequency set to 30m.
  1. Publish. Click Save and publish. Estuary creates the tables, and the initial backfill runs as fast as possible, ignoring the schedule.
  2. Verify data is landing.
    • Both tasks show a green dot (tooltip PRIMARY) on the Sources and Destinations pages; their details pages (capture, materialization) chart Data Written and Data Read.
    • Compare row counts with countDocuments(), excluding "_meta/op" = 'd' rows under soft deletes.
    • Run the update, delete and $unset checks from setup step 7 after one sync.
    • Set a Data Processing interval on the materialization's Alerts tab.
Snowflake materialization Overview: Connector Status Running and a Data Read chart.

How Estuary handles backfills, oplog loss and schema-less collections

Backfill in rounds alongside the change stream

The capture doesn't finish the snapshot before reading change streams. It alternates: catch the streams up to the cluster's latest operation time, backfill for up to five minutes, repeat. A collection's backfilled documents and change events are never emitted concurrently, so a scanned document can't overwrite a newer change. Resume tokens stay current throughout, so the oplog window must cover an outage, not the whole snapshot. Progress is checkpointed every 50,000 documents or 4 MiB with the last _id read, so a restart continues the scan.

When the oplog window is exceeded

If the capture stays stopped longer than the oplog retains, it fails on resume instead of silently skipping the gap. Recover with an incremental backfill of the affected bindings: it re-reads those collections while Snowflake tables stay in place, and merge bindings absorb re-read documents without duplicates (delta-update tables get duplicates). Documents deleted in MongoDB during the gap aren't removed: a backfill only emits the documents it finds, so their rows stay in Snowflake. Finding them means comparing _id values between MongoDB and Snowflake.

Schema inference for schema-less collections

MongoDB guarantees only _id, so the collection's write schema requires just a string _id; the read schema is inferred from captured documents. Until a real inferred schema is published, a placeholder rejects all documents, so nothing is materialized against a guess. With auto-discovery on, new fields become Snowflake columns without a backfill. Inference only widens: a field first seen as a string stays one until you fix the source and reset the dataflow. An incompatible change backfills the table by default (onIncompatibleSchemaChange can abort or disable instead); a key change creates a _v2 collection, keeping table names.

Rebuilding Snowflake tables without re-reading MongoDB

The collection is a durable log of JSON files in your cloud storage bucket, and the materialization reads it, not MongoDB. A materialization backfill truncates and repopulates Snowflake tables from the collection, with no MongoDB reads. A second materialization reads the same history, then follows changes, without another change stream. Limits: replay covers only what the collection retains (trial storage about 20 days; your own bucket keeps full history), and an oplog gap still needs the re-read above.

When another approach may make more sense

  • One-time migration:mongoexport to JSON, COPY INTO a single-VARIANT-column table, then views. One full scan per collection and nothing to operate afterwards; deleted documents disappear only when you replace the table.
  • Daily reporting, deletes irrelevant: a batch job on an indexed updatedAt has no oplog dependency, but queries every collection each run and needs updatedAt set on every insert and update (an ObjectId _id catches inserts, not updates). Estuary's Batch Incremental mode does this without a DAG.
  • An existing Kafka and Debezium stack: add the Snowflake Kafka connector and a merge layer (the connector only inserts) rather than a second CDC system. If Debezium's offset leaves the oplog, it re-snapshots only with snapshot.mode: when_needed (the default fails); Snowflake rebuilds are limited by topic retention.
  • No change streams (unwatchable collections, or DocumentDB with streams off): poll on an indexed, strictly increasing field and accept missed deletes. Estuary's DocumentDB connector captures DocumentDB this way, in Batch Snapshot or Batch Incremental mode only.

To try the route, follow the MongoDB CDC tutorial or create a MongoDB capture in the Estuary dashboard. Pricing: estuary.dev/pricing.

FAQs

    Can I move a MongoDB to Snowflake pipeline off another tool without re-syncing everything?

    Resume tokens don't move between tools. Usually you run the new pipeline into new tables beside the old ones, compare counts and samples, then repoint views. To keep the old tables, start the new pipeline from "now" without a backfill and reconcile the cutover gap. Estuary supports this with skipBackfills (per database:collection, or :).
    Open streams only on the databases you need, and filter by namespace: db.watch([{ $match: { "ns.coll": { $nin: ["audit_log", "sessions"] } } }]). The server still reads every oplog entry to evaluate the filter, but excluded events never leave it. In Estuary, set excludeCollections, or exclusiveCollectionFilter when only a few bindings are enabled.
    Yes, but it rarely pays. Applying updateDescription (updatedFields, removedFields, truncatedArrays) means ordered OBJECT_INSERT/OBJECT_DELETE calls on VARIANT, handling dotted paths and truncated arrays. Snowflake micro-partitions are immutable, so a merge rewrites affected partitions whether one field changed or fifty.
    No. A BSON document is capped at 16 MiB, and a VARIANT value holds up to 128 MB uncompressed. Files larger than 16 MiB live in GridFS as chunk documents, which replicate as ordinary rows.

Start streaming your data for free

Build a Pipeline

About the author

Picture of Seán Whelan
Seán WhelanSolutions Engineer

Seán Whelan is a Solutions Engineer at Estuary, helping customers design, troubleshoot, and run production data pipelines. Seán works across change data capture, database replication, and SaaS integrations, and writes about what actually comes up when real teams move data.

Real-Time & Batch Pipelines.Simple to Deploy.Simply Priced.

  • $0.50/GB of data moved + $100/month per connector instance
  • Simple usage-based pricing with volume discounts
  • Sub-100ms end-to-end latency for real-time pipelines