How to Monitor Data Distribution Across Shards Using getShardDistribution()
Use db.<collection>.getShardDistribution() from a mongosh session connected to a mongos router to retrieve per-shard document counts, data sizes, and storage metrics that reveal how evenly your data is split across the cluster.
The mongodb-cluster-docker-compose repository provides a complete Docker Compose environment for running a sharded MongoDB cluster with multiple replica sets. Monitoring how data distributes across these shards is critical for performance tuning and capacity planning, and the getShardDistribution() shell helper offers a direct way to inspect this balance without external monitoring tools.
What is getShardDistribution()?
getShardDistribution() is a MongoDB shell helper method available on collection objects that reports distribution statistics for a sharded collection. When invoked, it queries the cluster metadata and aggregates storage statistics from each shard that owns chunks of the target collection.
The method returns a document containing a raw field with sub-documents for each shard replica set. Each shard entry includes the following metrics:
| Metric | Description |
|---|---|
db |
Database name hosting the collection |
collections |
Number of collections on the shard (typically 1 for the target) |
objects |
Document count on that shard (may be 0) |
dataSize |
Uncompressed data size in bytes |
storageSize |
Allocated storage size on disk |
indexes |
Number of indexes on the shard |
indexSize |
Total size of indexes on the shard |
fsUsedSize / fsTotalSize |
Filesystem utilization on the shard node |
These values enable you to verify that your sharding key is effectively distributing data and to identify hotspots where a single shard holds disproportionate data.
Prerequisites for Monitoring Shard Distribution
To use getShardDistribution() in the context of the mongodb-cluster-docker-compose setup, ensure the following:
- The Docker Compose stack is running with the three replica-set shards (
rs-shard-01,rs-shard-02,rs-shard-03), config servers, andmongosrouters as defined indocker-compose.yml. - You have enabled sharding on your target database and sharded the collection using
sh.shardCollection(). - You have a
mongoshsession connected to one of themongosrouter containers (e.g.,router-01), as only routers possess the cluster metadata required to map chunks to shards.
How to Run getShardDistribution() in the Docker Cluster
Accessing the Router Container
Open a shell session on the first router container:
docker exec -it router-01 mongosh --port 27017
This connects you to the mongos instance that routes queries to the underlying shards.
Executing the Command
Switch to your sharded database and invoke the helper on the target collection:
use MyDatabase
db.MyCollection.getShardDistribution()
Replace MyDatabase and MyCollection with your actual database and collection names.
Interpreting the Output
The command returns a JSON document similar to this excerpt from the repository documentation:
{
"raw": {
"rs-shard-01/shard01-a:27017,shard01-b:27017,shard01-c:27017": {
"db": "MyDatabase",
"collections": 1,
"objects": 0,
"dataSize": 0,
"storageSize": 4096,
"indexes": 1,
"indexSize": 4096,
"fsUsedSize": 123456789,
"fsTotalSize": 1073741824
},
"rs-shard-02/shard02-a:27017,shard02-b:27017,shard02-c:27017": {
"objects": 15000,
"dataSize": 15728640,
...
},
"rs-shard-03/shard03-a:27017,shard03-b:27017,shard03-c:27017": {
"objects": 15000,
"dataSize": 15728640,
...
}
},
"objects": 30000,
"dataSize": 31457280,
"storageSize": 36864
}
In this example, rs-shard-01 holds zero documents while the other two shards hold 15,000 each, indicating a potential issue with the sharding key or chunk distribution that requires running the balancer or re-evaluating the shard key strategy.
Automating Shard Distribution Health Checks
For production monitoring, wrap getShardDistribution() in a shell script that alerts when data imbalance exceeds a threshold. The following script runs inside the router-01 container and exits with an error if any single shard holds more than 60% of total documents:
#!/bin/bash
THRESHOLD=0.6
OUTPUT=$(docker exec router-01 mongosh --quiet --port 27017 <<'EOF'
use MyDatabase
printjson(db.MyCollection.getShardDistribution())
EOF
)
TOTAL=$(echo "$OUTPUT" | jq '.objects')
MAX_SHARD=$(echo "$OUTPUT" | jq '[.raw[].objects] | max')
RATIO=$(awk "BEGIN {print $MAX_SHARD/$TOTAL}")
if (( $(echo "$RATIO > $THRESHOLD" | bc -l) )); then
echo "⚠️ Data imbalance detected: $MAX_SHARD of $TOTAL objects on one shard"
exit 1
else
echo "✅ Shard distribution is balanced"
fi
Save this as check-distribution.sh and schedule it via cron or your CI/CD pipeline to detect hotspots before they impact performance.
Key Files in the Repository
Understanding the cluster architecture helps interpret the distribution output. These files from minhhungit/mongodb-cluster-docker-compose define the environment:
| File | Purpose |
|---|---|
docker-compose.yml |
Defines the three replica-set shards (rs-shard-01, rs-shard-02, rs-shard-03), config servers, and mongos routers. |
readme.md |
Documents the getShardDistribution() usage and provides sample JSON output for verification. |
scripts/entrypoint-route.sh |
Entry-point script for router containers ensuring mongos processes are ready to accept shell commands. |
scripts/init-router.js |
JavaScript initialization for configuring sharding after router startup. |
Summary
getShardDistribution()is a MongoDB shell helper that reports per-shard document counts, data sizes, and storage metrics for a sharded collection.- Run the command from a
mongoshsession connected to amongosrouter container (e.g.,router-01) in the Docker Compose cluster. - The
rawfield in the output contains individual shard statistics that reveal whether your sharding key is creating balanced chunks or causing hotspots. - Automate health checks by wrapping the helper in shell scripts that alert when any single shard exceeds a defined percentage of total data.
Frequently Asked Questions
When should I run getShardDistribution()?
Run getShardDistribution() immediately after bulk imports to verify chunks migrated correctly, during troubleshooting when query performance degrades on specific shards, and periodically as part of automated health checks to ensure the balancer is maintaining an even distribution.
Can I run getShardDistribution() on a config server?
No. The helper must be executed from a mongos router because only routers maintain the cluster metadata that maps chunks to shards. Running it directly on a shard or config server will result in an error or incomplete data.
What does it mean if one shard has significantly more objects than others?
A disproportionate object count indicates a hotspot, where your sharding key is causing uneven chunk distribution. This can occur with monotonically increasing keys (like timestamps) or low-cardinality fields. You may need to adjust the sharding key or manually split and move chunks using sh.moveChunk().
How do I fix data imbalance detected by getShardDistribution()?
First, ensure the balancer is enabled and running via sh.getBalancerState(). If the balancer is active but chunks remain uneven, manually trigger a chunk migration with sh.moveChunk("<database>.<collection>", { <shardKey>: <value> }, "<targetShard>"). For severe cases, consider resharding the collection with a better distribution key.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →