Skip to main content
Create your own

Building a Sharded MongoDB Cluster

Hello! Welcome to the sixth module of our course on high-load distributed systems.

In our last two lessons, we focused on MongoDB's consistency model and the practical setup of replica sets to achieve high availability. You learned how a replica set uses an election protocol to automatically failover, ensuring the system remains operational even when a node goes down.

Today, we address a different challenge: horizontal scalability. While a replica set provides redundancy, it is still limited by the storage capacity and write throughput of a single primary node. To overcome this, we use sharding. This lesson will guide you through setting up the fundamental architecture of a MongoDB sharded cluster. We will assemble the core components: config servers, shard replica sets, and the mongos query router.

By the end of this ~60-minute lesson, you will be able to set up a complete, multi-component MongoDB sharded cluster architecture from scratch.

1. The Architecture of a Sharded Cluster

A sharded cluster distributes data across multiple replica sets, called shards. This allows the database to scale horizontally, handling massive datasets and high throughput workloads that would overwhelm a single replica set.

The architecture consists of three primary components:

  1. Shard Replica Sets: These are the workhorses of the cluster. Each shard is a full replica set that stores a subset of the total data. By using replica sets for each shard, you get both horizontal scaling (from sharding) and high availability (from replication) combined.
  2. Config Server Replica Set (CSRS): This is the cluster's metadata authority. It stores the mapping of data ranges (called "chunks") to their corresponding shards. Since this metadata is critical for the cluster's operation, the config servers are themselves deployed as a replica set to ensure they are highly available. If the CSRS is down, the cluster can still serve reads and writes for established connections, but no new connections can be made, and no chunk migrations can occur.
  3. Query Routers (mongos): These are lightweight, stateless proxies that act as the interface for your application. The application connects to a mongos instance, not directly to the shards. The mongos process consults the config servers to determine which shard(s) hold the requested data, routes the query accordingly, and aggregates the results if the query spans multiple shards. You typically run multiple mongos instances for load balancing and to avoid a single point of failure at the routing layer.

This diagram illustrates how these components interact:

Architecture of a MongoDB sharded cluster. Client applications connect to the `mongos` routers, which use metadata from the config servers to direct operations to the appropriate shard replica sets.

To solidify your understanding of these components, please read the following introductory section from a Percona blog post.

A Tutorial on MongoDB Sharding Best Practices & When ...

This section provides a concise overview of the key components in a MongoDB sharded cluster, defining the role of each piece before we dive into the setup.

Please read the section titled 'What is sharding in MongoDB?' up to the image. Focus on the definitions of Shard Servers, Config Servers, Query Routers (mongos), Shard Key, Chunk, and Balancer.

2. Practical Guide: Setting Up a Sharded Cluster

Now, let's build a sharded cluster on your local machine. We will simulate a production topology by running each process on a different port. The overall process is:

  1. Deploy the Config Server Replica Set.
  2. Deploy the Shard Replica Sets.
  3. Start the mongos Query Router.
  4. Add the shards to the cluster.

For this entire process, we will follow the official MongoDB documentation. It's detailed and serves as an excellent reference for production deployments.

Deploy a Self-Managed Sharded Cluster - Database Manual

This official MongoDB tutorial is the primary guide for our hands-on setup. We will walk through its steps together, but it's essential to have it open as your main reference.

Keep this document open for reference as we proceed through the steps below. You don't need to read it all at once; we will refer to specific parts as we go.

Step 1: Prepare the Environment

First, create the data directories for all the mongod instances we'll be starting. We'll set up one config server replica set (3 members) and two shard replica sets (3 members each).

Open your terminal and run:

# Directories for the Config Server Replica Set (CSRS)
mkdir -p /data/csrs1 /data/csrs2 /data/csrs3

# Directories for Shard 1 Replica Set (rs-shard-A)
mkdir -p /data/shardA1 /data/shardA2 /data/shardA3

# Directories for Shard 2 Replica Set (rs-shard-B)
mkdir -p /data/shardB1 /data/shardB2 /data/shardB3

Step 2: Deploy the Config Server Replica Set (CSRS)

The CSRS stores the cluster's configuration. It's a mongod process started with the --configsvr flag.

Open three separate terminal tabs and start one mongod process in each.

Terminal 1 (CSRS 1):

mongod --configsvr --replSet csrs --dbpath /data/csrs1 --port 27019

Terminal 2 (CSRS 2):

mongod --configsvr --replSet csrs --dbpath /data/csrs2 --port 27020

Terminal 3 (CSRS 3):

mongod --configsvr --replSet csrs --dbpath /data/csrs3 --port 27021

Now, connect to one of these instances to initiate the replica set.

mongosh --port 27019

Inside mongosh, define the configuration and initiate. Note the configsvr: true property, which is mandatory.

config = {
    _id: "csrs",
    configsvr: true,
    members: [
        { _id: 0, host: "localhost:27019" },
        { _id: 1, host: "localhost:27020" },
        { _id: 2, host: "localhost:27021" }
    ]
};
rs.initiate(config);

Your CSRS is now running.

Step 3: Deploy the Shard Replica Sets

Next, we'll deploy two shard replica sets. These are standard replica sets, but the mongod instances are started with the --shardsvr flag.

For Shard A (rs-shard-A):
Open three new terminal tabs.

