
MongoDB Atlas Data Federation: Querying Data Across S3 and Clusters
Data rarely lives in one place. Your application's live orders are in an Atlas cluster. Last year's event logs are sitting in an S3 bucket as compressed JSON or Parquet files. A partner drops CSV exports into another bucket every night. When someone asks a question that spans all three, the usual answer is an ETL pipeline: copy everything into one system, keep it in sync, and pay to store it twice.
Atlas Data Federation takes a different approach. It gives you a federated database instance, a query endpoint that looks and behaves like a MongoDB deployment, but whose collections are backed by files in cloud object storage, collections in your Atlas clusters, or both. You query it with the same MongoDB Query Language and aggregation pipelines you already know, and Data Federation reads the underlying data where it lives.
This guide covers how federated database instances work, how to write a storage configuration, how to map S3 paths to collections with partition attributes, how to combine cluster and S3 data in one query, how to write results back to S3 with $out, and the performance and cost considerations that decide whether federation is a good fit.
How Data Federation Works
A federated database instance has three pieces:
- Stores describe where data lives: an S3 bucket, an Azure Blob Storage container, a Google Cloud Storage bucket, an Atlas cluster, an Online Archive, or even public HTTP URLs.
- Virtual databases and collections are the names your queries use.
- Data sources map each virtual collection to one or more stores, with path patterns or database and collection names.
When you run a query, Data Federation figures out which underlying sources are relevant, reads them in parallel, and returns results as if they came from a single collection. Nothing is copied or loaded in advance.
Data in object storage is read-only through federation. You can't insertOne into an S3-backed collection. The exception is the $out stage, which can write query results to S3 or an Atlas cluster.
Supported File Formats
For object storage sources, Data Federation can read common analytics formats:
| Format | Notes |
|---|---|
| JSON | One document per line or arrays; gzip compression supported |
| BSON | Useful for mongodump output |
| CSV / TSV | Header row becomes field names; values are strings unless typed |
| Parquet | Columnar, compressed, and by far the most efficient to query |
| Avro, ORC | Supported for existing data lake pipelines |
If you control the format, choose Parquet. It's compressed, it's columnar (so a query that uses three fields doesn't read the other forty), and it's what $out produces by default when you export analytics data.
Setting Up Access to S3
Data Federation reads your bucket through an AWS IAM role that Atlas assumes. The setup, at a high level:
- In Atlas, open Data Federation and click Create New Federated Database.
- Add an S3 data source. Atlas walks you through Cloud Provider Access, which gives you an Atlas AWS account ARN and an external ID.
- In AWS, create an IAM role with a trust policy for that ARN and external ID, and attach a policy granting read access to the bucket (and write access if you'll use
$out). - Paste the role ARN back into Atlas and select the bucket.
The bucket policy should be as narrow as possible:
{
"Version": "2012-10-17",
"Statement": [
{
"Effect": "Allow",
"Action": ["s3:ListBucket", "s3:GetObject", "s3:GetObjectVersion"],
"Resource": [
"arn:aws:s3:::acme-analytics",
"arn:aws:s3:::acme-analytics/events/*"
]
}
]
}
Add s3:PutObject on the export prefix only if you plan to write back with $out.
Writing a Storage Configuration
The UI builds the storage configuration for you, but understanding the JSON makes it much easier to reason about what's happening. Here's a configuration with one S3 store, one Atlas cluster store, and a virtual database that uses both:
{
"stores": [
{
"name": "analyticsBucket",
"provider": "s3",
"bucket": "acme-analytics",
"region": "us-east-1",
"prefix": "events/"
},
{
"name": "prodCluster",
"provider": "atlas",
"clusterName": "Prod",
"projectId": "65f0c3e2a1b2c3d4e5f60718"
}
],
"databases": [
{
"name": "analytics",
"collections": [
{
"name": "events",
"dataSources": [
{
"storeName": "analyticsBucket",
"path": "/{year int}/{month int}/*"
}
]
},
{
"name": "orders",
"dataSources": [
{
"storeName": "prodCluster",
"database": "shop",
"collection": "orders"
}
]
}
]
}
]
}
With this configuration, analytics.events reads Parquet or JSON files under s3://acme-analytics/events/, and analytics.orders reads live data from the shop.orders collection on the Prod cluster.
Partition Attributes: The Key to Fast Queries
The most important part of that configuration is the path:
/{year int}/{month int}/*
This tells Data Federation that the bucket is laid out like events/2026/09/part-0001.parquet, and that the first two path segments are partition attributes. Each file's documents get virtual year and month fields with integer values taken from the path.
Why does that matter? Because when you query on those fields, Data Federation can skip entire folders without opening them:
db.events.aggregate([
{ $match: { year: 2026, month: { $in: [7, 8, 9] } } },
{ $group: { _id: "$eventType", count: { $sum: 1 } } },
{ $sort: { count: -1 } },
]);
This query only reads files under 2026/07, 2026/08, and 2026/09. Without partition attributes, it would scan every file in the bucket, which on a large data lake can mean terabytes of reads and a correspondingly large bill.
Path attributes support types like string, int, isodate, and objectid, and you can use wildcards for segments you don't care about. Design your S3 layout around the filters your queries use most often, usually a date and perhaps a tenant or region.
Connecting and Querying
A federated database instance has its own connection string. You connect with mongosh, Compass, or any driver exactly like a regular cluster:
mongosh "mongodb://federateddatabaseinstance0-abcde.a.query.mongodb.net/analytics" \
--tls --authenticationDatabase admin --username analyst
From there, you run normal read queries and aggregations:
db.events
.find(
{ year: 2026, month: 9, eventType: "checkout_started" },
{ userId: 1, ts: 1, _id: 0 },
)
.limit(5);
[
{ userId: "u_18422", ts: ISODate("2026-09-01T08:12:44Z") },
{ userId: "u_90117", ts: ISODate("2026-09-01T08:13:02Z") },
// ...
];
From application code, it's the same driver you already use. Here's Node.js:
import { MongoClient } from "mongodb";
const client = new MongoClient(process.env.FEDERATION_URI);
const events = client.db("analytics").collection("events");
const topEvents = await events
.aggregate([
{ $match: { year: 2026, month: 9 } },
{ $group: { _id: "$eventType", count: { $sum: 1 } } },
{ $sort: { count: -1 } },
{ $limit: 10 },
])
.toArray();
console.log(topEvents);
await client.close();
Keep federated queries out of latency-sensitive request paths. They're built for analytics, reporting, and batch jobs, not for sub-millisecond lookups.
Joining Cluster Data with S3 Data
The real payoff is querying across sources. Because both events and orders are collections in the same virtual database, you can $lookup between them.
Suppose you want to know how many September checkout events turned into paid orders, broken down by region:
db.events.aggregate([
{ $match: { year: 2026, month: 9, eventType: "checkout_completed" } },
{
$lookup: {
from: "orders",
localField: "orderId",
foreignField: "_id",
as: "order",
},
},
{ $unwind: "$order" },
{ $match: { "order.status": "paid" } },
{
$group: {
_id: "$order.region",
paidOrders: { $sum: 1 },
revenue: { $sum: "$order.total" },
},
},
{ $sort: { revenue: -1 } },
]);
The events come from S3, the orders come from the live cluster, and you wrote one pipeline. No export, no staging table, no sync job.
A word of caution: lookups against a live cluster put read load on that cluster. For heavy analytical joins, point the Atlas store at analytics nodes or secondaries using a read preference in the store configuration, so reporting doesn't compete with production traffic.
Combining Multiple Sources into One Collection
A single virtual collection can have multiple data sources. This is useful when data has moved over time, for example recent orders in the cluster and older orders exported to S3:
{
"name": "allOrders",
"dataSources": [
{
"storeName": "prodCluster",
"database": "shop",
"collection": "orders"
},
{
"storeName": "analyticsBucket",
"path": "/archive/orders/{year int}/*"
}
]
}
Queries against allOrders return the union of both. This is essentially how Online Archive works under the hood: when you connect with the combined cluster-and-archive connection string, you're querying a federated view over the live collection plus archived files. If you want the managed version of this pattern, see Atlas Online Archive: tiering cold data to cut costs.
Writing Results to S3 with $out
Data Federation can also write. The $out stage on a federated database instance supports S3 as a destination, which turns federation into a lightweight export tool.
This pipeline exports last month's paid orders from the live cluster to Parquet files in S3, partitioned by day:
db.orders.aggregate([
{
$match: {
status: "paid",
createdAt: {
$gte: ISODate("2026-08-01T00:00:00Z"),
$lt: ISODate("2026-09-01T00:00:00Z"),
},
},
},
{
$project: {
_id: 1,
customerId: 1,
region: 1,
total: 1,
createdAt: 1,
},
},
{
$out: {
s3: {
bucket: "acme-analytics",
region: "us-east-1",
filename: {
$concat: [
"exports/orders/",
{ $dateToString: { format: "%Y-%m-%d", date: "$createdAt" } },
"/",
],
},
format: { name: "parquet", maxFileSize: "100MiB" },
},
},
},
]);
The filename can be a string or an expression, and using an expression like this creates one prefix per day. Those Parquet files can then be read by Data Federation, Athena, Spark, or any other tool in your data stack.
This is a common way to feed a data lake from MongoDB without writing a custom export service. Schedule it with an Atlas scheduled trigger and you have a nightly export pipeline.
SQL Access
Data Federation also powers Atlas SQL, which lets BI tools that speak SQL query your federated data through a JDBC driver or connectors for tools like Tableau and Power BI. The SQL layer maps documents to relational-looking schemas and supports a read-only SQL dialect.
This is handy when your analysts live in a BI tool and don't want to learn aggregation pipelines. For complex nested documents, though, MQL usually expresses the query more naturally.
Performance Considerations
Federated queries are fundamentally scans over files. Performance depends mostly on how much data they have to read.
- Partition your data. Path-based partition attributes are the single biggest lever. A query that prunes to one month of data can be orders of magnitude faster than a full scan.
- Use Parquet. Columnar formats let Data Federation read only the fields your query needs.
- Right-size your files. Millions of tiny files add per-file overhead. Very large single files limit parallelism. Files in the tens to low hundreds of megabytes are a reasonable target.
- Filter early. Put
$matchon partition fields first, then other filters, then heavier stages. - Keep the federated instance near your data. Choose a region for the federated database instance close to your S3 bucket to reduce latency and cross-region transfer.
There are no indexes on object storage data. If a query needs indexed, low-latency access, that data belongs in a cluster.
Cost Considerations
Data Federation is billed primarily on the amount of data processed by your queries, plus data transfer where applicable. Check the current Atlas pricing page for rates, but the cost drivers are consistent:
- Bytes scanned is the main driver. Partitioning and Parquet reduce it dramatically.
- Data returned and transferred, especially across regions or out to the internet.
- Underlying storage is billed by your cloud provider (S3, GCS, or Azure), not Atlas.
- Load on Atlas clusters used as sources, which may push you toward analytics nodes.
You can set query limits on a federated database instance to cap the amount of data a single query or a time period can process. Set them before you hand the connection string to a team of analysts.
Common Pitfalls
No partition attributes in the path. Without them, every query scans every file. Structure buckets as prefix/{year}/{month}/{day}/ or similar, and declare those segments in the storage configuration.
Mismatched types in path attributes. If you declare {month int} but your folders are named 09, check that the values parse the way you expect, and query with the matching type. A query for month: "9" won't match an int attribute.
Using federation for application reads. It's an analytics endpoint. Serving user-facing requests from S3 through federation will be slow and costly.
Heavy joins against primaries. Cross-source $lookup reads from your live cluster. Route that load to analytics nodes or secondaries.
Forgetting query limits. One accidental full-bucket scan can be the largest line on your bill that month. Configure limits and alerts.
Treating CSV like typed data. CSV values arrive as strings unless typed. Convert with $toInt, $toDate, or $toDecimal in the pipeline, or better, convert the data to Parquet once.
Conclusion
Atlas Data Federation turns object storage and Atlas clusters into one queryable namespace. You describe where data lives in a storage configuration, map S3 paths to collections with partition attributes, and then query everything with MQL, including joins across live and historical data. With $out to S3, it also becomes a simple way to export cluster data into your data lake as Parquet.
Start small: pick one S3 prefix you already have, ideally one with a date-based folder layout, create a federated database instance over it with partition attributes, and run a $match plus $group against a single month. Check how much data the query processed, and you'll immediately see how partitioning shapes both speed and cost.


