Skip to main content
Create your own

Understanding Sharded Cluster Balancing

Hello! Welcome back.

In our last lesson, we focused on the strategic decision of selecting a shard key, analyzing the critical trade-offs between data distribution and query performance. You learned how properties like cardinality, frequency, and monotonicity dictate the effectiveness of a key.

Today, we move from strategy to mechanics. We will explore what happens after you've chosen a shard key: how MongoDB partitions data, how it maintains balance, and what happens when the system is under load. Understanding these internal dynamics is crucial for operating a high-performance sharded cluster, diagnosing bottlenecks, and predicting system behavior—skills that are essential when managing platforms handling hundreds of thousands of transactions per second.

This lesson covers the lifecycle of data within a sharded cluster. We will analyze how data is organized into chunks, when and why those chunks are split, and the precise behavior of the balancer process that migrates them.

By the end of this ~60-minute lesson, you will be able to analyze chunk distribution, understand the mechanics of chunk splits, and describe the behavior of the balancer in a sharded MongoDB cluster.

1. Data Partitioning with Chunks

The fundamental unit of data distribution in a sharded cluster is the chunk. A chunk is simply a contiguous range of shard key values. The mongos router maintains metadata that maps these chunks to their respective shards, allowing it to direct queries efficiently.

The initial creation and subsequent size of these chunks are the first steps in understanding data distribution.

Data Partitioning with Chunks - Database Manual

To begin, let's ground our understanding in the official MongoDB documentation. This reading defines what chunks are, how they are created, and the trade-offs associated with their size.

Please read the introduction and the sections 'Initial Chunks' and 'Range Size'. As you read, focus on the different initial states of the cluster depending on whether the collection is empty and whether you're using ranged or hashed sharding. Also, consider the performance implications of the chunk size.

Based on that reading, let's analyze the key points:

  • Initial Chunk Creation: The starting state of your cluster's balance depends heavily on your sharding strategy.

    • Ranged Sharding (on an empty collection): The process starts with a single, empty chunk covering the entire shard key range. This means that initially, all writes will go to a single shard until that chunk grows and splits.
    • Hashed Sharding (on an empty collection): To avoid the initial single-shard problem, MongoDB pre-emptively creates and distributes empty chunks. By default, it creates two chunks per shard. This provides a much better initial write distribution.
    • Sharding a Populated Collection: The system creates a single large chunk covering all existing data, which the balancer then begins to split and migrate.
  • Chunk Size: The default chunk size is 128 MB. This is a configurable parameter, but changing it involves a significant trade-off:

    • Smaller Chunks: Lead to more frequent migrations. This allows for a very even data distribution but increases the overhead on the balancer and the network. Query routing via mongos can also be more expensive as it has to manage more chunk metadata.
    • Larger Chunks: Result in fewer, less frequent migrations. This is more efficient from a networking and routing overhead perspective but can lead to a "lumpier" data distribution, where imbalances between shards can be larger.

For many high-load systems, the efficiency gained from fewer migrations often outweighs the benefit of a perfectly balanced dataset, making the default or even a larger chunk size preferable.

2. The Balancer and Chunk Migration

An uneven distribution of chunks is the natural state of a growing cluster. The balancer is the background process responsible for correcting this. It runs on the primary of the config server replica set and its sole job is to migrate chunks between shards to maintain equilibrium.

Sharded Cluster Balancer - Database Manual

The next resource provides a detailed look at the balancer's internal workings, including the precise conditions that trigger a migration and the step-by-step procedure it follows.

Please read the introduction and the sections 'Balancer Internals', 'Range Migration Procedure', and 'Migration Thresholds'. The 'Range Migration Procedure' is particularly important; focus on understanding each step of the process.

Let's dissect the balancer's behavior.

The Migration Trigger

The balancer doesn't run constantly. It initiates a balancing round only when the data distribution for a collection crosses a specific migration threshold. A balancing round is triggered for a collection when the difference in the amount of data between the most-loaded and least-loaded shard exceeds a certain number of chunks.

The documentation notes that for a migration to occur, the difference must be at least three times the configured range size. With the default 128 MB chunk size, this means a shard must have at least 384 MB more data than another for the balancer to act. This threshold prevents the balancer from performing costly migrations for trivial imbalances.

The Migration Process

