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:
Three key problems, whichever you pick:
- 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?
- How does the routing component learn about changes in the shard→node assignment?
- 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:
Consensus algorithms give the mapping fault tolerance plus protection against split brain (Ch 10).
| System | Coordination mechanism |
|---|---|
| HBase, SolrCloud | ZooKeeper |
| Kubernetes | etcd (tracks which service instance runs where) |
| MongoDB | Own config server implementation + mongos daemons as the routing tier |
| Kafka, YugabyteDB, TiDB, ScyllaDB | Built-in implementations of the RAFT consensus protocol |
| Riak | A 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).