Review cards · 18 cards
Sharding
Splitting data across machines: partition keys, consistent hashing and hot partitions.
Cards
- Placing each user's data in a region near them, which also helps meet data-residency laws, is called _____.
- You store 60 TB of data with 3 replicas of everything, and each node can hold 4 TB. How many nodes do you need, ignoring headroom?
- Peak load is 120,000 writes/s. One primary can safely take 8,000 writes/s. How many shards do you need?
- A Snowflake generator has a 12-bit sequence per millisecond. At most how many ids can one worker issue per second?
- Why use consistent hashing instead of hash(key) % N to pick a shard?
- Two tables you often join end up on different shards. What are your options?
- What makes a good partition key?
- Which query gets much more expensive if you partition by a hash of the key instead of by key range?
- A query that cannot be routed by the partition key is sent to every shard and the results merged. This is a _____ query, and its latency is set by the _____ shard.
- How does a Snowflake id fit time, worker and sequence into 64 bits?
- Sensor readings are range-partitioned by timestamp. Why does one shard take all the writes, and what is the fix?
- A document app stores pages as blocks. Most requests load many blocks of one workspace. Which partition key fits?
- Users are sharded by user_id, and you add a global secondary index on email, itself partitioned by email. What do you pay for fast email lookups?
- You shard view counts by video_id. One viral video now gets 200,000 writes a second, far more than one shard can take. What is the usual fix?
- Why create many more logical shards than machines, such as 480 logical shards on 32 databases?
- How do you move a shard's data to a new machine without stopping writes?
- A Snowflake id generator sees its clock jump backwards by 50 ms. What should it do?
- Why does consistent hashing give each server many virtual nodes on the ring instead of one position?
More topics
- Estimation 21 cards
- Networking 16 cards
- API design 17 cards
- Caching 21 cards
- Databases 22 cards
- Replication 15 cards
- Consistency 19 cards
- Queues 18 cards
- Streaming 18 cards
- Availability 14 cards
- Resilience 16 cards
- Storage 14 cards
- Realtime 15 cards
- Data structures 16 cards
- Security 17 cards
- Observability 18 cards
- Coordination 16 cards