7.1 Pros and Cons of Sharding
The primary reason is SCALABILITY — data volume or write throughput has become too great for a single node.
If READ throughput is the problem, you don't necessarily need sharding — you can use read scaling (Ch 6 followers).
Sharding is one of the main tools for HORIZONTAL SCALING (scale-out) — growing capacity not by moving to a bigger machine, but by adding more smaller machines. If you can divide the workload so each shard handles a roughly equal share, you can assign shards to different machines to process data and queries in parallel.
While replication is useful at BOTH small and large scale — it enables fault tolerance and offline operation — SHARDING IS A HEAVYWEIGHT SOLUTION THAT IS MOSTLY RELEVANT AT LARGE SCALE. If a single machine can handle your data volume and write throughput (AND A SINGLE MACHINE CAN DO A LOT NOWADAYS!), it's often better to avoid sharding and stick with a single-shard database.
The four costs
| Cost | Detail |
|---|---|
| The partition key choice is consequential and hard to change | All records with the same partition key go in the same shard. Accessing a record is fast if you know which shard it's in; if you don't, you have to do an INEFFICIENT SEARCH ACROSS ALL SHARDS. And the sharding scheme is difficult to change |
| Relational data is harder than key-value | Sharding works well for key-value data, where you can easily shard by key, but harder with relational data where you may want to search by a secondary index or JOIN records distributed across shards (§6) |
| Cross-shard writes need distributed transactions | A write may need to update related records in several shards. Single-node transactions are common; cross-shard consistency requires a DISTRIBUTED TRANSACTION — usually much slower, and may become a bottleneck for the system as a whole (Ch 8) |
| Operational complexity | Rebalancing, routing, monitoring skew — all new work |
A non-scalability use: sharding on ONE machine. Some systems run one single-threaded process per CPU core to exploit CPU parallelism or a NUMA (nonuniform memory access) architecture where some memory banks are closer to one CPU than others. Redis, VoltDB, and FoundationDB use one process per core and rely on sharding to spread load across cores in the same machine.