Type something to search...
MongoDB Sharding Explained: How to Scale Horizontally

MongoDB Sharding Explained: How to Scale Horizontally

Every database eventually runs into physics. A replica set can only hold as much data as its disks allow, only cache as much as its RAM allows, and only accept writes as fast as its single primary can process them. You can buy bigger machines for a while, and you should, because it's the simplest option. But at some point the next instance size is either absurdly expensive or doesn't exist.

Sharding is how MongoDB scales past a single machine. It splits a collection's data across multiple replica sets, called shards, and puts a routing layer in front so your application still talks to what looks like one database. Each shard holds a portion of the data and handles a portion of the traffic, so adding shards adds capacity for storage, memory, and writes all at once.

This guide covers the components of a sharded cluster, how data is divided and moved between shards, how queries are routed, how to shard a collection, and the operational realities that come with running a sharded cluster in production.

When You Actually Need Sharding

Sharding adds real complexity: more servers, more moving parts, and a permanent design decision (the shard key) that shapes your application's performance. Before you shard, make sure you've outgrown the alternatives.

Signs you might need it:

  • Data size. Your dataset is heading toward several terabytes, and backup, restore, and initial sync times on a single replica set are becoming unacceptable.
  • Working set. The indexes and frequently accessed documents no longer fit in RAM on the largest instance you're willing to pay for, and disk reads dominate.
  • Write throughput. The primary's CPU or disk is saturated by writes, and you've already optimized indexes and batching.
  • Data locality. You need to keep certain users' data in a specific region for latency or compliance reasons.

Signs you don't need it yet:

  • Slow queries caused by missing indexes. Sharding won't fix a collection scan; it'll just run it on several machines. Start with Using explain() to Analyze and Debug Slow MongoDB Queries.
  • Read load that could be handled by a larger instance or better caching.
  • A dataset of a few hundred gigabytes on hardware that's nowhere near its limits.

Vertical scaling first, sharding second. The best sharded clusters are built by teams who knew exactly which limit they were hitting.

The Components of a Sharded Cluster

A sharded cluster has three kinds of members:

  • Shards. Each shard is a full replica set that stores a subset of the data. Because each shard is a replica set, sharding and high availability work together: every piece of data still has multiple copies.
  • mongos routers. These are lightweight, stateless processes that your application connects to. A mongos looks at each query, figures out which shards hold the relevant data, forwards the query, and merges the results. You typically run several, often one per application server or behind a load balancer.
  • Config servers. A replica set (called the CSRS) that stores the cluster's metadata: which collections are sharded, their shard keys, and which ranges of data live on which shard. In MongoDB 8.0, the config server can also act as a regular data-bearing shard (a config shard), which reduces the minimum footprint for smaller clusters.

Your application's connection string points at the mongos routers, not the shards:

import { MongoClient } from "mongodb";

const client = new MongoClient(
  "mongodb://mongos-1:27017,mongos-2:27017/?retryWrites=true",
);

From the application's point of view, very little changes. Queries, updates, aggregations, transactions, and change streams all work through mongos. The differences show up in performance, which depends heavily on the shard key.

How Data Is Divided: Shard Keys and Chunks

When you shard a collection, you choose a shard key: one or more fields present in every document. MongoDB divides the range of possible shard key values into contiguous ranges called chunks, and each chunk lives on exactly one shard.

Say you shard orders on customerId. The config servers might record something like:

// Simplified view of config.chunks for shop.orders
{ min: { customerId: MinKey },  max: { customerId: "c_2500" }, shard: "shard-a" }
{ min: { customerId: "c_2500" }, max: { customerId: "c_5100" }, shard: "shard-b" }
{ min: { customerId: "c_5100" }, max: { customerId: MaxKey },  shard: "shard-c" }

A document with customerId: "c_3001" belongs to the second chunk, so it lives on shard-b. Every order for that customer lives on the same shard.

Ranged vs. Hashed Sharding

There are two main ways to map shard key values to chunks:

StrategyHow it worksStrengthsWeaknesses
RangedChunks are ranges of the actual key valuesEfficient range queries on the keyMonotonic keys send all inserts to one shard
HashedChunks are ranges of a hash of the key valueEven write distribution, even for ObjectIdsRange queries on the key hit every shard
// Ranged: good when queries filter by ranges of customerId and values are well spread
sh.shardCollection("shop.orders", { customerId: 1 });

// Hashed: good for spreading writes on a key like an ObjectId or timestamp
sh.shardCollection("telemetry.events", { deviceId: "hashed" });

// Compound: group by tenant, then spread within each tenant
sh.shardCollection("saas.documents", { tenantId: 1, _id: "hashed" });

Choosing the key well is the single most important sharding decision, and it gets a full treatment in How to Choose a Good Shard Key in MongoDB. The short version: you want high cardinality, values that are spread evenly, writes that don't all hit the same chunk, and a key that appears in most of your queries.

Sharding a Collection Step by Step

Assuming your cluster is running with shards added, connect to a mongos and shard a collection:

// Connect with mongosh to a mongos, not a shard

// Optional in recent versions: shardCollection enables sharding for the database
sh.enableSharding("shop");

// The shard key needs a supporting index; on an empty collection it's created for you
db.getSiblingDB("shop").orders.createIndex({ customerId: 1, orderDate: 1 });

sh.shardCollection("shop.orders", { customerId: 1, orderDate: 1 });

Then check the result:

sh.status();

The output lists shards, the balancer state, and each sharded collection with its key and chunk distribution:

shards
[
  { _id: 'shard-a', host: 'shard-a/sa-1:27018,sa-2:27018,sa-3:27018', state: 1 },
  { _id: 'shard-b', host: 'shard-b/sb-1:27018,sb-2:27018,sb-3:27018', state: 1 },
  { _id: 'shard-c', host: 'shard-c/sc-1:27018,sc-2:27018,sc-3:27018', state: 1 }
]
---
balancer
{ 'Currently enabled': 'yes', 'Currently running': 'no' }
---
databases
[
  {
    database: { _id: 'shop', primary: 'shard-a', ... },
    collections: {
      'shop.orders': {
        shardKey: { customerId: 1, orderDate: 1 },
        unique: false,
        balancing: true,
        chunkMetadata: [
          { shard: 'shard-a', nChunks: 4 },
          { shard: 'shard-b', nChunks: 3 },
          { shard: 'shard-c', nChunks: 3 }
        ]
      }
    }
  }
]

For a per-shard view of actual data size and document counts, use:

db.orders.getShardDistribution();
Shard shard-a at shard-a/sa-1:27018,...
{ data: '41.2GiB', docs: 51203344, chunks: 4, 'estimated data per chunk': '10.3GiB' }
Shard shard-b at shard-b/sb-1:27018,...
{ data: '39.8GiB', docs: 49551019, chunks: 3, 'estimated data per chunk': '13.2GiB' }
...
Totals
{ data: '120.6GiB', docs: 150012877, chunks: 10, 'Shard shard-a': [ '34.16 % data', ... ] }

Collections you don't shard still work. Each database has a primary shard, and unsharded collections live entirely on it. That's fine for small collections, but keep an eye on it: a database with many large unsharded collections can overload its primary shard. MongoDB 8.0 added moveCollection so you can place an unsharded collection on a specific shard.

How Queries Are Routed

The mongos uses the chunk map from the config servers to decide where to send each operation. There are two kinds of routing:

Targeted queries include the shard key (or a prefix of it), so mongos can send them to only the shards that own the relevant chunks:

db.orders.find({ customerId: "c_3001" }); // one shard
db.orders.find({ customerId: { $in: ["c_10", "c_9000"] } }); // at most two shards

Scatter-gather queries don't include the shard key, so mongos must ask every shard and merge the results:

db.orders.find({ status: "pending" }); // every shard

You can tell which one you got from explain(). A targeted query shows a SINGLE_SHARD stage; a broadcast query shows SHARD_MERGE with every shard listed:

db.orders.find({ customerId: "c_3001" }).explain().queryPlanner.winningPlan
  .stage;
// 'SINGLE_SHARD'

db.orders.find({ status: "pending" }).explain().queryPlanner.winningPlan.stage;
// 'SHARD_MERGE'

Scatter-gather isn't automatically bad. Analytics queries that need to scan a lot of data benefit from running in parallel on every shard. But for high-volume, latency-sensitive operations (loading a user's profile, a customer's orders, a cart), you want targeted queries. The latency of a scatter-gather query is set by the slowest shard, and it consumes resources on every shard for every request.

Writes follow the same logic. An insertOne always goes to exactly one shard, determined by the document's shard key. updateOne and deleteOne are most efficient when the filter includes the shard key; recent versions can handle single-document updates without it, but at extra cost.

The Balancer and Chunk Migrations

As data grows, some shards can end up with more than their share. The balancer, a background process running on the config server primary, watches the distribution of each sharded collection and moves chunks from overloaded shards to underloaded ones.

In recent versions, balancing is based on the amount of data per shard rather than the number of chunks, and the default chunk size is 128 MB. When a shard holds noticeably more data for a collection than the others, the balancer migrates ranges until they're roughly even.

A chunk migration works like this:

  1. The destination shard copies the documents in the range from the source shard.
  2. Writes to the range that happen during the copy are captured and replayed on the destination.
  3. In a brief critical section, the source stops accepting writes to the range, the final changes are transferred, and the config servers update the chunk's owner.
  4. The source shard deletes its copy of the range asynchronously.

Migrations are online, but they cost I/O and network bandwidth on both shards. On busy clusters, you can restrict the balancer to a quiet time window:

use config
db.settings.updateOne(
  { _id: "balancer" },
  { $set: { activeWindow: { start: "01:00", stop: "05:00" } } },
  { upsert: true }
)

sh.getBalancerState()     // true
sh.isBalancerRunning()    // { mode: 'full', inBalancerRound: false, ... }

