Senior (5+ years)System Design

What is the difference between database sharding and replication?

Quick answer

Replication copies the same data to several servers to improve read capacity and availability, while sharding splits the data across servers so each holds only a part, which increases write capacity and total storage.

With replication, one primary accepts writes and replicas follow it. Reads can be spread across replicas and a replica can be promoted if the primary fails, but replicas can lag, so reads from them may be slightly stale, and write throughput is still limited by the single primary.

Sharding partitions rows by a shard key, such as user_id, so each shard handles a fraction of the traffic. Choosing the key is the critical decision: a bad key creates hot shards, and cross-shard queries, joins and transactions become expensive or impossible. Range sharding supports range queries but can be uneven; hash sharding spreads load evenly; consistent hashing makes adding shards cheaper. Shard only when replication, caching, indexing and a bigger primary no longer suffice, because it adds permanent operational complexity.

Key points

  • Replication: copies, for reads and failover
  • Sharding: partitions, for writes and storage
  • Shard key choice decides hot spots and query cost
  • What is the CAP theorem and what does it mean in practice?

    The CAP theorem says that during a network partition a distributed system must choose between consistency (every read sees the latest write) and availability (every request gets a response); you cannot have both while the partition lasts.

  • Monolith vs microservices: which should you choose?

    Start with a well-structured monolith unless you have several teams and clear service boundaries, because microservices add distributed-system complexity such as network failures, data consistency and operational overhead.

  • How do you design a URL shortener like bit.ly?

    Generate a short unique code for each long URL (for example by base62-encoding a unique ID), store the mapping in a key-value or relational database, serve redirects through a cache because reads far outnumber writes, and record analytics asynchronously.

  • How do you design a rate limiter?

    Choose an algorithm such as token bucket or sliding window, store a counter per client (by user ID, API key or IP) in a fast shared store like Redis, and return HTTP 429 with a Retry-After header when the limit is exceeded.