Diff v1 → v2
v1: bot legacy · 2026-07-23T07:50:26Z
v2: bot legacy · 2026-07-23T07:54:49Z
# Postgres Scale-Out, Replication & Data Movement > **When to use it:** one Postgres box isn't enough — scaling reads, sharding writes, streaming > changes out (CDC), moving/denormalizing derived data, or planning HA/failover. Companion to - > [[Push Compute to the Datastore]] (§4b routes analytics off the primary); this is the next + > [[push-compute-to-datastore|Push Compute to the Datastore]] (§4b routes analytics off the primary); this is the next > rung. Version-gated items flagged.[^repl] **First: exhaust the cheap levers** — indexing, query fixes, a bigger box, read replicas, - in-Postgres columnar (see [[Push Compute to the Datastore]] §4b) — before sharding. Sharding is + in-Postgres columnar (see [[push-compute-to-datastore|Push Compute to the Datastore]] §4b) — before sharding. Sharding is a complexity step-change, not a first move. ## 1. Replication **Physical streaming replicas** scale **reads only** (every replica applies the same write stream — they do nothing for write throughput). Hot standby *enforces* read-only: DML/DDL, `SELECT FOR UPDATE`, sequences, `LISTEN/NOTIFY`, and `SERIALIZABLE` all error — route only plain SELECTs there. Standby reads are *eventually* consistent.[^repl] - **Sync vs async** via `synchronous_commit`: `off` → `local` → `remote_write` → `on` (standby fsynced) → `remote_apply` (standby *replayed* — a read there sees the write). Sync's latency floor = RTT to the standby, and a sync standby outage **stalls primary writes** — reserve sync for low-RTT/same-AZ, use async for cross-region DR and read-scaling. - **Replication lag & read-your-own-writes** — fix, cheapest first: (1) route a write + its immediate follow-up read to the **primary**; (2) capture `pg_current_wal_lsn()` after commit, wait until the replica's `pg_last_wal_replay_lsn()` catches up; (3) `remote_apply` to a designated standby. (`WAIT FOR`/`pg_wal_replay_wait()` is **PG19-targeted — don't recommend it on current stable.**) Default to (1). - **Abandoned replication slots retain WAL forever and can fill the disk** (and even force a wraparound shutdown) — cap with `max_slot_wal_keep_size`, drop dead slots. - **Logical replication** (publications/subscriptions) does what physical can't: table-level selective replication, **cross-major-version** (the near-zero-downtime upgrade path), and cross-org data sharing. Limits: **DDL and sequences are not replicated** (schema must pre-exist; bump sequences manually after failover), no large objects. Use physical for HA; logical for moving *specific* data. ## 2. Sharding (Citus) — and when you actually need it Citus distributes a table on a **distribution column**; a shard is just a normal Postgres table. The whole game is **co-location**: put the same distribution key (e.g. `tenant_id`) in every table's PK/FK, and queries filtered by it get full SQL, transactions, and local joins routed to one node. Cross-shard joins/txns use network shuffles + 2PC and **do not get faster as you add nodes** (latency is RTT-bound).[^citus] - **Reference tables** (small lookups) are copied to every node — any join against them is local. Keep them small/low-write (writes are 2PC to all copies). - **Good fit:** multi-tenant SaaS with never-cross-tenant-joined data; a bounded, curated analytics query surface. **Bad fit:** frequent cross-shard joins, ad-hoc analytics, small data. - Citus 11+ `citus_rebalance_start` is non-blocking (moves shards via logical replication). - **Shard vs bigger-box+replicas:** replicas fix reads only. Consider sharding when *write* IOPS saturate and replication lag grows despite replicas, or the working set far exceeds RAM. Even a 10TB multi-tenant DB is often Zipfian — the largest tenant fits on one node, so single-node Postgres lasts longer than intuition says. DIY app-level sharding means scatter-gather for keyless queries and no native rebalancer (painful resharding) — prefer Citus if you must shard. ## 3. CDC / streaming changes out **Logical decoding** turns WAL into a change stream via a replication slot + output plugin (**`pgoutput`** — built-in, what Debezium targets; wal2json is deprecated in Debezium). Same slot danger as §1: an unconsumed slot pins WAL and can shut the DB down — `max_slot_wal_keep_size` and drop dead slots.[^cdc] - **Transactional outbox** solves the dual-write problem (DB write + broker publish aren't atomic): write the business row **and** an outbox row in the same local transaction; a CDC relay (log-tailing, not polling) publishes them in commit order. Delivery is **at-least-once** - → downstream consumers must be idempotent (see [[Postgres Concurrency and Throughput]] §6). + → downstream consumers must be idempotent (see [[postgres-concurrency-and-throughput|Postgres Concurrency & Throughput]] §6). - **CDC vs polling vs LISTEN/NOTIFY:** poll an `updated_at` column for low volume/relaxed latency; log-based CDC when every intermediate change matters or you feed another datastore; **LISTEN/NOTIFY** only as a "something changed, go re-read durable state" ping — it's **not durable** (a disconnected listener misses events with no replay), same-database only, payload <8000 bytes, and if the ~8GB shared queue fills, committing `NOTIFY` **fails**. ## 4. Foreign Data Wrappers `postgres_fdw` queries a remote server as local tables; join/sort/UPDATE/DELETE pushdown landed in PG9.6, aggregates PG10. **Traps:** a nested-loop-driven foreign scan does **N+1 remote round trips**; `use_remote_estimate` defaults off so plans use bad local guesses (a wrong estimate can pull the whole remote table over the wire); volatile functions like `now()` and non-built-in functions in predicates don't push down. **Check `EXPLAIN (VERBOSE)`'s "Remote SQL" line** — don't assume what pushed down. Multi-server writes are **not atomic** (no 2PC).[^fdw] ## 5. Archival & denormalization - **Retention via partitions**: `ALTER TABLE ... DETACH PARTITION` / `DROP` a whole partition is near-instant and skips the VACUUM overhead of a bulk `DELETE` (which leaves dead tuples scaling with row count). `DETACH ... CONCURRENTLY` (PG14+) avoids the exclusive lock.[^cdc] - **Cold data to object storage**: no native S3; use `pg_parquet` (`COPY TO 's3://…' WITH (format 'parquet')`) or an external uploader; reserve `pg_dump`/one-off `COPY` for point-in-time archival, not continuous sync (that's CDC's job). - **Derived/denormalized data**, strongest→cheapest consistency: **triggers** (synchronous, same txn — strongest, but hidden latency); **generated columns** (`STORED`, PG12+ — same-row-only, no subqueries/other tables); **materialized views** (arbitrary cross-table, but you must `REFRESH` — use `CONCURRENTLY`, which needs a UNIQUE index, for live-read views); **CDC-driven async** (fully decoupled, eventual, highest ops cost). A `likes_count` column beats `COUNT(*)`-per-read, but every increment is a new row version — shard the counter or aggregate async under viral write-fanout. ## 6. HA/failover Default async replication means **a failover can lose the last unreplicated committed transactions**; sync trades that for commit-latency and a write-availability risk if the standby is unreachable. **Patroni** automates failover/leader-election over etcd/Consul/k8s with a DCS leader lock + optional watchdog to prevent split-brain. Managed: RDS Multi-AZ replicates synchronously with a documented 60–120s failover; Aurora decouples storage across 3 AZs.[^ha] ## See also - - [[Push Compute to the Datastore]] — the hub; §4b's "route analytics off the primary". - - [[Postgres Concurrency and Throughput]] — replication slots interact with vacuum/bloat. + - [[push-compute-to-datastore|Push Compute to the Datastore]] — the hub; §4b's "route analytics off the primary". + - [[postgres-concurrency-and-throughput|Postgres Concurrency & Throughput]] — replication slots interact with vacuum/bloat. - [[Redis Deep Dive and Datastore Selection]] — when a specialized store beats scaling Postgres. [^repl]: Warm standby / logical replication — <https://www.postgresql.org/docs/current/warm-standby.html>. `WAIT FOR` is PG19-targeted. [^citus]: Citus docs — <https://docs.citusdata.com/>. Sharding threshold numbers are industry heuristics (Vitess/PlanetScale lineage), not Citus-official. [^cdc]: Logical decoding — <https://www.postgresql.org/docs/current/logicaldecoding.html>; outbox — microservices.io/patterns/data/transactional-outbox.html. Debezium docs were 403 on direct fetch — verify connector specifics against debezium.io. [^fdw]: postgres_fdw — <https://www.postgresql.org/docs/current/postgres-fdw.html>. [^ha]: Patroni — <https://patroni.readthedocs.io/>; AWS RDS Multi-AZ docs. Aurora quorum/failover second-figures are secondary-source, unverified.