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, and mongos routers as defined in docker-compose.yml.
  • You have enabled sharding on your target database and sharded the collection using sh.shardCollection().
  • You have a mongosh session connected to one of the mongos router 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 mongosh session connected to a mongos router container (e.g., router-01) in the Docker Compose cluster.
  • The raw field 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →