Sharding
In this chapter, we divide Shopend’s data across several database servers, so that its writes are no longer limited to what one primary can accept.
After reading this chapter, you should be able to:
- Define sharding and shard, and explain why replication does not add write capacity but sharding does
- Explain why database operations need a shard router
- Define shard key, and explain why Shopend shards by
shopId - Explain how a directory maps shops to shards, and trace a database operation through the router and the directory
- Assign shops to shards based on their workloads, and describe how to move a shop to another shard without losing its writes
- Explain how a report across all shops is computed, and how a slow or unavailable shard affects it
- Trace an operation to its shard and then to that shard’s primary or a replica
- Define partitioning, partition key, and partition pruning, and choose between partitioning and sharding
- Compare a central counter, ranges per shard, UUIDs, and Snowflake IDs for creating unique IDs across shards