Terminal 4 (Shard A1):

mongod --shardsvr --replSet rs-shard-A --dbpath /data/shardA1 --port 27022

Terminal 5 (Shard A2):

mongod --shardsvr --replSet rs-shard-A --dbpath /data/shardA2 --port 27023

Terminal 6 (Shard A3):

mongod --shardsvr --replSet rs-shard-A --dbpath /data/shardA3 --port 27024

Connect to one member and initiate it:

mongosh --port 27022
config = {
    _id: "rs-shard-A",
    members: [
        { _id: 0, host: "localhost:27022" },
        { _id: 1, host: "localhost:27023" },
        { _id: 2, host: "localhost:27024" }
    ]
};
rs.initiate(config);

For Shard B (rs-shard-B):
Repeat the process for the second shard in three more terminal tabs.

Terminal 7 (Shard B1):

mongod --shardsvr --replSet rs-shard-B --dbpath /data/shardB1 --port 27025

Terminal 8 (Shard B2):

mongod --shardsvr --replSet rs-shard-B --dbpath /data/shardB2 --port 27026

Terminal 9 (Shard B3):

mongod --shardsvr --replSet rs-shard-B --dbpath /data/shardB3 --port 27027

Connect and initiate:

mongosh --port 27025
config = {
    _id: "rs-shard-B",
    members: [
        { _id: 0, host: "localhost:27025" },
        { _id: 1, host: "localhost:27026" },
        { _id: 2, host: "localhost:27027" }
    ]
};
rs.initiate(config);

At this point, you have three independent replica sets running. Now we need to tie them together.

Step 4: Start the mongos Query Router

The mongos process is the gateway to the cluster. It needs to know the location of the config servers.

Open a new terminal tab and start the mongos instance. Note that mongos is a different binary from mongod.

Terminal 10 (mongos):

mongos --configdb csrs/localhost:27019,localhost:27020,localhost:27021 --bind_ip_all --port 27017

The --configdb argument tells mongos where to find the CSRS. The format is <replica_set_name>/<host1:port1>,<host2:port2>,.... We run it on the default MongoDB port 27017 for convenience.

Step 5: Add Shards to the Cluster

All components are running, but the mongos router doesn't know about the shards yet. We need to add them.

Connect to the mongos instance using mongosh. This is the entry point for all administrative commands and application traffic from now on.

mongosh --port 27017

Your prompt should now look like mongos>.

Use the sh.addShard() command to add each shard replica set to the cluster.

// Add the first shard
sh.addShard("rs-shard-A/localhost:27022,localhost:27023,localhost:27024");

// Add the second shard
sh.addShard("rs-shard-B/localhost:27025,localhost:27026,localhost:27027");

Step 6: Verify the Cluster Status

Your sharded cluster is now fully assembled. You can verify its status using sh.status().

mongos> sh.status()

The output will be verbose, but look for these key sections:

  • sharding version: Information about the cluster's metadata.
  • shards: A list of the shards added to the cluster, confirming they are recognized.
  • databases: Will be empty for now, but will later show which databases are sharded and the distribution of their data.

You have successfully built the infrastructure for a MongoDB sharded cluster.

3. Enabling Sharding for a Collection

The final step, which we will only touch on briefly today, is to enable sharding for a specific database and collection. This involves choosing a shard key.

Deploy a Self-Managed Sharded Cluster - Database Manual

This final part of the MongoDB documentation shows the commands to enable sharding on a database and a collection. It also introduces the two main sharding strategies.

Read the section 'Shard a Collection'. Pay attention to the sh.shardCollection() command and the distinction between Hashed and Range-based sharding.

To shard a collection, you would run the following commands in the mongos shell:

  1. Enable sharding on a database:

    sh.enableSharding("myAppDB")
    
  2. Shard a collection within that database:

    // This requires an index on the shard key field if the collection has data
    db.getSiblingDB("myAppDB").myCollection.createIndex({ userId: 1 })
    
    // Shard the collection based on the 'userId' field
    sh.shardCollection("myAppDB.myCollection", { userId: 1 })
    

The choice of the shard key (userId in this example) is arguably the most critical decision in designing a sharded cluster. It determines how data is distributed and directly impacts query performance, write scalability, and the efficiency of the cluster's balancer.

Conclusion

In this lesson, we constructed a complete MongoDB sharded cluster from its constituent parts. You now have a practical understanding of how the different components are deployed and connected.

Key Takeaways:

  • A sharded cluster is composed of three main components: Shard Replica Sets (to store data), a Config Server Replica Set (to store metadata), and one or more mongos Query Routers (to route client requests).
  • Each component plays a distinct role, and all are designed for high availability using replica sets.
  • The setup process is sequential: deploy config servers, deploy shards, start mongos, and finally add the shards to the cluster via sh.addShard().
  • All application traffic and administrative tasks for the cluster should be directed through a mongos instance.

Next Lesson Preview:

We have built the house, but now we need to decide how to organize the rooms. In our next lesson, "Select an appropriate shard key for a given data model and query pattern," we will dive deep into the art and science of choosing an effective shard key. We'll analyze the trade-offs between different strategies, such as ranged vs. hashed sharding, and explore how the choice of key impacts the performance and scalability of your entire system.

Can't find a good explanation? Sign up and we'll make it for you

Sign up