You run one schema per customer organization, 3,000 schemas in total, and need to add a column plus backfill it. Describe how you would roll that change out across all tenants, and what is harder about it than making the same change to one shared set of tables.
answer
- expand/contract, never one breaking step
- registry table: tenant -> version, resumable
- one transaction per tenant, not per fleet
- canary internal and small tenants first
- reconcile afterwards to catch drift
basics
~20 sTreat it as an orchestrated job, not one DDL: backward-compatible expand/contract steps, per-tenant version tracking, idempotent and resumable execution, batched throttled backfills, canary tenants first, and application code that tolerates tenants at both old and new versions mid-rollout.
solid answer
~50 sWith 3,000 schemas the change is a **fleet operation**. - **Expand/contract.** Add the column nullable with no table rewrite, deploy code that writes it where present and never requires it, backfill, then tighten the constraint and drop the old path. Never one breaking step. - **Per-tenant state.** A registry table records each tenant's schema version; the runner is idempotent and resumable, so a crash at tenant 1,800 continues rather than restarts. - **Mixed versions are normal.** The rollout takes hours, so the deployed application must work against both shapes; that is the hard constraint, not the DDL. - **Canary and throttle.** Internal and small tenants first, then bounded-concurrency batches so backfills do not saturate I/O or blow out replication lag. One transaction per tenant, never one for the fleet. - **Verify.** Reconcile afterwards: every schema at the target version, no drift. Versus one shared table: a single DDL, but its blast radius is every tenant at once, no canary is possible, and the backfill is one enormous job.
code
sql · 14 linesCREATE TABLE tenant_registry (
tenant_id uuid PRIMARY KEY,
schema_name text NOT NULL,
schema_version int NOT NULL,
last_error text,
updated_at timestamptz NOT NULL
);
-- runner picks the next batch of laggards
SELECT tenant_id, schema_name
FROM tenant_registry
WHERE schema_version < 42
ORDER BY schema_version, tenant_id
LIMIT 50;go deeper
Recognise that the change must be applied per tenant by an automated job with progress tracking, not by hand.
Describe expand/contract, idempotent resumable execution, and batched backfills committed per tenant.
Add canary ordering, lock timeouts, throttling against replication lag, per-tenant failure reporting, and post-run drift reconciliation.
Treat migration cost as a property of the tenancy model: budget the mixed-version window into the release process, and weigh fan-out cost against the blast radius of a single shared-table DDL when choosing the layout.
## Why fan-out is a different problem With shared tables, a migration is one statement. With schema-per-tenant (or database-per-tenant) it becomes an orchestrated job over N targets that will be *partially applied* for a meaningful window. Everything below follows from that single fact: for hours, some tenants have the new shape and some do not. ## Make the change backward compatible (expand/contract) Because old and new shapes coexist, the change must be decomposed: 1. **Expand** - add the column as nullable with no volatile default, so the DDL is a fast catalog change rather than a table rewrite holding a strong lock. Create new indexes concurrently where the engine supports it. 2. **Deploy tolerant code** - the application writes the new column when it exists and never requires it. If it cannot tolerate absence, gate behaviour per tenant on the recorded schema version. 3. **Backfill** - in bounded batches, committing frequently, throttled to protect replication lag and I/O. 4. **Contract** - only once every tenant has reached the target version: add NOT NULL, drop the legacy column, remove the compatibility branch. The old-fashioned alternative - take the product down, run everything, come back - is defensible at 30 tenants and untenable at 3,000, because total time scales with tenant count. ## The runner Treat it as a job with durable state, not a shell loop: - **A tenant registry** holding, per tenant, its location and current schema version. It is the source of truth for progress, and the same directory that later lets you move a tenant. - **Idempotency.** Every step must be safe to re-run: conditional DDL (add-if-missing), backfills driven by a watermark or a "not yet migrated" predicate. A failure at tenant 1,800 must resume, not restart. - **Bounded concurrency.** Run K tenants in parallel, tuned so the instance is not saturated; a single-threaded loop may not finish in a maintenance window, while unbounded parallelism turns the database over. - **One transaction per tenant**, never one transaction for the fleet: a fleet-wide transaction holds locks for hours, bloats undo and write-ahead log, and cannot make partial progress. - **Timeouts and lock guards.** A short lock timeout means a DDL waiting behind a long-running transaction fails that tenant fast and retries later, instead of queueing every subsequent query behind it. - **Canary order.** Internal tenants, then small ones, then the whales, whose tables are largest and riskiest. This canary ability is the one genuine advantage over the shared-table model. - **Reconciliation.** Afterwards, diff every schema against the target definition. Drift - a tenant that failed silently, or one hand-patched during an incident - is the chronic failure mode of this layout, and only automated comparison catches it. ## What else bites at 3,000 schemas - **Catalog pressure:** 3,000 x (tables + indexes + constraints) objects slow catalog lookups, plan caching and dump/restore, and multiply background maintenance work. Every migration adds to that count. - **Deploy duration:** even one second per tenant is nearly an hour; a rewrite-heavy change can span days, so the mixed-version window becomes a design input rather than an edge case. - **Observability:** you need per-tenant success and failure reporting, or one failed tenant is discovered by a customer. - **Connection churn:** the runner opens sessions across many schemas; reuse pooled connections and select the target schema per transaction. ## Contrast with the shared-table model One DDL and one backfill - simpler orchestration, but: the blast radius is every tenant simultaneously with no canary; the backfill is a single job over a table containing every tenant's rows, so it must be batched and throttled anyway; and a long-lock DDL stalls the whole product rather than one customer. Database-per-tenant behaves like schema-per-tenant plus extra deployment plumbing (many endpoints, credentials, possibly regions), but without the shared-catalog pressure and with genuinely independent per-tenant schedules. ## The rule to state in an interview Schema changes in a fan-out layout must be **backward compatible, idempotent, resumable, per-tenant tracked, throttled, and reconciled** - and the application must stay correct while the fleet is half migrated.
- Half the tenants have the new column and half do not. What does that require of the deployed application?It must run correctly against both shapes for the whole rollout window - writing the new column only where it exists and never depending on it for reads - or gate behaviour per tenant on the recorded schema version. That is exactly why the change is split into expand, backfill and contract phases, with constraint tightening deferred until every tenant is migrated and verified.
- Why not wrap the whole fan-out in a single transaction so it is atomic?It would hold DDL locks across thousands of objects for hours, block ordinary traffic, accumulate enormous undo and log volume, and lose all progress on any failure. Per-tenant transactions give partial progress, bounded lock duration and the ability to resume. Atomicity across tenants is not needed, because tenants are independent - which is the whole point of the layout.
saying these in an interview costs you the question
- Running the fan-out as a non-resumable script with no per-tenant state
- Assuming the application only ever sees one schema version at a time
- Adding a NOT NULL column with a default in the same step on large tenant tables
- Unbounded parallelism that saturates I/O or blows out replication lag
- Hand-patching one tenant during an incident and never reconciling the drift