skip to content

When a dataset is split across shards or owned by separate services so the engine cannot enforce a foreign key across the boundary, how do you decide where that invariant lives and how do you keep the data trustworthy?

level: principalimportance: nice to knowfreq 30%

answer

  1. Fix the shard key before accepting the loss
  2. Immediate+total → eventual+measured
  3. One owner, enforce at owner's write path
  4. Outbox + idempotent consumers
  5. Reconciliation metric + repair runbook

basics

~20 s

Choose a partition key that keeps related rows together so most invariants stay engine-enforced. For the ones that genuinely cross a boundary, name an owning service, enforce on write in that owner, and add continuous reconciliation plus a repair path — the invariant becomes eventual, so design for detecting and fixing violations.

solid answer

~1 min

First, avoid the problem where you can: pick a partition or ownership boundary that keeps an aggregate's rows co-located, so referential integrity, uniqueness and check rules remain engine-enforced within each shard. Most 'we can't use foreign keys' situations are a shard-key choice, not a law of physics. For invariants that truly cross a boundary, accept that they become **eventual** and design accordingly: - **Name an owner.** Exactly one service is authoritative for the rule; others read a copy and must not assume it. - **Enforce at the write path of the owner**, so the only way to create a violation is a bug or an out-of-band write. - **Make violations detectable.** A reconciliation job continuously compares the two sides and emits a metric; the number of orphans is a monitored signal, not a surprise found during a migration. - **Make them repairable.** A documented remediation — backfill, quarantine, or compensating action — with a decision about which side wins. - **Prefer designs that tolerate the gap**: idempotent consumers, soft references that render gracefully when the target is missing, and outbox-based propagation so the write and the notification share a transaction. The judgement call is how much integrity you trade for the scaling or autonomy you gained, stated explicitly rather than discovered later.

go deeper

for a junior

Recognize that foreign keys work only inside one database, and that crossing databases means the check moves into code.

for a middle

Describe the mitigation basics: validate in the owning service, propagate changes reliably, and run a job that finds orphans.

for a senior

Design the whole loop — ownership, outbox with idempotent consumers, reconciliation metric, repair runbook, deletion policy.

for a principal

Lead with boundary design (make most invariants stay engine-enforced), then state the explicit trade for those that cannot, including why two-phase commit is usually the wrong price.

## First: is the boundary real? The common failure is treating loss of referential integrity as unavoidable when it was a design choice. Two questions come first. **Did the partition key have to split this?** If orders and order lines land on different shards, the shard key is wrong. Choosing the aggregate root (customer id, tenant id, order id) as the partition key keeps the rows that constrain each other in one engine, where foreign keys, unique constraints and checks work normally. Split only across genuinely independent aggregates. **Did service ownership have to split this?** A rule that spans two services is a signal that the boundary cuts through a single consistency domain. Sometimes the right answer is moving the data, not building distributed enforcement. A principal-level answer starts here, because the cheapest cross-boundary invariant is the one that never crosses. ## When it genuinely crosses Accept that the invariant changes character. Inside one engine it was **immediate and total**: no committed state ever violates it. Across a boundary it becomes **eventual and probabilistic**: it holds most of the time, violations are possible, and your job is to bound their frequency, detect them fast, and repair them cheaply. That reframing drives the design: ### Ownership Exactly one component is authoritative. Other components hold caches or projections and must render sensibly when the referenced thing is missing — an order UI that shows 'unknown product' beats one that throws. Write the ownership down; ambiguity here is what produces two half-enforcing services and no one accountable. ### Enforcement at the write path The owner validates before writing, and does so in a way that is race-free within its own database (constraint, unique index, or a serialized write). Cross-boundary existence checks are inherently stale — the parent could be deleted a millisecond later — so the goal is 'no writes we could have prevented', not 'zero violations by construction'. ### Propagation When one side's change must reach the other, share a transaction between the state change and the notification: write the event into an outbox table in the same transaction, and publish from there. This removes the 'wrote the row, crashed before publishing' class of drift. Consumers must be idempotent, because at-least-once delivery is the realistic guarantee. ### Detection Run continuous reconciliation — periodic anti-joins or checksum comparisons between the two sides — and expose the violation count as a metric with an alert threshold. The important cultural move: a nonzero orphan count is a normal, monitored quantity, not an emergency, until it crosses the threshold. Systems that never measure this discover the drift during a migration, years of bad data later. ### Repair Decide in advance which side wins and what remediation means: backfill the missing parent, delete or quarantine the orphan, or raise it to a human queue. Money and compliance data usually get a human queue; caches and denormalized copies get automatic repair. ### Deletion policy Most cross-boundary breakage originates in deletes. Prefer soft delete or a two-phase retire (mark inactive, wait out the propagation window, then remove) so consumers get a chance to react. A hard delete on the parent side with no notification is the classic orphan factory. ## The tradeoff to state out loud What you bought: independent scaling, independent deploys, blast-radius isolation, no cross-shard write coordination. What you paid: a class of correctness that used to be free is now an operational program — reconciliation jobs, metrics, repair runbooks, and a permanent design tax on every consumer that must handle missing references. Distributed transaction protocols (two-phase commit) can restore immediacy but couple availability across the boundary, which usually defeats the reason you split. An interview answer that names this trade honestly — and refuses to pretend an eventual invariant is the same as a declared one — is the differentiator. The weak answer says 'we just check it in the service'; the strong one says which invariants stayed engine-enforced by construction, which became eventual, who owns them, and how a violation is noticed within minutes rather than years.

  • Why not use two-phase commit to keep the invariant immediate across the boundary?
    Two-phase commit restores atomic cross-boundary writes but couples the availability of both sides: a coordinator or participant failure leaves locks held and can block progress, and every write pays extra round trips. That coupling usually defeats the reason the boundary exists. It stays defensible in a small, low-throughput, same-datacenter case with a mature coordinator, and is a poor fit for high-volume service-to-service work.
  • How do you keep reconciliation from becoming a job nobody looks at?
    Emit the violation count as a first-class metric with an alert threshold and an owner, the same way you treat error rate. Give it a runbook that says which side wins and how to repair, and review the trend in the same forum as other reliability metrics. A reconciliation job whose only output is a log line will be silently broken within a quarter.

Inside one engine an invariant is a lock on the door. Across a boundary it becomes a security patrol: you cannot prevent every breach, so you shorten the time between breach and discovery and rehearse the response.

saying these in an interview costs you the question

  • Accepting the loss of foreign keys without questioning the shard or service boundary
  • Claiming service-code checks give the same guarantee as a constraint
  • No reconciliation or metric — drift is invisible until a migration finds it
  • Reaching for two-phase commit without acknowledging the availability coupling
  • Hard-deleting parent rows with no propagation window or soft-delete step

context