Once data is sharded across multiple databases, how do systems handle queries that need to touch multiple shards - such as a join across shards or a transaction spanning two shards?
answer
- scatter-gather fan-out
- co-locate shard key for local joins
- denormalize to avoid joins
- 2PC blocking risk
- sagas plus compensating actions
basics
~20 sIf you need data from more than one shard, the app or a routing layer must ask each shard and combine the results itself, since the database can't do it in one step. Transactions across shards need extra coordination, like a two-step commit protocol, or are avoided by design.
solid answer
~40 sCross-shard reads are handled by scatter-gather: the querying layer fans a request out to every relevant shard in parallel, then merges, sorts, or aggregates the partial results (as in Vitess or Citus), costing latency bounded by the slowest shard. Cross-shard joins are the hardest case - either co-locate related data on the same shard by choosing a shared shard key, denormalize so the join isn't needed, or accept a slow scatter-gather join. Cross-shard transactions require a distributed-transaction protocol like two-phase commit (2PC), or a database with built-in distributed transactions (Spanner, CockroachCB use consensus-backed atomic commit across ranges); most teams instead avoid multi-shard transactions by designing shard boundaries so related writes land on one shard, or by using sagas/compensating transactions for cross-shard processes that can tolerate eventual consistency.
go deeper
Knows that a query needing two shards' data isn't a single simple query anymore.
Can describe scatter-gather and denormalization as basic mitigations.
Designs shard keys for co-location, explains 2PC's blocking failure mode, and can propose a saga for a cross-shard business process.
Chooses between building/adopting a distributed-SQL layer versus application-level sagas for an entire platform, weighing consistency guarantees against operational complexity and vendor lock-in.
## What sharding gives up Sharding trades a single database's convenience - one engine that can join, aggregate, and transact across all your data - for horizontal scale. The moment data lives on separate physical databases, any query or transaction that needs data from more than one shard becomes the application's or a middleware layer's problem to solve, because no single database engine can natively see across the shard boundary. ## Reads: scatter-gather For reads, the standard mechanism is **scatter-gather** (fan-out query): a routing layer sends the same or a shard-specific version of the query to every shard that might hold relevant data, waits for each shard's partial result, and then merges them - concatenating, re-sorting, deduplicating, or applying an aggregate function (SUM, COUNT, AVG) across the partial results. This is exactly what Vitess (MySQL sharding for Slack/YouTube), Citus (Postgres sharding), and Elasticsearch's cross-shard search do under the hood. The cost is real: - total query latency is bounded by the slowest shard that must respond (**tail-latency amplification** - the more shards you fan out to, the higher the chance one is briefly slow); - any query needing global ordering (e.g. 'top 10 by score across all shards') requires pulling a candidate set from every shard rather than just the top 10, because a locally-top-10 result on one shard might not be globally top 10. ## Joins: the sharpest pain point Cross-shard joins are the sharpest pain point, because a relational join fundamentally assumes both tables are reachable from the same query engine. Three practical strategies exist. 1. **First, co-location:** choose the shard key so data usually joined together lands on the same shard - e.g. shard both `orders` and `order_items` by `customer_id` so a customer's order and its line items are always on the same physical shard, turning the join into a normal local join. This is the most common and highest-performing fix, which is why shard-key selection is a schema-design decision, not an infrastructure afterthought. 2. **Second, denormalization:** duplicate the data you'd otherwise join, accepting eventual consistency and extra storage in exchange for avoiding the join entirely (e.g. store a denormalized `customer_name` on every order row instead of joining to a customers table on every read). 3. **Third, application-side join:** fetch both sides separately, often from different shards, and join them in application code - workable for small result sets, prohibitively slow and memory-heavy for large ones. ## Transactions: two-phase commit and its blocking cost Cross-shard transactions raise the stakes further, because now you need **atomicity** - either both shards commit the change or neither does - across physically separate databases that don't share a transaction log. The classical mechanism is **two-phase commit (2PC)**: a coordinator asks every participating shard to 'prepare' (lock the rows and confirm it can commit), and only after every shard says yes does the coordinator tell them all to actually commit; if any shard says no, or times out, everyone rolls back. 2PC gives real atomicity but at a real cost: it's synchronous and blocking - if the coordinator crashes after shards have prepared but before it sends commit/abort, those shards are stuck holding locks in an undetermined state until the coordinator recovers, which can stall unrelated transactions on the affected rows. Modern distributed SQL systems (Google Spanner, CockroachDB, YugabyteDB) build atomic cross-shard transactions on top of a consensus protocol (Raft/Paxos) per shard/range plus a coordination layer, and use techniques like Spanner's TrueTime (a globally synchronized clock with bounded uncertainty) to give externally consistent timestamps to cross-shard transactions without a slow, blocking 2PC round-trip for every case. ## The pragmatic default: design the problem away Given this cost, the pragmatic default in most sharded systems is to avoid multi-shard transactions by design rather than solve them technically: - choose shard boundaries so any single business transaction's writes land on one shard (the same co-location principle as for joins); - and for genuinely cross-shard business processes (e.g. transferring a balance between two users on different shards, or a multi-step order-fulfillment workflow), use a **saga** - a sequence of local transactions, each on one shard, with a compensating action defined for each step to undo it if a later step fails. Sagas trade strict atomicity for eventual consistency and require the business logic to tolerate a temporarily inconsistent intermediate state (e.g. money debited from one account but not yet credited to the other for a brief window), which is why they're paired with idempotency keys and retry logic, and why 'transfer between shards' remains one of the genuinely hard problems in sharded system design.
- Why does two-phase commit (2PC) across shards create an availability risk even when all shards are individually healthy?Because during the prepare phase, shards hold locks and wait for the coordinator's final commit/abort decision. If the coordinator crashes or the network partitions after prepare but before that decision reaches all shards, those shards are stuck holding locks indefinitely until the coordinator recovers, which can stall unrelated transactions touching the same rows.
- If you shard an orders system by order_id, why might a 'get all orders for this customer' query become expensive, and how would you fix it?Because a customer's orders are scattered across every shard by order_id's hash or range, so the query has to scatter-gather across all shards instead of hitting one. The fix is usually to shard by customer_id instead, so a customer's data co-locates on one shard for that access pattern.
- What does a saga give up compared to a true cross-shard ACID transaction, and why is that often an acceptable trade?A saga gives up atomicity and isolation - intermediate steps are visible and a failure partway through leaves a temporary inconsistent state that must be reversed by compensating actions rather than rolled back atomically. It's often acceptable because many business processes can tolerate a brief, bounded inconsistent window as long as it's eventually corrected and idempotent.
Querying across shards is like asking ten separate branch libraries, instead of one central library, whether they have a book - you have to call every branch, wait for the slowest one to answer, and combine the answers yourself. Doing a 'transaction' across two branches, like transferring a book from branch A's catalog to branch B's, needs a phone call where both branches promise to update their records at the same time, or the whole thing has to be undone.
saying these in an interview costs you the question
- assumes SQL joins 'just work' across shards without any design effort
- doesn't know 2PC exists or confuses it with simple retries
- proposes cross-shard transactions with no discussion of coordinator failure
- picks a shard key without considering the queries that need co-located data
- thinks sagas provide the same guarantees as ACID transactions