In a MongoDB sharded cluster, what does mongos do with a find() the driver sends it?
answer
- The application never talks to a shard directly
- It caches a map fetched from elsewhere
- The query's filter decides how many shards see it
- One shard versus every shard
- A stale map triggers a refresh and retry
basics
~20 smongos is the query router. It uses the chunk map cached from the config servers to pick which shards can hold matches — one when the filter names the shard key, all shards otherwise — and merges their replies into a single cursor.
solid answer
~40 s`mongos` is a stateless router, and drivers connect to it rather than to the shards. It keeps a cached routing table describing which ranges of the shard key live on which shard; the authoritative copy of that table lives on the config server replica set. When a `find()` arrives, `mongos` inspects the filter. If the filter contains the shard key, it computes which chunks can match and sends the query only to the shards owning them — a **targeted** query, often a single shard. If the shard key is absent, it broadcasts the query to every shard holding data for that collection — **scatter-gather** — and merges the responses. Either way the client sees one ordinary cursor. Because the router holds no durable state, you run several of them, typically next to the application.
go deeper
Be able to name the three process types — shards, config servers, mongos — and say that the driver connects to mongos, which forwards the query and merges what comes back.
Explain where the routing table comes from, that it is a cache of config-server metadata, and exactly what makes a query go to one shard instead of all of them.
Show that you treat routers as disposable, stateless capacity you place near the application, and that you know a stale routing table produces a retry rather than a wrong answer.
Own the topology question: how many routers, where they run, and how the team keeps the share of broadcast traffic low enough that adding shards still buys throughput.
## What mongos is A sharded MongoDB deployment has three kinds of process: the **shards** (each a replica set holding a subset of the data), the **config server replica set** (CSRS, which stores cluster metadata), and one or more **`mongos`** routers. Applications never connect to a shard directly in normal operation; the driver's connection string points at the routers, and `mongos` presents the cluster as if it were a single logical database. Every command — find, insert, update, aggregate — goes through it. `mongos` is *stateless*: it stores no user data and no durable metadata of its own. It only caches. That is why routers are cheap to run, cheap to restart, and usually deployed several at a time, often as a sidecar on each application host so the extra network hop is a loopback call. ## The routing table A sharded collection's documents are grouped into **chunks** — contiguous ranges of shard-key values — and each chunk is owned by exactly one shard. The mapping from chunk range to shard is cluster metadata, stored on the config servers. `mongos` fetches that mapping and caches it in memory, keyed by collection. The cached table is what makes routing possible without asking anyone. Given a filter, the router can answer "which shards could possibly hold a document matching this?" purely locally. ## Targeted versus scatter-gather The decision hinges on one thing: **is the shard key in the filter?** - If yes, `mongos` maps the shard-key value (or range) to chunks, and from chunks to shards. It contacts only those shards. With an equality match on the whole shard key that is exactly one shard. This is a **targeted** query. - If no, the router has no way to exclude any shard, so it sends the query to **all** shards owning data for the collection. This is **scatter-gather** (a broadcast). Every shard runs the query — using its own indexes, hopefully — and returns its results. A broadcast is not an error; plenty of workloads run them deliberately. It is simply much more expensive: the work done per query grows with the number of shards, and the response is only as fast as the slowest shard that answered. ## Merging and the cursor When more than one shard responds, `mongos` merges the streams before handing results to the client. For an unsorted query that merge is just interleaving. For a sorted one it is a merge-sort over per-shard streams that are each already in order. The client's driver sees one cursor and iterates it normally; batches are fetched through the router as the application reads. ## When the cached table goes stale Chunks move between shards (the balancer migrates them), so a router's cached map can be out of date. MongoDB handles this with versioning rather than locking: each request carries the shard version `mongos` believes is current, and a shard that sees a stale version rejects the request with a stale-config error. `mongos` then refreshes its routing table from the config servers and retries. The application does not see this; it shows up only as an occasional latency blip and as refresh entries in the router's log. ## Practical consequences Three things follow from this design and are worth carrying into an interview: 1. **Routing quality is a property of your queries, not of the cluster.** The same cluster serves one query from one shard and the next from twenty, based only on the filter. 2. **The router is a merge point.** Sorting, limiting and combining results across shards happen there (or on a designated shard for aggregations), so the router does real work on fan-out queries. 3. **Routers are disposable.** Losing one costs you in-flight cursors, not data. Scale them horizontally and put them close to the application. A useful mental check when reading someone's design: for each hot query in the workload, say out loud whether `mongos` can target it. If most of the traffic cannot be targeted, the cluster is doing N times the work per request for N shards, and adding shards will not make those requests faster.
- Does mongos store anything durably of its own?No. It is stateless: it caches the chunk-to-shard routing table and holds in-flight cursors, nothing more. The authoritative metadata lives on the config server replica set. That is why routers are typically run several at a time, often colocated with the application, and can be restarted freely.
- What happens if a chunk migrates to another shard after mongos cached the routing table?The router sends the request tagged with the shard version it believes is current. The shard notices the version is stale and returns a stale-config error instead of wrong results. mongos refreshes its routing table from the config servers and retries the operation transparently, so the application sees at most extra latency.
- Can you still connect directly to one shard's replica set?You can, and it is sometimes done for diagnostics, but a direct read sees only that shard's slice of the collection — and it can also see orphaned documents left behind by an interrupted migration, which the router-level path filters out. Application traffic should always go through mongos.
saying these in an interview costs you the question
- Says the application connects to the shards directly
- Thinks mongos stores the chunk metadata durably
- Believes every query goes to every shard
- Claims mongos routes by hashing _id
- Assumes the shard, not the filter, decides routing