Learn Labs
7. Sharding

7.10 Self-test

Self-test21 questions

—/21
  1. Distinguish sharding from replication in one sentence each. How do they combine under single-leader replication?

  2. Why is sharding recommended only at large scale, when replication is useful at any scale?

  3. Name four costs sharding imposes. Which one is hardest to reverse?

  4. Give three reasons a SaaS product might shard per tenant that have nothing to do with scalability.

  5. Define skew, hot shard, and hot key. Give an example where hashing eliminates the first but not the third.

  6. Why must key-range shard boundaries "adapt to the data"? Give the encyclopedia illustration.

  7. A time-series database keyed by timestamp has one shard at 100% and the rest idle. Diagnose it, fix it, and state precisely what the fix costs you.

  8. Why is splitting a shard both necessary and dangerous?

  9. Work out how many keys move when going from 4 nodes to 5 under mod N versus under a fixed-shard scheme.

  10. In the fixed-shard scheme, three things could change during rebalancing — which one actually does?

  11. Why choose a shard count that is "divisible by many factors"? Why can you not have more nodes than shards?

  12. State the Goldilocks problem of fixed shard counts.

  13. How does hash-range sharding get the benefits of both key-range and hash sharding? What does it still lose, and what partially rescues it?

  14. What two properties define a consistent hashing algorithm? What does "consistent" not mean here?

  15. Explain why salting a hot key helps writes but not reads, and what bookkeeping it forces on you.

  16. Draw the cascading-failure loop caused by automatic rebalancing plus timeout-based failure detection.

  17. Name the three request-routing architectures and the three problems common to all of them.

  18. Why is ZooKeeper/etcd used for shard assignment rather than a single coordinator? What does Riak do instead, and what does it give up?

  19. For local and global secondary indexes: which is cheap on write, which is cheap on single-condition read, and why does neither give you cheap multi-condition reads?

  20. Why does adding shards not increase query throughput for a local secondary index?

  21. Design question

    you're building an event-tracking system: 500k events/s, 30-day retention, queried as (a) "all events for user X in the last hour" (high volume, low latency) and (b) "count of event type Y per hour across all users" (analyst, minutes acceptable). Choose the sharding scheme, the partition and clustering keys, the secondary-index strategy, and the routing architecture. Justify each against a specific trade-off from this chapter, and name the failure mode you're most worried about.