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

Sections

  1. Dividing the Data
  2. Keeping a Shop Together
  3. Finding the Right Shard
  4. Balancing Shops Across Shards
  5. Operations Across Shops
  6. Replication and Sharding Together
  7. Partitioning a Table
  8. Unique IDs Across Shards