When the balancer decides to move a chunk from a source shard to a destination shard, it follows a well-defined, seven-step procedure:

  1. The balancer sends a moveRange command to the source shard.
  2. The source shard begins the process but continues to accept all reads and writes for the chunk being moved. It is still the authority for that data range.
  3. The destination shard receives the documents for the chunk and builds any necessary indexes.
  4. The destination shard begins requesting documents from the source shard.
  5. After receiving the last document, the destination shard enters a synchronization phase, applying any changes that occurred on the source shard during the copy process.
  6. Once fully synchronized, the source shard updates the cluster metadata on the config servers. This is the atomic cutover step where ownership of the chunk officially transfers to the destination shard.
  7. After the metadata is updated and any in-flight operations on the source shard complete, the source shard deletes its copy of the chunk's documents.

This process is designed for safety, but it has performance costs: network bandwidth for the data transfer, I/O load on both source and destination shards, and a brief period where operations on the collection may be paused on the source shard during the metadata update.

To optimize this, MongoDB performs the final deletion step asynchronously (as detailed in a003b, part 5). This allows the balancer to start the next chunk migration without waiting for the resource-intensive deletion from the previous migration to complete, which is critical when quickly rebalancing a cluster (e.g., after adding a new shard).

3. Chunk Splits and Distribution Pathologies

Now that we understand how chunks are created and moved, let's look at how they are split and what happens when the system encounters a distribution problem.

Modern Chunk Splitting Behavior

This is an area where MongoDB's behavior has evolved significantly.

Sharded Cluster Balancer - Database Manual

This final reading covers two critical operational aspects: the modern behavior of chunk splitting and the problematic scenario of 'jumbo' chunks.

Please read the section 'Chunk Size and Balancing'. Pay close attention to the description of how chunk splitting behavior differs in MongoDB 6.0 and later.

As you just read, the logic for splitting has changed:

  • Before MongoDB 6.0: Chunks were automatically split when they grew larger than the configured chunkSize.
  • MongoDB 6.0 and later: Chunks are only split when they are being migrated.

This is a subtle but important shift. It means that chunks on a shard can grow far beyond the chunkSize as long as the cluster remains balanced. The benefit is a reduction in the total number of chunks in the cluster. A smaller number of chunks means the routing table cached by mongos is smaller, which can improve query routing performance. Defragmentation and splitting now happen as a function of balancing, not just growth.

Pathological Case: Jumbo Chunks

The entire balancing mechanism relies on the ability to split and move chunks. What happens when a chunk cannot be split?

Data Partitioning with Chunks - Database Manual

Let's revisit our first resource to understand what happens when a chunk becomes indivisible.

Please read the section 'Indivisible/Jumbo Chunks'. This directly connects back to our previous lesson on shard key selection.

A jumbo chunk is a chunk that has grown beyond the configured chunk size but cannot be split. This typically occurs when a single shard key value has a very high frequency, and the total size of all documents with that single key value exceeds the chunk size.

Since a chunk is defined by a range of shard key values, a chunk containing only a single value cannot be split further. The balancer is then unable to move this oversized chunk, and it becomes a permanent hotspot. The shard containing the jumbo chunk will continue to attract all writes for that high-frequency key, while other shards remain underutilized.

This is the direct, mechanical consequence of choosing a shard key with poor frequency characteristics, as we discussed in the last lesson. It highlights why a theoretical understanding of shard key properties is so critical for practical operations. Modern MongoDB versions provide tools like reshardCollection and refineCollectionShardKey to recover from this state, but it's an expensive process that is best avoided through proper initial design.

Conclusion

In this lesson, we moved from the "what" of sharding to the "how." You've seen the machinery that underpins a MongoDB sharded cluster, from the initial partitioning of data into chunks to the automated process of balancing and the potential failure modes.

Key Takeaways:

  • Data is partitioned into chunks, which are ranges of shard key values. The initial distribution depends on whether you use ranged or hashed sharding.
  • The balancer is a background process that migrates chunks to maintain equilibrium, triggered only when the data imbalance between shards exceeds a specific threshold.
  • The chunk migration process is a robust, multi-step operation designed for consistency, involving data copying, synchronization, and an atomic metadata update. It has performance overhead that must be considered.
  • In modern MongoDB versions (6.0+), chunks are only split during migration, an optimization that reduces the total chunk count and can improve routing performance.
  • Jumbo chunks are a critical failure mode caused by a poor shard key choice (high frequency). They are indivisible and cannot be moved by the balancer, creating permanent hotspots.

This concludes our deep dive into MongoDB's sharding architecture. We've covered setting up a cluster, the strategic decision of choosing a shard key, and the dynamic mechanics of balancing.

Next Lesson Preview:

In our next module, we will shift our focus to a different paradigm of NoSQL databases: Apache Cassandra. We'll begin by exploring its fundamentally different, decentralized "shared-nothing" architecture, starting with core concepts like consistent hashing and virtual nodes that enable its impressive fault tolerance and linear scalability.

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

Sign up