Learn Labs
7. Sharding

7.4 Request Routing

The question: to read or write a particular key, which node — which IP address and port — do you connect to?

Very similar to SERVICE DISCOVERY (Ch 5). The biggest difference: with application services, each instance is usually STATELESS and a load balancer can send a request to any instance. With sharded databases, A REQUEST FOR A KEY CAN BE HANDLED ONLY BY A NODE THAT IS A REPLICA FOR THE SHARD CONTAINING THAT KEY.

Three approaches:

① Any node — round-robin load balancer
forwards, then repliesreplyclientnode Anode Bowns the shardclientNode A forwards the request if it doesn't own the shard,and passes the reply back — so there is one extra hop.
② Routing tier
clientrouting tiershard-aware LBnode BThe routing tier handles no requests itself.
③ Shard-aware client
direct — no intermediaryclientknows the mappingnode B
Figure 7.4.1Three approaches

Three key problems, whichever you pick:

  1. Who decides which shard lives on which node? Simplest is a single coordinator — but then how do you make it fault-tolerant if that node goes down? And if the coordinator role can fail over, how do you prevent SPLIT BRAIN where two coordinators make CONTRADICTORY shard assignments?
  2. How does the routing component learn about changes in the shard→node assignment?
  3. The cutover window: while a shard is being moved, the new node has taken over but REQUESTS TO THE OLD NODE MAY STILL BE IN FLIGHT. How do you handle those?

The standard answer — a coordination service:

subscribe / notify on changeeach noderegisters itselfZooKeeper / etcdshard 1 → node A · shard 2 → node B · shard 3 → node Crouting tiershard-aware client

Consensus algorithms give the mapping fault tolerance plus protection against split brain (Ch 10).

Figure 7.4.2The standard answer — a coordination service
SystemCoordination mechanism
HBase, SolrCloudZooKeeper
Kubernetesetcd (tracks which service instance runs where)
MongoDBOwn config server implementation + mongos daemons as the routing tier
Kafka, YugabyteDB, TiDB, ScyllaDBBuilt-in implementations of the RAFT consensus protocol
RiakA GOSSIP PROTOCOL among nodes — much weaker consistency than consensus; SPLIT BRAIN IS POSSIBLE, with different parts of the cluster having different node assignments for the same shard. Leaderless databases can tolerate this because they make weak consistency guarantees anyway

And for finding IPs in the first place: shard→node assignment is fast-changing, but node IP addresses are not — so DNS is often sufficient for that layer.

This discussion focused on finding the shard for an INDIVIDUAL KEY — most relevant for sharded OLTP. Analytical databases shard too, but with a very different query execution: rather than executing in a single shard, a query commonly needs to AGGREGATE AND JOIN DATA FROM MANY SHARDS IN PARALLEL (Ch 11).