skip to content

In a sharded relational database, why can't you run a JOIN across two tables that live on different shards the same way you would in a single-instance database, and how do teams typically handle a transaction that must update rows on two different shards?

level: seniorimportance: must knowfreq 65%

answer

  1. query engine sees only its own shard
  2. scatter-gather = fan out + merge in app
  3. 2PC = atomic but locks + coordinator risk
  4. saga = local txns + compensating actions
  5. co-locate by shard key is the preferred fix

basics

~20 s

Each shard is a separate database that doesn't know about the others, so it can't run a join or a single all-or-nothing transaction spanning two of them. The application has to query each shard separately and combine results, and cross-shard writes need special handling like two-phase commit or a saga, or the schema should be redesigned so related data lives on the same shard.

solid answer

~50 s

A JOIN runs inside one database engine's query planner using local indexes and local transactional guarantees; when the two tables sit on physically separate shard instances, no single engine can see both, so the 'join' has to be done at the application or middleware layer — query each shard separately (scatter-gather), then merge in memory, which is slower and loses the engine's join optimizations. For writes spanning shards, a plain local transaction (BEGIN/COMMIT) can't span two processes, so teams either implement two-phase commit or a saga/compensating-transaction pattern to get atomicity (or eventual consistency) across shards, at the cost of more complex failure-handling code, or — the preferred fix whenever possible — redesign the schema so operations that must be transactional together live on the same shard by co-locating related data under the same shard key. The general principle: sharding trades cross-entity transactional and query convenience for horizontal write and storage scalability.

go deeper

for a junior

Should understand that data on different shards can't be joined the way data in one database can.

for a middle

Should be able to describe scatter-gather at a basic level and know cross-shard writes aren't simply atomic.

for a senior

Should compare two-phase commit and saga trade-offs concretely and identify shard-key co-location as the preferred first fix.

for a principal

Should reason about schema/shard-key design proactively to minimize cross-shard operations organization-wide, and weigh when a saga's eventual consistency is acceptable for the business versus when it isn't.

