# How to Monitor Data Distribution Across Shards Using getShardDistribution()

> Monitor data distribution across shards using getShardDistribution() in mongosh. Gain insights into document counts, data sizes, and storage metrics for balanced sharding.

- Repository: [Jin/mongodb-cluster-docker-compose](https://github.com/minhhungit/mongodb-cluster-docker-compose)
- Tags: how-to-guide
- Published: 2026-03-07

---

**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`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/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:

```bash
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:

```javascript
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:

```json
{
  "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:

```bash
#!/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`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/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`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/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`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/readme.md) | Documents the `getShardDistribution()` usage and provides sample JSON output for verification. |
| [`scripts/entrypoint-route.sh`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/scripts/entrypoint-route.sh) | Entry-point script for router containers ensuring `mongos` processes are ready to accept shell commands. |
| [`scripts/init-router.js`](https://github.com/minhhungit/mongodb-cluster-docker-compose/blob/main/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.