> Markdown version of https://archtenet.dev/blog/mongodb-geosharding-hidden-cost — the same page without the site chrome.
> Index of everything published here: https://archtenet.dev/llms.txt

# The Hidden Complexity of "Easy" Geo-Sharding in MongoDB

> Geo-sharding looks trivial to adopt — add a shard key, keep data close to users. That simplicity is exactly where the expensive failures hide. A field story about symptoms, root cause, and the bill that arrives monthly.

- **HTML version:** https://archtenet.dev/blog/mongodb-geosharding-hidden-cost
- **Published:** 2026-08-03
- **Authors:** Vladyslava Prykhodko
- **Tags:** mongodb, sharding, distributed-systems, architecture, mongodb-atlas, cost-optimization

Geo-sharding is one of the features that makes MongoDB attractive for a global product. You keep each user's data physically close to them, one logical collection still reads as a single dataset in Compass or Atlas, and the adoption cost looks tiny: add a shard key, and queries route to the right region.

The simplicity is real. So is the failure mode hiding inside it. This is a field story about how "easy" geo-sharding quietly doubled a database bill — and why that was never MongoDB's fault.

## Start with what the team actually saw

Before any explanation, here is what showed up on the dashboards:

- **CPU pinned near 95% across every shard at once**, not one hot node — all of them.
- **Atlas autoscaling flipping the cluster tier up and down twice a day** (M40 → M50 during peak, back down overnight).
- **Latency complaints and a rising ticket count** from users far from the database.
- **A support thread** suggesting the fix was to move the maintenance window.
- **A monthly bill that jumped by four figures** with no change in traffic or features.

Not one of these was the cause. Every item on that list is a *symptom* of a single decision made without understanding. The rest of this article is how those dots connect.

## What geo-sharding actually is

MongoDB's geo-distribution is built on **zone sharding**. You associate ranges of the shard key (for example, a `region` value) with zones, and you pin each zone to shards running in a specific cloud region. A user in India writes and reads data on shards in India; a user in Brazil stays on shards in Brazil. It is still one collection — your admin tools show everything together — but the bytes live where the people are.

Two things make it appealing: lower latency for a global user base, and data residency when you need it. Both are genuinely useful, and both are why it got chosen on the project in this story.

## Why it *looks* easy

The pitch is honestly this simple: add the location shard key to your queries, and MongoDB routes each read and write to a single region. That's a **targeted query** — the router (`mongos`) knows exactly which shard owns the data and goes straight there.

And it really is that easy — *as long as your queries carry the shard key*. That single condition is the whole game. Miss it, and the tool does something very different.

## Where it goes wrong

### 1. Unsharded collections don't live where you think

In a sharded cluster, every database has a **primary shard**, and every collection you *haven't* sharded lives entirely on it. The primary shard is chosen when the database is created — MongoDB picks whichever shard holds the least data at that moment. That can be any region, including the one farthest from most of your users.

So a collection that was never sharded ends up physically parked in a single region by accident. Now picture the service behind your app's main screen — the feed or activity read that fires on almost every user action. For users on the other side of the world, every one of those reads crosses an ocean, adding roughly 150–250 ms of round-trip latency to an interaction that's supposed to feel instant. (A checkout read or a per-request session lookup behaves the same way — anything latency-sensitive and read-on-every-action.) Users feel it immediately, and the tickets start.

Worth knowing: this default is no longer a hard restriction. Since MongoDB 8.0, moveCollection lets you place an unsharded collection on any shard in the cluster, independent of the database's primary shard — including a shard in the user's own region. Before 8.0, unsharded collections were pinned to the primary shard, and that constraint alone pushed a lot of teams into expensive tier upgrades to buy their way out. Keep this in your pocket; it becomes the cheaper fix later.

### 2. Sharding under pressure without fixing the queries

Under fire, the team does the reasonable-sounding thing: they shard the collection and add `region` as the shard key. Good instinct.

But they don't yet pass the shard key in their queries — no time, no shared understanding, and the field was never designed into the access patterns. Here's the trap closing:

```js
// No shard key → mongos broadcasts to EVERY shard (scatter-gather)
db.feed.find({ userId })

// Shard key present → mongos routes to ONE shard (targeted)
db.feed.find({ region, userId })   // region is the shard key
```

The moment the collection is sharded, any query *without* the shard key becomes a **scatter-gather**: `mongos` fans it out to all shards, waits for every one of them, and merges the result. Eight shards means eight queries per request instead of one. You can confirm exactly this with `explain()` — a healthy query hits a single shard; a broken one shows the query fanned out across all of them.

So the "fix" made things worse, not better. Before sharding, one shard was busy. After sharding without shard-key-aware queries, *every* shard is busy on *every* request.

### 3. Missing indexes turn scatter-gather into a CPU fire

Now add the other design gap: no indexes tuned to the real access patterns. Each broadcast query lands on each shard and does a collection scan. Multiply that by every request and every shard, and CPU climbs to ~95% **across the whole cluster at once** — which is exactly the first symptom on the dashboard.

### 4. Autoscaling turns a design bug into a recurring bill

MongoDB Atlas autoscaling reacts to sustained high CPU by bumping the cluster up a tier (M40 → M50). Overnight the load drops and it scales back down. Two tier changes a day, on schedule.

Each tier change reconfigures the replica set — a rolling restart and a primary election — so during the change, clients see brief connection blips and latency spikes. And this is where the team took a wrong turn: they read those blips as the *cause* of the slowness. A support engineer confirmed that tier changes and maintenance do cause short slowdowns, and suggested moving the maintenance window.

That advice isn't wrong. It's just aimed at a symptom. The autoscaling churn is real, but it only exists because the queries are scatter-gathering across an under-indexed cluster and pinning CPU. Move the maintenance window and you still pay for M50; fix the queries and the CPU spikes disappear, autoscaling stops firing, and the whole problem dissolves.

## Now the money

This is the part that makes the "easy feature" expensive. Rough illustrative Atlas list prices (AWS, dedicated tier — your real numbers depend on cloud provider, region, storage, IOPS, and node count):

| Tier | ~Hourly | ~Monthly (×730h) |
|------|---------|-----------------|
| M40  | ~$1.04  | ~$760           |
| M50  | ~$2.00  | ~$1,460         |

Those rates already cover a standard three-node replica set — that's how Atlas quotes a dedicated cluster (an M10 sits at ~$0.08/hr ≈ ~$57/month for the whole three-node set, not per node). So the M40 → M50 jump is roughly a **$700/month** delta **per shard**.

The multiplier that bites is the shard count. A geo-distributed deployment isn't one replica set — it's one replica set *per shard*, and autoscaling moves them together. Multiply that ~$700/month by the number of shards, and even a modest cluster clears four figures a month from this one jump — for a workload a correctly-routed, properly-indexed query would run on a fraction of the hardware.

The bill isn't a MongoDB price problem. It's the receipt for a query-routing decision nobody reviewed.

## The actual lesson

MongoDB didn't fail here. It did exactly what it was told: it broadcast queries that didn't name a shard, scanned collections that had no useful index, and scaled up hardware when CPU stayed hot. Every layer behaved correctly.

The failure was upstream, and it was organizational as much as technical:

- Sharding decisions were made without an architect in the room, on a POC that quietly became production.
- A shard key was added without the query changes that make a shard key mean anything.
- Indexes and schema were never designed for the real access patterns.
- The team debugged the symptoms it could see (autoscaling, maintenance windows) instead of the cause it couldn't.

The easiness of geo-sharding is precisely what makes it dangerous: it lets you ship the shard key without understanding query routing, and the invoice for that missing understanding arrives every month, quietly, as an autoscaling event.

## How to avoid the whole thing

- **Review before you shard.** Sharding is close to a one-way door; a 30-minute design conversation is cheaper than a tier upgrade.
- **Pick a shard key that matches how you actually query,** and make that key mandatory in those queries.
- **Verify targeting with `explain()`.** Confirm the hot queries hit one shard, not all of them — before you're under load.
- **Design indexes for the access patterns first,** not after CPU is at 95%.
- **Hunt for unsharded collections** sitting on a primary shard in the wrong region.
- **Treat autoscaling as a smoke detector, not a fix.** If the tier is flapping, go find the query.

Most "MongoDB is slow" or "MongoDB is expensive" incidents I've run into aren't MongoDB problems at all — they're query-routing and design problems wearing a MongoDB costume. They're also almost always far cheaper to prevent than to autoscale around.