Part VII — Advanced Topics and Applications
Chapter 26. Distributed Databases and Cloud Deployment
One server, one data directory — everything so far has lived on one machine, and that design carried the whole book. But real systems face the limits of one machine: the disk fills, the write rate exceeds one server's capacity, the machine fails and the university stops registering, and the datacenter itself is a single point of failure. Distributed databases answer by placing data on multiple machines, and the cloud answers by making those machines someone else's operational problem. This chapter is the systems view: replication, sharding, distributed transactions, failover, and the managed services that run PostgreSQL and MySQL for a living.
The chapter is deliberately placed near the end: it consumes the whole book. Transactions (Chapter 18), backup and recovery (Chapter 21), security (Chapter 20), and connection management (Chapter 23) all reappear, each gaining a second copy and a new failure mode.
After studying this chapter you will be able to:
- Explain why systems distribute data and what distribution costs.
- Compare replication modes and consistency guarantees.
- Design sharding schemes and know their operational price.
- Explain distributed transactions and the saga alternative.
- Build failover architectures and define RPO/RTO for them.
- Deploy PostgreSQL and MySQL in the cloud, self-managed and managed.
- Evaluate managed database services and their trade-offs.
- Secure distributed access: networks, TLS, and IAM.
- Reason about scaling, cost, and operational ownership.
26.1 Distributed database concepts
Distribution buys three things, each with a bill. Scale: one machine's CPU, memory, disk, and IOPS have ceilings; N machines multiply them. Availability: machines fail — replicas survive the failure (Chapter 21's RPO/RTO, made structural). Locality: data close to its users (a regional replica for each campus).
The bill has four lines. Latency: every network hop costs milliseconds — a distributed query is slower than a local one; only independent parallelism wins. Partial failure: some machines work while others do not, and the system must behave correctly in every subset — the defining hard problem. Consistency: copies must be managed — Chapter 21's async replication means brief windows of disagreement, and how a system explains those windows is its consistency model. Operational complexity: N machines, N versions, N clocks, N failure dashboards — Chapter 21's routine, multiplied.
The theory's one-line summary is the CAP theorem: under a network partition (P), a system must choose between consistency (C — all copies agree, always) and availability (A — every request answered) — the two cannot both hold through a partition. The engineering lesson is softer and more useful: choose per workload, per feature, on the spectrum between "strongly consistent, briefly unavailable" and "always answering, briefly stale" — Chapter 25's freshness-versus-load trade, at system scale.
26.2 Replication and consistency
Chapter 21 introduced replication as the administrator's insurance; the distributed view makes the modes precise:
- Synchronous: the primary waits for the replica's acknowledgment before confirming the commit — RPO zero (no committed data lost on failover), at the cost of every commit paying the slowest replica's round-trip. Used for the one standby that matters (PostgreSQL's
synchronous_standby_names; MySQL's semi-sync). - Asynchronous: commit confirms locally; the replica follows from the log — fast, but failover can lose the un-replicated tail (seconds of RPO). The default for scale-out replicas (Chapter 24's read-replica rung).
- Consistency models, the vocabulary: strong (a committed read is visible everywhere, immediately after the guarantee — synchronous or read-your-writes from the primary); eventual (replicas converge; a reader may briefly see yesterday's value — async replicas); read-your-writes (a user always sees their own changes — the practical middle: route a user's reads to the primary, or to the replica verified past their last write).
The operational disciplines: measure lag (Chapter 21's replication-lag metric — the honest number behind "how stale is the report"); choose the failover candidate deliberately (the sync standby, not the laggiest async); and remember the transaction identity problem — a promoted replica must not re-apply what the old primary already acknowledged (PostgreSQL's timelines; MySQL's GTIDs — Chapter 21's identifiers, doing exactly this work).
26.3 Partitioning and sharding
When reads outgrow replicas, writes need their own answer: sharding (horizontal partitioning across machines) — different rows on different servers, each shard a complete database for its subset (students 211000xx on shard A, 213000xx on B). Contrast with Chapter 25's table partitioning (one machine, many storage slices — pruning, not distribution).
The shard key is the decision that owns everything downstream:
- Hash sharding (
hash(student_id) % N) spreads load evenly; range queries cross every shard. - Range sharding (by ID prefix, date, or region) keeps related rows together; hot ranges concentrate load.
- The iron laws: the key must distribute load evenly; queries carrying the key touch one shard (fast), queries without it touch all shards (scatter-gather — the distributed full scan, the price of a bad key); and cross-shard joins and foreign keys degrade —
enrollmentandstudenton different shards cannot share a FK, which is why shard sets are usually co-located by the key (student_id shards both the student rows and their enrollments).
The three architectures: application-sharded (the application routes by key — full control, full responsibility), proxy/middleware-sharded (Vitess for MySQL — the proxy routes and merges; YouTube's provenance at planetary scale), and distributed-SQL engines (CockroachDB, YugabyteDB, Spanner-lineage — automatic sharding and distributed transactions, the Chapter 27 convergence story). The migration truth: sharding is a one-way door — design the key first, because changing it is a full data movement.
26.4 Distributed transactions
Chapter 18's transactions were single-server; a transaction touching two shards needs a distributed transaction, and the mechanisms are the theory's hardest pages:
- Two-phase commit (2PC) — the classic: a coordinator asks every participant to prepare (promise and durably hold the change), then commit or abort. Correct and slow — two network rounds while locks are held, plus the coordinator-failure limbo that recovery protocols must resolve. Both book platforms can be driven by external 2PC coordinators; PostgreSQL's
PREPARE TRANSACTIONis the hook. - Replication is the same problem, mirrored: the sync-replica commit of Section 26.2 is a two-party 2PC — the local commit and the replica's acknowledgment must both survive.
- The saga pattern — the practical alternative when 2PC is too heavy: split the distributed work into local transactions chained by compensations (book the seat locally; charge payment locally; if charging fails, run the compensating cancel-seat transaction). Sagas trade atomicity for availability — the intermediate states are visible, so the workflow must be designed for them (the airline seat-hold is the canonical saga).
The choosing rule: prefer single-shard transactions by design (co-locate by shard key — most "distributed transaction" problems dissolve at the modeling step); use 2PC where true atomicity is contractual; use sagas where the business already tolerates intermediate states.
26.5 High availability and failover
High availability is architecture, not a setting: enough redundancy that any single failure is a non-event, plus the machinery to route around failures. The stack, in layers:
- Redundant data: a sync standby (RPO 0) plus async replicas (read scale and disaster recovery).
- Redundant routing: a virtual IP or DNS-based service name that points clients at the healthy primary — clients connect to
university-db.example.edu, never to a hostname. - Automated failover: a quorum of witnesses (Patroni for PostgreSQL, InnoDB Cluster/Orchestrator for MySQL — Chapter 21's ecosystems) votes on the primary's health, promotes the best replica, and swings the virtual IP — in seconds, not page-response minutes.
- Split-brain discipline: the classic HA failure — two nodes each believing they are primary, both accepting writes, diverging irreconcilably. The cure is the quorum: promotion requires majority agreement, so a partitioned minority cannot promote (CAP's consistency choice, operationalized).
The design discipline is the RPO/RTO budget, stated before the architecture: the university registration system decides "lose zero committed enrollments, be down under a minute" — and that sentence is the architecture: sync standby, quorum, automated promotion. Chapter 21's numbers, doing their real job.
26.6 Cloud-hosted PostgreSQL and MySQL
The cloud offers both platforms at every level of management, and the levels are the decision:
- IaaS, self-managed: your PostgreSQL on a cloud VM. You keep everything from Chapter 21 (backups, patching, replication, monitoring) and gain elastic hardware and the provider's network. Right when you need full control (extensions, tuning, compliance) and have the operations team.
- Managed databases (next section): the provider runs the operational layer.
- Anything-as-a-service beyond: serverless endpoints, autoscaling clusters, distributed SQL — the Chapter 27 frontier.
Cloud PostgreSQL specifics worth knowing: the major providers (RDS/Aurora for PostgreSQL, Cloud SQL/AlloyDB, Azure Database) all offer version-current PostgreSQL with managed backups, PITR, multi-AZ sync replication, and read replicas; extensions are the usual constraint (PostGIS yes everywhere; niche extensions vary) — check before committing a schema to a provider. Cloud MySQL mirrors the story (RDS/Aurora MySQL, Cloud SQL), with Aurora's storage layer as the platform-specific twist (a distributed storage service under the MySQL-compatible engine — the write path scales differently than vanilla InnoDB). The self-managed chapter disciplines all carry over: the connection ladder, the security checklist, the monitoring thresholds — the cloud changes who executes the runbook, not the runbook.
26.7 Managed database services
Managed services — RDS, Cloud SQL, Azure Database, Aurora, AlloyDB — run the Chapter 21 calendar for you: automated backups with PITR, minor-version patching, replica management, failover on provider-detected failure, metrics and alarms wired to the console. What you trade for that is exactly what you would expect from this book's structure: operational control (maintenance windows, extension availability, superuser access — restricted; you are the tenant of a hardened postgresql.conf), performance ceilings (instance classes, IO credits — throttles that show up as the Chapter 19 symptoms), and migration gravity (managed snapshots are platform-shaped; logical dumps — Chapter 21's portable layer — remain your exit).
The decision, honestly: managed services win overwhelmingly for teams whose product is not the database — the operational floor (a tested backup, a surviving failover) exists on day one, at the price of the platform's decisions. Self-managed wins when the database is strategically shaped (extensions, extreme tuning, cost at scale, regulatory isolation) and the team exists to hold Chapter 21's calendar. The professional habit either way: keep the portability layer deliberate — logical backups, the Chapter 17 migration method rehearsed, no provider-proprietary features in the core schema without a written escape hatch.
26.8 Connection security and network access
Distribution multiplies the network surface, and Chapter 20's checklist scales with it. The cloud-native layering: the database lives in a VPC (a private network — Chapter 14's "ports private" rule, structural); access is governed by security groups (firewall rules by source — the application subnets may reach 5432; the internet may not, full stop); TLS is mandatory at these distances (verify-full; managed services ship certificates and often require encryption in transit); and IAM-based authentication joins password auth as a cloud pattern (cloud IAM identities issuing short-lived database credentials — Chapter 20's vault discipline, provided as a service; PostgreSQL's IAM auth on the managed platforms, MySQL's equivalent IAM plugin flows).
Two additions the distributed world forces: connection routing resilience — the application's Chapter 23 pool must survive a failover (point it at the reader/writer endpoints, the service names that re-resolve after promotion; a pool full of dead connections to the old primary is the classic post-failover incident, cured by pool health checks and pool_pre_ping); and cross-region discipline — replication between regions crosses the public internet (encrypted, verified) and drags the Chapter 20 audit questions with it (who can reach the DR replica? The answer must be "fewer than the primary").
26.9 Scaling and operational considerations
The closing synthesis — the scale ladder, in the order its rungs are actually climbed, each with its trigger:
- Vertical scaling (a bigger machine): the first rung — simple, downtime-adjacent, and the honest default until a measured ceiling (the Chapter 19 slow log and saturation metrics name the moment).
- Replicas for reads (Chapter 24's rung): write rate intact, read scale out.
- Caching (Redis/memcached at the application tier): removes reads entirely; introduces invalidation — the classic "cache and database disagree" bug class.
- Functional split (the reporting rung — Chapter 25's warehouse; heavy jobs to their own systems).
- Sharding for writes (Section 26.3): the last resort, designed around the key — the one-way door.
The operational bill, honestly stated: distributed systems need distributed observability (per-node Chapter 21 metrics plus lag, quorum health, and end-to-end probes); cost review (cloud costs are consumption curves — storage, IO, and instance-hours, each with its own surprise; the monthly review is a database feature); and rehearsed failure (the Chapter 21 restore drill becomes the failover drill: kill the primary on purpose, quarterly, and time the promotion — an HA architecture that has never failed over is a hypothesis).
The chapter's closing judgment, and Part VII's bridge to Chapter 27: distribution is not a feature you add; it is a set of trades you make when the single machine's measured limits arrive — and every trade has been this book's material all along: transactions, recovery, security, and the operational calendar, now with more than one copy.
Chapter Summary
- Distribution buys scale, availability, and locality; it costs latency, partial failure, consistency management, and operational multiplication — CAP names the partition-time choice between consistency and availability.
- Replication: synchronous (RPO 0, slower commits), asynchronous (fast, lag-tolerant), with strong/eventual/read-your-writes as the consistency vocabulary; lag is the honest freshness metric; timelines and GTIDs prevent re-application after promotion.
- Sharding: the key distributes rows across machines (hash vs range), key-carrying queries touch one shard, others scatter-gather, FKs need co-location; architectures are app-sharded, proxy-sharded (Vitess), and distributed SQL — and resharding is a full data movement.
- Distributed transactions: 2PC (prepare/commit, correct, slow), replication as two-party 2PC, sagas with compensations when atomicity is unaffordable — prefer single-shard design.
- HA: sync standby + async replicas + quorum-witnessed automated failover on a virtual service name; split-brain is cured by majority quorum; RPO/RTO stated first, architecture second.
- Cloud PostgreSQL/MySQL: IaaS self-managed keeps the Chapter 21 calendar; managed services run it (backups, patching, failover) at the price of control, ceilings, and gravity — keep the logical, portable exit.
- Security at distribution: VPC, security groups, mandatory TLS, IAM-issued credentials, failover-aware pools (endpoints + health checks), tighter DR access than primary.
- The scale ladder: vertical → replicas → cache → functional split → sharding, each by measured trigger; the operational bill is distributed observability, cost review, and rehearsed failure.
Key Terms
| Term | Definition |
|---|---|
| Distribution costs | Latency, partial failure, consistency, operational complexity |
| CAP theorem | Under partition, choose consistency or availability |
| Synchronous / asynchronous replication | Acknowledged commit (RPO 0) / local-first (lag window) |
| Strong / eventual / read-your-writes | Consistency models for replicas |
| Replication lag | The measured freshness of async copies |
| Timeline / GTID | Post-merge re-application guards (PG / MySQL) |
| Sharding (horizontal) | Rows split across machines by shard key |
| Shard key (hash / range) | The routing decision: even spread vs locality |
| Scatter-gather | Keyless queries touching all shards |
| Co-location | Shard-related tables together to keep FKs and joins local |
| Vitess / distributed SQL | Proxy sharding / automatic sharding-plus-transactions engines |
| Two-phase commit (2PC) | Prepare-then-commit distributed atomicity |
| Saga pattern | Local transactions + compensations |
| Quorum / split-brain | Majority promotion / dual-primary divergence |
| Virtual IP / service endpoint | The routed name clients actually connect to |
| RPO/RTO budget | The stated losses-and-downtime decision that is the architecture |
| VPC / security group | Private network / source-filtered firewall rules |
| IAM authentication | Short-lived cloud-issued database credentials |
| Pool failover health | Endpoints, health checks, pre-ping after promotion |
| Scale ladder | Vertical → replicas → cache → split → shard |
| Rehearsed failover | The quarterly deliberate-primary-kill drill |
Laboratory Exercises
- Build streaming replication (two PostgreSQL instances, or Docker/VMs): one primary, one async standby; verify data arrives (
SELECTon the standby), measure lag under write load, and promote the standby by the documented manual steps (the automation's manual rehearsal). Expected results: replication working; lag in milliseconds; promotion produces a writable new primary; timeline incremented. - Consistency in practice: with async replication, write on the primary and immediately read on the replica until it appears (record the iterations); repeat with
synchronous_commitsettings and compare behaviors. Expected results: eventual visibility observed (a brief window); the sync setting removes it at measured commit-cost. - Scatter-gather study: shard the Chapter 19 synthetic data by
student_id % 2across two dev databases; query by student (one shard) and by date range (both), and compare round-trips and plans. Expected results: key queries hit one shard; keyless queries visit both — the shard-key lesson, executed. - Saga walkthrough (paper and dev): the enroll-and-charge workflow as two local transactions with a compensation; run the happy path, then force the payment failure and verify the compensation cancels the seat — and note what a 2PC version would have guaranteed instead. Expected results: happy path commits both; failure path leaves the seat correctly cancelled; the visible intermediate state is the saga's documented trade.
- Failover drill: with your Laboratory 1 pair, kill the primary (deliberately); promote, repoint the Chapter 23 application's connection (endpoint or DNS), and time the whole sequence — then repeat the drill a second time and compare. Expected results: measured RTO in seconds; the application recovers by reconnecting to the new primary; the second drill faster — the rehearsal effect, in your own numbers.
- Managed-service evaluation: create the free/trial tier of one managed PostgreSQL or MySQL; exercise the operational features (automated backup, restore to a point, read replica, patching window); document what Chapter 21's calendar items you no longer run — and what you could no longer control. Expected results: a two-column ledger (provider-runs vs you-lose-control) from your own account, including at least one extension or superuser limitation.
Review Questions and Exercises
- Name the four costs of distribution, each with a chapter that already foreshadowed it. Latency (Chapter 4's network tiers); partial failure (Chapter 21's incidents); consistency (Chapter 21's replication lag); operations (Chapter 21's calendar — multiplied.
- State CAP and its practical engineering translation. Under partition, consistency or availability must give; in practice, choose per workload on the strong-vs-eventual spectrum, and state the choice.
- Which replication mode is right for the failover candidate, and why? Synchronous (RPO 0 — promotion loses nothing), while async replicas serve reads and disaster recovery.
- Why does the shard key decide both performance and schema? It routes queries (key-carrying = one shard) and must co-locate related tables for FKs and joins — a bad key taxes every query and breaks every constraint.
- Why are cross-shard foreign keys impossible, and what is the design response? The two machines cannot share an atomic constraint; co-locate related tables by the shard key (and let the application enforce what remains cross-shard).
- Contrast 2PC and the saga pattern in two sentences. 2PC atomically commits across participants — correct, slow, coordinator-fragile; sagas chain local transactions with compensations — always available, with visible intermediate states the workflow must own.
- How does a quorum prevent split-brain? Promotion requires majority agreement — a partitioned minority cannot win the vote, so only one primary can exist.
- Why must clients connect to a virtual endpoint rather than the primary's hostname? Failover re-points the endpoint; clients that hardcode a hostname keep dialing a dead machine — the classic post-failover outage.
- What three things do managed services run for you, and what three do they take? Run: backups/PITR, patching, failover/replication management; take: superuser/control, extension and tuning latitude, and easy exit (provider-shaped snapshots).
- Why is the connection pool a failover component, not just a performance one? A pool full of connections to the dead primary survives nothing; endpoint re-resolution, health checks, and pre-ping make the application itself failover-capable.
- Your read traffic has outgrown the primary. Walk the ladder to the first rung that solves it, and why not the next. Replicas — the read-heavy trigger; sharding is a write-rate solution with a one-way-door cost.
- Why rehearse failover quarterly even with automated tooling? The tool's assumptions drift (endpoints, credentials, versions); the drill measures real RTO and finds the broken link before the incident does — an untested failover is Chapter 21's untested backup.
Mini-Project
Design and document the availability architecture, HA_DESIGN.md, for the university registration system, then build its core in the lab: (1) the stated RPO/RTO budget with its business justification (one paragraph each — the sentences the architecture must satisfy); (2) the topology diagram — primary, sync standby, two async replicas (one cross-region), the service endpoints, the quorum witnesses — with each component's role labeled; (3) the build: primary + one standby (Laboratory 1's pair), with replication verified and lag measured; (4) the failover runbook — manual promotion steps (the automation's rehearsal), endpoint repointing, and the application's reconnect behavior — executed once, timed, with the before/after session behavior documented; (5) the security addendum: security-group matrix (who may reach each node), TLS verification, and the DR replica's stricter access list; (6) the managed-service comparison: the same architecture drawn twice — self-managed vs one managed service — with the Chapter 21 calendar mapped to "yours" versus "theirs," and the escape hatch stated; (7) the scale-ladder forecast: current metrics, the trigger thresholds for each rung, and the honest paragraph on when sharding would actually arrive. This design is the infrastructure half of Chapter 30's capstone deployment story.