0Pricing
MongoDB Academy · Lektion

Skalierungsplan: Vom Replikatsatz zum Sharded Cluster

Lernende erstellen einen Kapazitätsplan und wählen einen Shard Key, der die Verteilung der Lese- und Schreibzugriffe der Anwendung unterstützt, ohne Hotspots zu erzeugen.

Skalierungsplan: Vom Replikatsatz zum Sharded Cluster ist eine kostenlose MongoDB Academy-Lektion auf CoddyKit. Dies ist Lektion 3 von 4. Du kannst die komplette Lektion unten kostenlos lesen – dann übst du sie direkt im Browser mit einem integrierten Code-Editor und einem KI-Tutor rund um die Uhr. Sie ist Teil des MongoDB Academy-Lernpfads, und dein Fortschritt wird über Web und CoddyKit-App synchronisiert. Der MongoDB Academy-Kurs umfasst insgesamt 4 Lektionen.

Teile dieser Lektion wurden noch nicht übersetzt und werden auf Englisch angezeigt.

When Do You Need to Scale?

Most applications start with a single MongoDB replica set and never need to shard. Sharding adds significant complexity and should be the last resort, not the first choice. Consider sharding when: your data volume exceeds what a single replica set can store affordably; write throughput exceeds what a single primary can handle; or specific collections are too large to index efficiently in memory. Before sharding, always try vertical scaling (larger instances) and read scaling (distribute reads to secondaries).

Phase 1: Single Replica Set

The standard starting point for any MongoDB deployment is a 3-member replica set: one primary, two secondaries. This provides high availability (automatic failover if the primary fails), data durability (writes replicated to multiple members), and read scaling (send reads to secondaries for reporting workloads). For our e-commerce capstone, a 3-member M30 Atlas cluster handles millions of daily orders comfortably. Start here and measure before considering sharding.

// Capacity metrics to monitor on a single replica set
// (via Atlas Metrics or db.serverStatus())
const metricsToWatch = [
  'connections.current',         // approaching maxIncomingConnections?
  'opcounters.insert',           // write ops/sec approaching primary limit?
  'mem.resident',                // working set fitting in RAM?
  'wiredTiger.cache.bytesCurrentlyInCache', // cache utilisation
  'replicationLag'               // secondaries keeping up?
]

Phase 2: Read Scaling With Secondary Reads

Before sharding, scale reads by routing analytics and reporting queries to secondaries using readPreference: 'secondary'. This offloads read pressure from the primary without adding operational complexity. In Atlas, Analytics Nodes are dedicated secondaries (never elected primary) that absorb heavy aggregation workloads without impacting primary performance. This approach works well until write throughput itself becomes the bottleneck.

// Route heavy analytics to secondary nodes
const { MongoClient } = require('mongodb')

const client = new MongoClient(process.env.ATLAS_URI, {
  readPreference: 'secondary'   // global default for this client
})

// Or per-operation
const result = await db.collection('orders').aggregate(
  [ /* heavy reporting pipeline */ ],
  { readPreference: 'secondary' }  // does not compete with primary writes
)

Phase 3: When to Shard

Shard when you hit a wall that vertical scaling cannot solve: the primary's write throughput is saturated even with the largest available instance; a collection's working set (data + indexes actively in use) does not fit in the cluster's RAM even on the largest tier; or a specific collection is too large to store on one cluster's disk. In practice, most applications hit RAM limits before write throughput limits — monitor the WiredTiger cache utilisation and working set size monthly.

// Indicator: working set exceeding cache
// db.serverStatus().wiredTiger.cache
const cache = db.serverStatus().wiredTiger.cache
const cacheHitRatio = 1 - (cache['pages read into cache'] / cache['pages requested from the cache'])
console.log('Cache hit ratio:', (cacheHitRatio * 100).toFixed(1) + '%')
// Below 95%: working set is not fitting in cache — time to scale

Choosing Which Collection to Shard

Only shard the collection(s) that are causing the bottleneck. In our e-commerce platform, the orders collection will grow fastest and generate the most write traffic. The products collection may be large but receives mostly reads that can be served from secondaries. Sharding orders while keeping products unsharded (on every shard as a broadcast collection) is a common, practical approach.

// Enable sharding on the database
sh.enableSharding('ecommerce')

// Shard the orders collection
sh.shardCollection('ecommerce.orders', { userId: 'hashed' })

// Verify shard distribution
sh.status()
db.orders.getShardDistribution()

Shard Key Selection for Orders

The shard key for orders must distribute writes evenly across shards and support the most common query patterns. Using userId as a hashed shard key distributes writes uniformly because user IDs are high-cardinality and random. The downside: queries scoped to a single user scatter across all shards. A ranged shard key on { userId: 1, _id: 1 } keeps one user's orders on the same shard (faster user-specific queries) but risks hot shards if a few users generate most activity.

// Option 1: Hashed shard key — uniform write distribution
sh.shardCollection('ecommerce.orders', { userId: 'hashed' })
// Pros: even distribution
// Cons: user order history queries scatter across all shards

// Option 2: Compound ranged shard key — user orders co-located
sh.shardCollection('ecommerce.orders', { userId: 1, _id: 1 })
// Pros: all orders for a user are on one shard — fast history queries
// Cons: may create hot shards if a few users dominate traffic

Zone Sharding for Data Residency

If the e-commerce platform serves multiple regions with data residency requirements (EU data must stay in Europe), use zone sharding to pin documents to specific shards based on shard key ranges. Create zones for each region, assign shards to zones, and define which shard key ranges map to each zone. Data for EU users stays on EU-region shards, satisfying GDPR without maintaining separate clusters.

// Zone sharding for regional data residency
// Assign shards to zones
sh.addShardToZone('shard0001', 'EU')
sh.addShardToZone('shard0002', 'US')
sh.addShardToZone('shard0003', 'APAC')

// Define shard key ranges for each region
// (Assuming userId prefix encodes region: 'EU-', 'US-', 'APAC-')
sh.updateZoneKeyRange('ecommerce.orders',
  { userId: 'EU-' }, { userId: 'EU-zzz' }, 'EU'
)
sh.updateZoneKeyRange('ecommerce.orders',
  { userId: 'US-' }, { userId: 'US-zzz' }, 'US'
)

mongos and the Config Server

In a sharded cluster, mongos instances are the query router layer. Application drivers connect to mongos (not directly to shards). The mongos reads the config server replica set (which holds the cluster's metadata, chunk maps, and zone assignments) to determine which shard(s) hold the data for each query. Always deploy at least two mongos instances for high availability — they are stateless and can be restarted without data loss.

Targeted vs Scatter-Gather Queries

In a sharded cluster, a targeted query includes the shard key in its filter — mongos routes it to exactly one shard. A scatter-gather query lacks the shard key — mongos must broadcast it to all shards and merge the results. Scatter-gather queries are expensive and should be avoided on hot paths. Design your shard key and query patterns so the most frequent queries include the shard key field.

// TARGETED: includes shard key (userId) — goes to one shard only
db.orders.find({ userId: 'user-123', status: 'pending' })

// SCATTER-GATHER: no shard key — hits ALL shards (expensive!)
db.orders.find({ status: 'pending', total: { $gt: 100 } })

// Verify with explain in sharded cluster
db.orders.find({ userId: 'user-123' }).explain('executionStats')
// Look for 'SINGLE_SHARD' vs 'SHARD_MERGE' in the winning plan

Capacity Planning and Monitoring

Build a capacity model: estimate daily document growth rate × average document size to project when each shard will fill. In Atlas, use Cluster Autoscaling to automatically add storage or upgrade instance tiers when usage thresholds are crossed. Set Atlas alerts for: disk usage above 80%, CPU utilisation above 70% for more than 1 hour, and replication lag above 10 seconds. Proactive monitoring prevents reactive scrambles at 3 AM.

// Capacity projection script
const avgDocBytes = 1024  // 1 KB average order document
const dailyOrders = 50000
const retentionDays = 365 * 3  // 3 years

const totalOrders = dailyOrders * retentionDays
const totalBytes = totalOrders * avgDocBytes
const totalGB = totalBytes / 1e9

console.log('Projected orders:', totalOrders.toLocaleString())
console.log('Projected storage:', totalGB.toFixed(0), 'GB')
// Add 3x for indexes + WiredTiger overhead
console.log('Recommended disk:', (totalGB * 3).toFixed(0), 'GB')

Atlas Sharded Clusters vs Self-Hosted

Atlas manages the entire sharding infrastructure automatically: provisioning mongos routers, config servers, and shard replica sets; balancing chunks; and patching the cluster. For self-hosted deployments, each component must be provisioned, monitored, and maintained manually — a substantial operational burden. For most teams, the operational savings of Atlas sharding justify the premium over self-hosted, unless compliance requirements mandate on-premises deployment.

Quick Check

Test your understanding of MongoDB & NoSQL Databases concepts from this lesson.

Lesson Recap

In this lesson you learned: scale vertically and leverage secondary reads before sharding — sharding adds significant complexity that most applications never need, shard key selection must balance write distribution (hashed) against query targeting (ranged), and targeted queries that include the shard key go to a single shard while scatter-gather queries hit all shards and are expensive. Next up we complete the capstone with security hardening and the production readiness checklist.

Häufig gestellte Fragen

Ist die Lektion „Skalierungsplan: Vom Replikatsatz zum Sharded Cluster“ kostenlos?

Ja — der vollständige Text von „Skalierungsplan: Vom Replikatsatz zum Sharded Cluster“ ist hier im Web kostenlos zu lesen. Um sie interaktiv zu üben (integrierter Code-Editor und 24/7 KI-Tutor) und den Rest des MongoDB Academy-Kurses freizuschalten, upgrade auf CoddyKit PRO. Der MongoDB Academy-Kurs umfasst insgesamt 4 Lektionen.

Was lerne ich in „Skalierungsplan: Vom Replikatsatz zum Sharded Cluster“?

Lernende erstellen einen Kapazitätsplan und wählen einen Shard Key, der die Verteilung der Lese- und Schreibzugriffe der Anwendung unterstützt, ohne Hotspots zu erzeugen. Du übst MongoDB Academy mit praktischem Code, den du direkt im Browser ausführst, und ein 24/7 KI-Tutor beantwortet deine Fragen während du die Lektion bearbeitest.

Brauche ich Erfahrung, um MongoDB Academy zu starten?

Keine Vorkenntnisse erforderlich. MongoDB Academy auf CoddyKit ist für Anfänger bis fortgeschrittene Lernende strukturiert, sodass du hier starten oder von Anfang an beginnen und in deinem eigenen Tempo voranschreiten kannst. Dies ist Lektion 3 von 4.

Wie lange dauert die Lektion „Skalierungsplan: Vom Replikatsatz zum Sharded Cluster“?

Die meisten CoddyKit-Lektionen dauern etwa 5–10 Minuten. Jede ist kompakt und interaktiv, sodass du stetig Fortschritte machst und genau dort weitermachst, wo du aufgehört hast – im Web und in der App.

Kann ich in dieser MongoDB Academy-Lektion Code schreiben und ausführen?

Ja. Jede MongoDB Academy-Lektion enthält einen integrierten Code-Editor, sodass du echten Code direkt in deinem Browser schreibst und ausführst und sofort KI-Feedback erhältst — ohne lokale Einrichtung erforderlich.

Alle Lektionen in diesem Kurs

  1. Anforderungsanalyse und Schemaentwurf
  2. Indexstrategie und Validierung des Abfrageplaners
  3. Skalierungsplan: Vom Replikatsatz zum Sharded Cluster
  4. Sicherheits-Härtung und Checkliste für den Produktivbetrieb
← Zurück zu MongoDB Academy