Design a Distributed Document Database

Problem

Design MongoDB's core: a distributed document database with sharding and replica sets.

Requirements

Functional:

  • Store/query JSON documents with indexes
  • Shard by key for horizontal scale
  • Replica sets for HA with automatic failover
  • Tunable read/write concern

Non-functional:

  • Petabyte scale, high throughput
  • Strong durability; tunable consistency
  • Automatic failover

Discussion points

  1. Sharding (shard key, chunk balancing)
  2. Replica sets + Raft-like election
  3. Storage engine (WiredTiger, MVCC, B-tree)
  4. Read/write concerns and consistency
  5. Secondary indexes and query routing
added …
LeaderboardSalaryAccount