Replication gives you copies; sharding gives you capacity, splitting one collection across several replica sets so the working set need not fit on one machine. Each shard is a replica set holding a slice, the config server set records which slice lives where, and mongos routers fan operations out to the shards that can answer them.

A cluster needs three config servers, three processes per shard and a router — more than one machine can honestly demonstrate, so the commands below were not run for this book, unlike everything else here.
mongod --configsvr --replSet cfgrs --port 27019 --dbpath /data/cfg1
mongod --shardsvr --replSet shard1 --port 27018 --dbpath /data/s1
mongos --configdb cfgrs/hostA:27019,hostB:27019,hostC:27019 --port 27017
sh.addShard('shard1/hostA:27018,hostB:27018,hostC:27018') # from mongosh on mongos
sh.shardCollection('shop.orders', { customerId: 'hashed' })
db.orders.getShardDistribution()The shard key is the decision you cannot casually undo, and three properties decide whether it works. Cardinality caps how finely data can split: a key with four values gives four ranges, so a fifth shard never receives data. Frequency decides whether those ranges are even; if 40% of orders belong to one customer, that shard stays hot. Monotonicity decides where writes land: an ObjectId _id sends every insert to the top range, so one shard absorbs the write load.
Hashed keys spread writes evenly but destroy range targeting, turning find({ createdAt: { $gte: ... } }) into a scatter-gather. Compound keys buy both — a high-cardinality prefix for distribution, a range field after it for targeting — which is why { tenantId: 1, createdAt: 1 } suits multi-tenant data. Shard only when one replica set cannot hold the working set or absorb the write rate: below that it buys a router hop, scatter-gather queries and a key you live with.