Balancing also only works if chunks can be split. If a huge number of documents share the same shard key value, they must stay in one chunk, which can't be split or moved easily. That's a jumbo chunk, and it's almost always a sign of a low-cardinality shard key.

Zone Sharding for Data Locality

Zones let you pin ranges of shard key values to specific shards. The classic use is geographic: keep European customers' data on shards in a European region.

sh.addShardToZone("shard-eu-1", "EU");
sh.addShardToZone("shard-us-1", "US");

// Shard key: { region: 1, customerId: 1 }
sh.updateZoneKeyRange(
  "shop.customers",
  { region: "eu", customerId: MinKey },
  { region: "eu", customerId: MaxKey },
  "EU",
);
sh.updateZoneKeyRange(
  "shop.customers",
  { region: "us", customerId: MinKey },
  { region: "us", customerId: MaxKey },
  "US",
);

The balancer then ensures chunks in each range live only on shards in the matching zone. For this to work, the zone field must be the leading part of the shard key. Zones are also handy for tiered storage (recent data on fast shards, old data on cheaper ones) and for isolating a large tenant on dedicated hardware.

Changing Your Mind: Resharding

Early versions of MongoDB made the shard key permanent. That's no longer true:

  • refineCollectionShardKey adds suffix fields to an existing key, which is useful when the original key has too little cardinality: { customerId: 1 } becomes { customerId: 1, orderId: 1 }.
  • reshardCollection rewrites a collection under a completely new shard key while it stays online. It needs spare disk space and I/O capacity, and it takes time on large collections, but it's a supported path out of a bad choice.
db.adminCommand({
  reshardCollection: "shop.orders",
  key: { customerId: "hashed" },
});

MongoDB 8.0 also added unshardCollection to turn a sharded collection back into an unsharded one. These tools make a poor shard key recoverable, but they're expensive operations, not something to plan around.

Running Sharded Clusters in Production

Size each shard like a production replica set. Every shard needs three data-bearing members across failure domains. A cluster with three shards is at least nine data nodes, plus the config server replica set and several mongos instances.

Backups need cluster-wide consistency. Running mongodump against each shard separately doesn't produce a consistent point in time. Use a backup method designed for sharded clusters, like Atlas cloud backups or Ops Manager, or filesystem snapshots coordinated with the balancer stopped.

Transactions spanning shards cost more. A transaction that touches documents on several shards requires a two-phase commit. Design your shard key so related documents, like an order and its line items, live together.

On Atlas, sharding is a cluster setting. Dedicated clusters (M30 and above) can be deployed as sharded clusters, and you can add shards later from the UI or API. Atlas manages config servers and mongos for you, which removes much of the operational burden.

Common Pitfalls

Sharding too early. Sharding a 50 GB collection on a lightly loaded cluster adds latency and complexity with no benefit. Scale up and optimize first.

Choosing a monotonically increasing key. A ranged shard key on a timestamp or ObjectId sends every insert to the chunk with the highest values, so one shard takes all writes while the others sit idle. Use a hashed key or a compound key with a better-distributed leading field.

Picking a key your queries don't use. If most queries filter on customerId but the collection is sharded on orderId, nearly every read becomes scatter-gather. The key should match your most frequent query patterns.

Low cardinality. Sharding on country or status gives you a handful of possible values, which means a handful of chunks that can never be split. The largest one becomes a jumbo chunk and a permanent hot spot.

Connecting directly to shards. Applications must go through mongos. Writing directly to a shard bypasses routing and can put documents on the wrong shard, where queries through mongos won't find them.

Forgetting the unsharded collections. Large unsharded collections all live on the database's primary shard and can make it much busier than the rest of the cluster.

Conclusion

Sharding spreads a collection across multiple replica sets so that storage, memory, and write capacity grow with the number of shards. mongos routers direct each operation using a chunk map held by the config servers, the balancer keeps data evenly distributed by migrating ranges in the background, and zones let you control where specific data lives. Most of the performance of a sharded cluster comes down to the shard key: it decides whether writes spread evenly and whether queries stay targeted.

Before your next scaling discussion, collect three numbers for your largest collection: its data size, its index size compared to available RAM, and the primary's peak write rate. If none of them is close to the limit of your current hardware, you don't need sharding yet. If one is, start designing your shard key now, well before you're forced to.

Tags :
Share :

Related Posts

A Complete Guide to MongoDB Query Operators

A Complete Guide to MongoDB Query Operators

Your first MongoDB queries are usually simple equality filters: find the user with this email, find orders with this status. That covers a surprising

Continue Reading
Async MongoDB in Python with Motor and FastAPI

Async MongoDB in Python with Motor and FastAPI

FastAPI runs your endpoints on an event loop. That's what lets a single worker juggle hundreds of concurrent requests: while one request waits on the

Continue Reading
Atlas Online Archive: Tiering Cold Data to Cut Costs

Atlas Online Archive: Tiering Cold Data to Cut Costs

Look at almost any production database and you'll find the same shape. A small slice of recent data gets nearly all the reads and writes: this week's

Continue Reading