SD Core

Sharding

Splits data across multiple database nodes by a shard key so write/read capacity scales beyond one machine.

Interview tip Lead with a 30-second definition, then one real system example and name 2–3 designs where Sharding is non-negotiable.

① What it is (30 seconds)

Splits data across multiple database nodes by a shard key so write/read capacity scales beyond one machine.

② How it works in system design

Choose shard key (user_id, tenant_id). hash(key) → shard index. Cross-shard queries expensive — design access patterns per shard. Rebalancing moves ranges when adding nodes.
Typical placement
ClientEdge / GatewayShardingServicesData stores

③ Concrete system design example

Scenario: Twitter timelines: tweets sharded by user_id. Write tweet to author shard; fan-out workers read follower lists and write to follower shard inboxes.

④ Important interview Q&A

QuestionAnswer
Shard key selection?High cardinality, even distribution, aligns with query patterns — avoid cross-shard joins.
Hot shard?Split range, add read replicas, or reshard by sub-key.
Shard vs read replica?Replica scales reads on same data; sharding splits data for write scale.

⑤ Seen in these system designs

In interviews, after explaining the concept, say: "This shows up directly in …" and link two designs.

⑥ Revision checklist

  • Shard key justified
  • Cross-shard query plan
  • Rebalance strategy
  • Hot shard mitigation
shardingdatabasescale