## Why a JOIN is a single-engine operation A JOIN, in a single-instance relational database, is executed entirely inside one query planner: the engine has direct access to both tables' storage, statistics, and indexes, and it can choose an efficient join strategy (nested loop, hash join, merge join) because everything it needs is local and visible within one transaction context. Once two tables — or even two rows of the same logical table — live on physically separate shard instances, that assumption breaks completely: shard A's query engine has no visibility into shard B's data, no shared index, and no way to execute a single query plan spanning both. The same is true for **ACID transactions**: a local transaction's atomicity and isolation guarantees are enforced by one engine's transaction manager and write-ahead log, and that mechanism has no native way to make a write on shard A and a write on shard B commit or roll back together. ## Cross-shard reads: scatter-gather Because of this, cross-shard reads have to be reimplemented above the database layer. The typical pattern is **scatter-gather**: a routing/middleware layer sends the equivalent query to every shard that might hold relevant data, each shard executes it locally and returns its own partial result set, and the middleware merges, sorts, or aggregates the partial results into a final answer. This works, but it is strictly worse than a native join on several dimensions: - it can't use cross-table indexes the way a single engine could, - it pays the network and serialization cost of moving partial results back to a coordinator, - and its latency is bound by the slowest responding shard rather than by one engine's local execution time. ## Cross-shard writes: two strategies Cross-shard writes are the harder problem, because "sometimes atomic, sometimes not" is not an acceptable outcome for money, inventory, or anything with a correctness requirement. Two broad strategies exist. 1. The first is **distributed transaction protocols** like **two-phase commit (2PC)**: a coordinator asks every participating shard to "prepare" (lock and stage the change, promising it can commit if told to), and only after every shard confirms it's ready does the coordinator tell them all to actually commit; if any shard fails to prepare, everyone rolls back. This gives real atomicity but at a real cost — it holds locks across all participating shards for the duration of the protocol, is vulnerable to the coordinator crashing mid-protocol (leaving participants blocked, unsure whether to commit or roll back), and doesn't scale well as a general-purpose pattern because it serializes on the slowest participant and the network round trips involved. 2. The second strategy is the **saga pattern**: instead of one atomic transaction, you break the operation into a sequence of local transactions, each on a single shard, with a corresponding compensating action defined for each step (e.g., "cancel reservation" undoes "reserve inventory"). If a later step fails, the saga runs the compensations for the steps that already succeeded, restoring an equivalent-to-rolled-back state without ever holding a distributed lock. Sagas trade strict atomicity (there's a window where the system is visibly in an intermediate state) for much better availability and throughput, and they require the application to write explicit compensation logic for every step — real engineering work that a native ACID transaction gave you for free. ## The preferred fix is co-location Because both of these are genuinely expensive, the most common and most preferred fix in practice is architectural, not protocol-level: choose the shard key so that data which must be transactionally consistent together is co-located on the same physical shard in the first place. If an e-commerce system shards by `customer_id` and ensures a customer's orders and that customer's cart both live on the shard determined by their `customer_id`, then "create an order from a cart" is a single-shard, single-engine, ordinary local transaction — the cross-shard problem never arises for that operation, because the schema design anticipated which operations needed atomicity and shaped the shard key around them. ## Failure modes Failure modes in production tend to concentrate around **partial completion** and **latency amplification**. - With naive cross-shard writes lacking any coordination protocol, a crash or timeout after shard A's write succeeds but before shard B's write is attempted leaves the system in a permanently inconsistent state — an order exists with no corresponding inventory decrement, for instance — with no automatic mechanism to detect or fix it. - With scatter-gather reads, tail latency degrades as the shard count grows, since the overall response time is bound by whichever shard is slowest to respond, and a single overloaded or partially-down shard can make an otherwise-healthy cluster's cross-shard queries look broken. ## Where it shows up A concrete real-world pattern: Vitess's VTGate router transparently executes scatter-gather queries across MySQL shards for queries that can't be satisfied by a single shard, while explicitly documenting that cross-shard JOINs and transactions come with real performance and consistency caveats compared to same-shard operations — which is exactly why schema design guidance for Vitess (and for most sharded systems) pushes teams toward choosing shard keys that keep related, jointly-transacted data together, treating true cross-shard transactions as an exception to design around rather than a routine operation to rely on.

  • What is a 'scatter-gather' query, and what does it do to tail latency?
    It's when the routing layer sends the same (or an equivalent) query to every shard in parallel and merges the results afterward. Overall response time is bound by the slowest shard's reply, so as shard count grows, the probability that at least one shard is momentarily slow increases, which pushes tail (p99-style) latency up even if the average per-shard latency looks fine.
  • What's the usual first fix teams reach for before resorting to distributed transactions across shards?
    Re-picking the shard key so that data which must be transactionally consistent together — like an order and its line items — is co-located on the same shard, turning what would be a cross-shard transaction into an ordinary local one. This avoids the cost and risk of two-phase commit or sagas entirely for that operation, and is generally preferred whenever the access pattern allows it.

It's like trying to cross-reference two separate filing cabinets kept in different office buildings — you can't just flip between drawers the way you would in one cabinet; someone has to physically go pull files from both buildings and manually collate them, and if you need a change to apply to both buildings atomically, you need an explicit handshake (or a plan to undo one building's change if the other fails) rather than a single lock on one cabinet.

saying these in an interview costs you the question

  • Believes cross-shard joins work identically to single-instance joins
  • Never mentions scatter-gather or application-level merging for cross-shard reads
  • Jumps straight to two-phase commit without mentioning co-location as the preferred fix
  • No mention of partial-failure/atomicity risk for uncoordinated cross-shard writes
  • Thinks sagas provide the same isolation guarantees as a local ACID transaction

context