In a sharded relational deployment, something must map each request's shard key to a specific database server. Compare putting that logic in the application/client library, in a dedicated proxy tier, and in a directory or lookup service, and say what each buys you.
answer
- App computes shard = no hop, shared topology
- Proxy = one implementation, merges, hides resharding
- Directory = placement as data, pin the whale
- Cache placement, shard rejects stale routes
- Buckets in the directory, not raw keys
basics
~20 sClient-side routing embeds the mapping in the application: fastest, no extra hop, but every service must agree on the topology and redeploy when it changes. A proxy tier centralises routing, pooling and cross-shard merging at the cost of a hop and a component to run. A directory service stores explicit per-key placement, allowing arbitrary moves and pinning, at the cost of a lookup and a critical dependency.
solid answer
~1 min**Client-side routing**: the app computes `shard = f(key)` locally and connects directly to that database. Lowest latency, no extra failure domain, but the topology becomes shared state across every service and language; adding a shard or moving a bucket means a coordinated config rollout, and each app instance holds pools to every shard, so connection counts multiply. **Proxy / routing tier** (the Vitess `vtgate` model): the app speaks the ordinary wire protocol to a proxy that parses the query, extracts the key, routes it, and — importantly — can also fan out, merge, sort and aggregate cross-shard results, and hide resharding behind itself. Costs: an extra network hop, a stateless but critical tier to scale and monitor, and a component that must understand your SQL. **Directory / lookup service**: a mapping table from key (or key range) to shard, consulted before routing and heavily cached. It decouples placement from any hash function, so you can pin an oversized tenant to dedicated hardware or migrate a single key. Costs: an extra lookup on the hot path, cache-invalidation correctness during moves, and a service whose unavailability blocks everything. In practice these compose: a proxy tier that consults a cached directory is the common mature shape.
go deeper
Know that something must map a key to a server, and that it can live in the app or in a proxy in front of the databases.
Compare the three placements with concrete costs — extra hop, config rollout, connection counts — and name a real system that uses each style.
Recommend a layered design and justify it by operational needs: online resharding, pinning hot tenants, and safe behaviour when a client's topology view is stale.
Frame it as where you want the coupling: a proxy trades latency and an owned component for the ability to evolve topology without touching applications, which is usually the decisive factor over a multi-year lifetime.
## The problem Every request in a sharded system must be resolved to a physical server before a byte of SQL executes. Three architectural placements exist for that resolution step, and they are not mutually exclusive. ## Client-side routing A library inside the application holds the topology — shard count, hash function, bucket-to-server table, connection details — and computes the destination itself. The app then opens an ordinary connection to that one database. *Strengths.* No extra hop, so latency is the raw database latency. No additional tier to run, capacity-plan or page someone about. Debugging is direct: the query you see is the query the shard runs. *Weaknesses.* The topology becomes distributed shared state. Every service, in every language, must implement the same hash and hold the same bucket map; a mistake in one implementation silently writes rows to the wrong shard, which is among the nastiest bugs in this space. Changing the topology is a coordinated deploy across all clients, and during the rollout different clients disagree about placement, so moves must be carefully staged. Connection fan-out is a real constraint: with C app instances and S shards you can approach C×S pooled connections, which relational engines handle poorly. Finally, anything that spans shards — a fan-out aggregate, an ordered page — has to be written by hand in each application. Client routing suits a small number of services, a stable shard count, and workloads that are essentially always single-shard. ## A dedicated routing tier A stateless proxy speaks the database's own wire protocol, so applications connect to it as if it were a single database. It parses incoming SQL, finds the shard key in the predicate, routes the statement, and returns results. Vitess's `vtgate` is the canonical example; Citus's coordinator plays a comparable role from inside the database. *Strengths.* One place implements routing, so all clients are automatically consistent. The proxy can multiplex thousands of client connections onto a modest pool per shard, solving the connection-fan-out problem. It can execute cross-shard work — scatter the query, merge sorted streams, combine partial aggregates, enforce a limit — so applications keep writing something close to normal SQL. Above all, it is the layer that makes **online topology change** possible: during a bucket move or resharding it can hold, retry or redirect traffic so clients never see the transition. *Weaknesses.* One extra network round trip on every query, which matters for sub-millisecond point reads. A new tier to deploy, scale, upgrade and monitor, and one that is on the critical path for 100% of traffic. It must parse your SQL dialect, so exotic statements may be unsupported or routed conservatively (i.e. scattered). Query cost attribution gets murkier because everything comes from the proxy's IP. ## Directory / lookup service Here placement is *data*, not a formula: a table maps key (or key range, or virtual bucket) to shard. Callers — client library or proxy — read it, cache it aggressively, and refresh on change or on a routing error. *Strengths.* Total flexibility. Any single key can be moved anywhere; an outsized tenant can be pinned to dedicated hardware; a migrating key can be marked read-only or dual-written for the duration of its move. Placement decisions become an operational action rather than a code change, which is exactly what you want when skew is discovered in production. Directory-based routing also makes shard-count changes uneventful, since nothing derives placement from the number of shards. *Weaknesses.* The lookup is on the hot path, so it must be cached — and caches of placement data are dangerous during moves, because a stale entry sends a write to the old shard. Robust designs make the shard itself authoritative: if a request arrives for a key the shard no longer owns, it rejects with a "moved" response that triggers a cache refresh and retry, the same discipline redirect-based clustered systems use. The directory is also a hard dependency; it needs its own replication and failover, and it grows with key count if you map individual keys rather than buckets (mapping buckets, not keys, keeps it small). ## How they compose The mature shape is layered: a **proxy tier** for protocol compatibility, connection multiplexing and cross-shard execution; a **directory of virtual buckets** for placement, so moves are metadata operations; and a **hash** to assign keys to buckets, so the directory stays small. Client-side routing survives where latency budgets are tight and topology is static, often as a fast path with proxy fallback for anything non-trivial. A related question is where routing metadata lives and how changes propagate — typically a consensus-backed store (etcd, ZooKeeper, or the coordinator's own catalog) with versioned topology, so clients can detect that their view is stale rather than acting on it blindly. ## How to answer Name the three placements, give one sentence of strength and one of cost for each, and then say what you would actually build: usually a proxy over a bucket directory, justified by the ability to reshard and to pin hot keys without touching application code.
- With client-side routing, how do you change the number of shards without a period where different application instances disagree about placement?You avoid deriving placement from the shard count: keep a fixed set of virtual buckets and version the bucket-to-shard map. Clients fetch the versioned map from a shared store and the shards themselves validate ownership, rejecting requests for buckets they no longer hold so the client refreshes and retries. Without that ownership check, a rolling config deploy inevitably has a window where two clients route the same key to different servers.
- What extra capability does a routing proxy give you that a client library realistically cannot?Cross-shard query execution and connection multiplexing. The proxy can scatter a statement, merge sorted streams, combine partial aggregates and apply the final limit, presenting a single result set; doing that consistently in every client language is impractical. It also collapses thousands of client connections into small per-shard pools, which matters because relational engines degrade badly under connection counts that grow as clients times shards.
saying these in an interview costs you the question
- Treating routing as trivial and ignoring the connection fan-out of client-side pools
- Assuming a config push can change topology safely without shard-side ownership checks
- Forgetting that a directory lookup must be cached and that stale cache entries misroute writes
- Presenting the proxy as free, ignoring the extra hop and that it is on the critical path for all traffic
- Implementing the hash independently in several services and expecting them to agree