
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
| Approach | Best fit | Freshness | Operational burden | Example implementation |
|---|---|---|---|---|
| One-time export | Migration or rare refresh | Point in time | Low, but manual | mongoexport → stage → COPY INTO a VARIANT table |
| Scheduled incremental batch | Deletes don't matter; you run an orchestrator | At the schedule | Medium: DAGs, watermark queries, MERGE | Airflow querying updatedAt past a watermark |
| Self-managed streaming CDC | You already operate Kafka Connect | Continuous to Kafka; Snowflake at your load cadence | High: Kafka, Debezium, your own merge | Debezium MongoDB connector + Snowflake Kafka connector |
| Managed CDC service | Production collections, no pipeline code | Continuous capture, scheduled loads | Low: configuration | Estuary |
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
$matchand$projectbefore events leave the server. - Full documents on update. Update events carry only changed fields (
updateDescription) unless you requestfullDocument: "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.
MongoDB requirements
- A replica set or sharded cluster on WiredTiger. A standalone
mongodhas no change streams. - A user with
findandchangeStreamon the replicated databases:
jsuse 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,configandsystem.*collections, plus time series collections and views, which emit no change events. - Network access: allowlisted IPs or an SSH bastion.
| Where MongoDB runs | Change streams | Oplog retention control | Watch out for |
|---|---|---|---|
| Self-managed replica set | Yes | replSetResizeOplog with minRetentionHours | Run it on every member; it doesn't replicate |
| MongoDB Atlas | Yes | Minimum Oplog Window: M10+, storage auto-scaling, 24 h default | Use the mongodb+srv:// string, not the Atlas SQL interface; ?authSource=admin if the user authenticates against admin |
| Sharded cluster | Yes, through mongos | Per shard | The effective window is the smallest shard's |
| Amazon DocumentDB | Off by default; enable with modifyChangeStreams | No oplog; change_stream_log_retention_duration, 3 h default, 7 days max | Events 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.
jsrs.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 valueWith 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 = 60on the warehouse; Snowflake's minimum billing granularity is 60 seconds. - Don't enable
QUOTED_IDENTIFIERS_IGNORE_CASEfor the pipeline user. - BSON dates are UTC instants: use
TIMESTAMP_NTZholding UTC, orTIMESTAMP_TZ, and match any account-levelTIMESTAMP_TYPE_MAPPING.
Setting up MongoDB to Snowflake replication
- Confirm change streams: a replica set or sharded cluster, or DocumentDB with
modifyChangeStreams. - Size the oplog window for your worst outage.
- Create the
find+changeStreamuser and open network access. - Prepare Snowflake as above.
- Choose the landing model per collection:
VARIANTplus promoted columns (below), andMERGEfor mutable collections or append for immutable ones. - Snapshot, then stream. Record the stream position, load existing documents, then merge each batch of events, keeping the last event per
_id:
sqlMERGE 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.
- Validate. Compare counts per collection. Update a document, delete one, and
$unseta 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 withupdateLookup, orupdateDescriptionapplied by your own code without it. - Deletes arrive as a
deleteevent carrying only_id. The merge removes the row (hard delete) or flags it (soft delete). For the deleted contents, enablechangeStreamPreAndPostImages(MongoDB 6.0+) for pre-images. - Removed fields.
$unsetproduces an update event, not a delete. - New fields land inside
VARIANTwith no DDL, and become typed columns only when you add them to a view orALTER 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 }sqlCREATE 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:Customerwon't matchcustomer. - 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
VARIANTare 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.
sqlCREATE 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 => TRUEkeeps orders with an empty or missing array as one row ofNULLs. Without it, those orders vanish.- Double counting comes from summing a parent measure across exploded rows:
SUM(order_total)overorder_itemscounts each order once per line. Aggregate each grain separately, then join onorder_id. - Never explode two sibling arrays (say
itemsandpayments) 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:
sqlSELECT 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 type | Snowflake type | Notes |
|---|---|---|
| ObjectId | STRING (24 hex chars) | First 4 bytes are a creation timestamp |
| Date | TIMESTAMP_NTZ (UTC) or TIMESTAMP_TZ | Millisecond precision, UTC |
| Int32, Int64 | NUMBER | |
| Double | FLOAT | Watch for NaN and Infinity |
| Decimal128 | NUMBER(p,s) or STRING | STRING beyond 38 digits |
| String, Boolean | STRING, BOOLEAN | VARCHAR defaults to a 16 MB maximum |
| Object, Array | VARIANT | Flatten as above |
| Binary | BINARY or base64 STRING | |
| Null / missing | JSON null vs SQL NULL | Snowflake 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
_idbetween 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.
updateLookupadds one MongoDB read per update.
Production considerations
- Retries. A restarted pipeline resumes from its last checkpointed token.
MERGEtables 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:
- Re-snapshot only the affected collections. Enlarge the oplog window first, then start a new stream.
- 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_idsets. This is only as correct as your application'supdatedAtdiscipline.
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
updateLookupreturnnull, 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;updateDescriptionis projected out to keep events under 16 MB. Nooplog.rsaccess 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,
_idas a string (ObjectId as hex). Top-level fields become columns (Depth 1 field selection); nested objects, arrays and multi-type fields becomeVARIANT. 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),
NaNandInfinityare captured as strings; inference marks numeric strings, so Decimal128 lands asFLOAT(INTEGERif every value is whole) unless you setcastToStringor aDDLoverride (see Choose how documents land below). - Deletes and updates: soft (
_meta/op) unlesshardDelete: true; merge by default, or delta updates that append. - Load interval: the sync schedule defaults to 30 minutes when caught up. Updates to one
_idbetween 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
- Confirm the prerequisites.
- MongoDB meets the MongoDB requirements, with the oplog window sized for your worst outage.
- Estuary's IP addresses are allowlisted, or an SSH tunnel through a bastion is ready.
- Snowflake meets the Snowflake requirements, using the connector's setup script and a key pair registered on the user.
- 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=adminif the user authenticates againstadmin. - User and Password: the
find+changeStreamuser. - 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.
- Address: host and port, or a
- 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.
- 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.
- Host (Account URL):
- 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.
- Keep Decimal128 amounts exact: in the specification editor, add
price: { castToString: true }under the binding'sfields.includeto keep the digits asTEXT, orprice: { DDL: "NUMBER(38,2)" }for a numeric column, which holds values up to about 15 significant digits exactly (custom column types). ADDLoverride applies when the table is created.
- 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.
- Publish. Click Save and publish. Estuary creates the tables, and the initial backfill runs as fast as possible, ignoring the schedule.
- 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
$unsetchecks from setup step 7 after one sync. - Set a Data Processing interval on the materialization's Alerts tab.
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:
mongoexportto JSON,COPY INTOa 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
updatedAthas no oplog dependency, but queries every collection each run and needsupdatedAtset on every insert and update (an ObjectId_idcatches 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.
Related resources
- Estuary docs: MongoDB capture · Amazon DocumentDB · Azure Cosmos DB · Snowflake materialization · Sync schedule · Customize materialized fields · Schema inference · Schema evolution · Backfilling data
- MongoDB: Change streams · Replica set oplog · replSetResizeOplog · Atlas oplog settings · AWS DocumentDB change streams
- Snowflake: Semi-structured types · FLATTEN · VARIANT vs columns · Kafka connector
- Estuary blog: MongoDB CDC · Snowpipe Streaming
To try the route, follow the MongoDB CDC tutorial or create a MongoDB capture in the Estuary dashboard. Pricing: estuary.dev/pricing.
FAQs
How do I exclude specific collections from a MongoDB to Snowflake sync?
Can I stream only changed fields instead of the whole document?
Can a MongoDB document be too large for a Snowflake VARIANT column?

About the author
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.



















