Skip to content

第 15 章 复制与高可用:让 PostgreSQL 永不掉线

目标读者:已经能独立运维一个单机 PG 实例,但不知道「主库挂了怎么办」「读 QPS 上不去怎么办」「跨机房灾备怎么搭」的同学。

学完你会:能搭出一套主从流复制 + 自动 Failover 的高可用集群,能用逻辑复制做不停机版本升级 / CDC 数据管道,能看懂 Patroni / repmgr / pg_auto_failover 的工作原理。


0. 导读:为什么要复制?

单台 PG 跑得好好的,为什么要拷贝一份?三个最常见的理由:

  1. 灾备(Disaster Recovery):主机房光纤被挖断、机器硬盘炸了、运维手抖 DROP DATABASE,你需要一个「随时可顶上去」的备份实例。
  2. 读扩展(Read Scaling):报表 / BI / 搜索这种重读轻写的业务,把流量打到只读从库,主库专心写。
  3. 零停机维护(Zero Downtime Upgrade):升级大版本(比如 PG 14 → PG 17)时,先把数据通过逻辑复制同步到新版本,再做切换。

PostgreSQL 把「复制」分成两大体系:

体系别名复制单元主从是否可异构典型场景
物理复制流复制 / 块级复制WAL 字节流(页面级)❌ 主从必须字节级一致主备灾备、读写分离
逻辑复制publication/subscription逻辑事件(行级 INSERT/UPDATE/DELETE)✅ 可跨大版本、跨表结构跨版本升级、CDC、双向同步

📌 与 MySQL 的区别

  • MySQL 主从复制基于 binlog(逻辑日志),原生就是「逻辑复制」;
  • PostgreSQL 默认走物理流复制(基于 WAL 字节流),从 PG 10 起才内置逻辑复制;
  • PG 物理复制几乎是「字节复印机」,从库连页面布局都跟主库一模一样;
  • MySQL 的「半同步」(rpl_semi_sync_master_enabled) 大致对应 PG 的 synchronous_commit = remote_write

1. PG 复制全景图

先一张图把整个体系印在脑子里:

                    ┌────────────────────────────────────────────────┐
                    │              PostgreSQL 复制体系               │
                    └─────────────────────┬──────────────────────────┘

              ┌───────────────────────────┴────────────────────────────┐
              │                                                        │
        ┌─────▼──────┐                                          ┌──────▼──────┐
        │ 物理复制   │                                          │  逻辑复制   │
        │ (基于 WAL) │                                          │ (基于解码)  │
        └─────┬──────┘                                          └──────┬──────┘
              │                                                        │
   ┌──────────┼────────────┐                          ┌────────────────┼──────────────┐
   │          │            │                          │                │              │
┌──▼───┐  ┌──▼────┐  ┌─────▼───────┐         ┌────────▼─────┐  ┌───────▼──────┐  ┌────▼────────┐
│ 文件 │  │ 流复制│  │ Hot Standby │         │ Publication  │  │ Subscription │  │pg_recvlogical│
│ 复制 │  │ Stream│  │ (从库可读)  │         │   (主库)     │  │   (从库)     │  │ (CDC 工具)   │
└──┬───┘  └───┬───┘  └─────────────┘         └──────────────┘  └──────────────┘  └─────────────┘
   │          │
   │  walsender ──→ 网络 ──→ walreceiver
   │  (主库)             (从库)

归档式:pg_wal/ 文件归档到共享存储(NFS/S3),从库 restore_command 拉取

简单说:

  • 物理 + 文件复制:主库把写满的 WAL 文件归档到共享存储,从库定期拉取重放。慢、有延迟,已逐渐被流复制取代,但 PITR(时间点恢复)还在用。
  • 物理 + 流复制(Streaming Replication):主库的 walsender 进程实时把 WAL 字节流推给从库的 walreceiver这是当下最主流的「主备」方案。
  • 物理 + Hot Standby:从库一边 redo WAL,一边对外提供只读查询。这个能力默认是开的(hot_standby = on)。
  • 逻辑复制:基于「逻辑解码 Logical Decoding」把 WAL 解析成行级变更事件,发给订阅端。可以跨版本、按表选择、甚至双向同步。

2. 物理流复制(Streaming Replication)原理

2.1 进程角色

主库和从库各起一组进程,它们的协作流程是流复制的「灵魂」:

                       ┌──────────── 主库 (Primary) ─────────────┐
   客户端 INSERT ──→   │                                          │
                       │   backend ──写─→ WAL Buffer ──flush─→  │
                       │                       │                 │
                       │                       ↓                 │
                       │                  pg_wal/0000001...      │
                       │                       │                 │
                       │   walwriter (异步刷盘)│                 │
                       │                       │                 │
                       │                  walsender ─────┐       │
                       └────────────────────────────────│────────┘

                                                  TCP/IP 流

                       ┌────────────────────────────────│────────┐
                       │   walreceiver  ←──────────────┘         │
                       │       │                                 │
                       │       ↓ 写入                             │
                       │   pg_wal/0000001...                     │
                       │       │                                 │
                       │   startup 进程持续 redo                 │
                       │       │                                 │
                       │       ↓ 应用                             │
                       │   shared_buffers(可被只读查询读到)    │
                       │                                         │
                       │   只读 backend ←—— 客户端 SELECT        │
                       └──────────── 从库 (Standby) ─────────────┘

逐步拆解:

  1. 客户端 commit:在主库,客户端发起 COMMIT,对应的 backend 把这次事务的修改写到 WAL Buffer,再 flush 到 pg_wal/ 下的 WAL 段文件。
  2. walsender 推送:每个连接上来的从库,主库都会 fork 一个 walsender 进程,从 pg_wal/ 里读 WAL 字节流,通过 TCP 推给从库。
  3. walreceiver 接收:从库的 walreceiver 进程负责接收 WAL,落到从库本地 pg_wal/,并通知 startup 进程。
  4. startup 进程 redo:从库的 startup 进程不断把 WAL 应用到从库自己的数据文件 + shared_buffers。
  5. Hot Standby 提供只读:因为应用到了 shared_buffers,从库的只读 backend 可以直接 SELECT

2.2 一个 commit 在主从的完整时间轴

时间 ──→
主库:   T1 backend写WAL ──→ T2 commit返回 ──→ T3 walsender推送

从库:                                            T4 walreceiver收 ──→ T5 落盘 ──→ T6 startup redo ──→ T7 可见

「主从延迟」就是 T7 - T2,可能因为:

  • 网络慢(T3→T4 慢)
  • 从库磁盘慢(T5 慢)
  • 从库重放慢(单进程 redo,遇到大事务卡死)
  • 从库被长查询挡住(hot_standby_feedback

2.3 同步级别:synchronous_commit

主库的 backend 在 commit 时到底等不等从库?由参数 synchronous_commit 控制:

含义性能数据安全
off不等 fsync,不等从库最快单机 crash 可能丢最后几百毫秒
local主库 fsync,不等从库单机安全,从库可能落后
on(默认)主库 fsync + 等所有同步从库收到WAL主库挂了,被同步的从库一定有这条 WAL
remote_write等同步从库写入OS 缓存(未必落盘)中-慢从库 OS crash 可能丢
remote_flush等同步从库fsync 落盘强一致
remote_apply等同步从库redo 完毕最慢切换无延迟读

synchronous_standby_names 决定哪几个从库参与同步

ini
# 类型 1:FIRST k (s1, s2, s3) —— 优先级模式:列表中前 k 个可用从库参与同步
synchronous_standby_names = 'FIRST 1 (replica1, replica2)'

# 类型 2:ANY k (s1, s2, s3) —— 法定人数模式(quorum):任意 k 个 ACK 就算成功
synchronous_standby_names = 'ANY 2 (replica1, replica2, replica3)'

Quorum 复制(法定人数复制):要求 N 台从库里至少 K 台 ACK,提供更高可用性。例如 3 台从库要求 ANY 2,则任意一台从库挂了仍能提交事务,不会卡主库。

📌 与 MySQL 的区别

  • MySQL 半同步只有「至少一个从库 ACK」一种模式;
  • PG 的 quorum 模式可以精确指定 K,灵活度高。

3. 流复制配置完整步骤(PG 12+ 新机制)

PG 12 之前用 recovery.conf 做从库配置,PG 12 起这个文件被废弃,改用「主配置文件 + standby.signal 空文件」标识从库身份。

3.1 主库配置

编辑 postgresql.conf

ini
# 必须是 replica 或 logical(要做逻辑复制就上 logical)
wal_level = replica

# 最多允许多少个 walsender 同时连进来(= 从库数 + 备份工具数 + 一点 buffer)
max_wal_senders = 10

# 至少保留多少 WAL 段在 pg_wal/ 不删(防止从库还没拉就被覆盖;推荐用 replication slot 替代)
wal_keep_size = 1024MB        # PG 13+;老版本用 wal_keep_segments

# 允许从库连过来同步 commit
synchronous_commit = on
synchronous_standby_names = 'FIRST 1 (standby1)'   # 可选,不设就是异步

编辑 pg_hba.conf,加 replication 专用行:

# TYPE  DATABASE      USER         ADDRESS          METHOD
host    replication   repl_user    192.168.1.0/24   scram-sha-256

注意:DATABASE 字段必须是字面量 replication,这是 PG 给复制连接保留的「假数据库名」。

创建复制专用账号:

sql
CREATE ROLE repl_user WITH REPLICATION LOGIN PASSWORD 's3cret';

重启主库(修改 wal_level 必须重启,其他参数 pg_reload_conf() 即可)。

3.2 从库配置:用 pg_basebackup 一键完成

pg_basebackup 是 PG 自带的物理基础备份工具,能:

  1. 把主库的整个 data 目录拷贝过来;
  2. 同步流式拷 WAL;
  3. -R 参数自动生成从库所需的两个文件:standby.signal + postgresql.auto.conf 里的 primary_conninfo

完整命令:

bash
# 在从库机器上执行
sudo -u postgres pg_basebackup \
    -h 主库IP -p 5432 -U repl_user \
    -D /var/lib/postgresql/data \
    -Fp -Xs -P -R \
    -S ch15_standby1_slot     # 可选:绑定一个复制槽

参数解读:

参数含义
-D输出目录(必须为空)
-Fpplain 格式(不打 tar)
-Xsstream 模式获取 WAL(-Xf 是 fetch 一次性拉)
-P显示进度
-Rrecovery 模式:自动写 standby.signalprimary_conninfo
-S使用复制槽

完成后启动从库:

bash
sudo -u postgres pg_ctl -D /var/lib/postgresql/data start

3.3 主库验证从库连上来了

sql
SELECT
    application_name,
    client_addr,
    state,                  -- streaming 表示正常
    sync_state,             -- async / potential / sync / quorum
    pg_wal_lsn_diff(sent_lsn, replay_lsn) AS lag_bytes
FROM pg_stat_replication;

实操示例输出:

 application_name | client_addr  |   state   | sync_state | lag_bytes
------------------+--------------+-----------+------------+-----------
 standby1         | 10.0.0.21    | streaming | sync       |       0
 standby2         | 10.0.0.22    | streaming | async      |    8192

3.4 从库可读吗?

sql
-- 在从库执行
SHOW transaction_read_only;     -- on,从库不能写
SELECT pg_is_in_recovery();     -- t,处于 recovery 模式
SELECT * FROM users LIMIT 10;   -- 可以读!

如果尝试在从库 INSERT

ERROR:  cannot execute INSERT in a read-only transaction

4. 复制槽 Replication Slot

4.1 为什么需要复制槽?

想象一个场景:从库网络断了 1 小时,期间主库一直在产生 WAL。这些 WAL 文件如果被 wal_keep_size 限制清掉了,从库再连回来时就追不上了,只能重新做一次 base backup。

**Replication Slot(复制槽)**就是给主库装的「书签」:

主库一旦发现某个 slot 的 restart_lsn 还没推到从库,就坚决不删对应的 WAL,直到从库消费完。

类比:你订了一份报纸,邮递员不能把今天的报扔掉,必须等你看完才能回收。

4.2 物理槽(给流复制用)

主库创建:

sql
SELECT pg_create_physical_replication_slot('ch15_standby1_slot');

SELECT slot_name, slot_type, active, restart_lsn FROM pg_replication_slots;

从库的 postgresql.auto.conf 里配置:

ini
primary_conninfo = 'host=主库IP port=5432 user=repl_user password=s3cret'
primary_slot_name = 'ch15_standby1_slot'

或在 pg_basebackup 时直接 -S ch15_standby1_slot --create-slot

4.3 风险:从库长期不消费 → 主库 pg_wal 爆盘

这是物理复制最经典的运维事故!

某天监控告警:「主库磁盘 95%」。一查发现 pg_wal/ 目录占了 500GB,而正常应该几百 MB。原因:

  • 某个从库下线/出问题了 1 周;
  • 复制槽不允许主库回收 WAL;
  • 主库 WAL 越攒越多,最后磁盘炸了,主库也挂了。

防御措施

ini
# PG 13+:限制单个 slot 最多保留多少 WAL,超了就强制清理
max_slot_wal_keep_size = 64GB

监控告警语句:

sql
SELECT
    slot_name,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots
WHERE active = false OR
      pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn) > 10737418240;  -- > 10GB

4.4 逻辑槽(给逻辑复制用)

逻辑槽创建时要指定 decoding 插件

sql
SELECT pg_create_logical_replication_slot('cdc_slot', 'pgoutput');
-- pgoutput 是 PG 内置的逻辑复制输出插件
-- 也可以用 wal2json 输出 JSON 格式

5. 级联复制(Cascading Replication)

从库 1 → 从库 2 → 从库 3,挂成一条链:

   ┌─────────┐
   │ Primary │
   └────┬────┘
        │ stream

   ┌─────────┐
   │Standby A│
   └────┬────┘
        │ stream(A 既是 standby 又是 sender)

   ┌─────────┐
   │Standby B│
   └─────────┘

好处:减轻主库 walsender 压力(多机房场景,每个机房只有一台从库连主库)。

代价:B 的延迟 = A 的延迟 + B 自己的延迟。

配置只需让 B 把 primary_conninfo 指向 A 即可。


6. 故障切换 Failover vs Switchover

操作中文触发原因主库状态
Switchover主备切换主动操作(升级、迁移)主库健康,主动让位
Failover故障切换被动操作(主库挂了)主库已死 / 失联

6.1 手动切换:pg_promote()

PG 12 起,从库可以用 SQL 函数提升为主库:

sql
-- 在从库执行
SELECT pg_promote(wait => true, wait_seconds => 60);
-- 返回 t 表示提升成功

或命令行:

bash
pg_ctl promote -D /var/lib/postgresql/data

提升后从库会:

  1. pg_wal/ 里所有未应用的 WAL redo 完;
  2. 删掉 standby.signal 文件;
  3. 写一条「Time Line ID 切换」记录,进入新的 timeline;
  4. 开始接受写入。

6.2 时间线 Timeline

每次 promote 都会让 PG 进入新 timeline,pg_wal/ 文件名变化:

原 timeline 1:000000010000000000000050
新 timeline 2:000000020000000000000050.partial   ← 旧 timeline 残留
新 timeline 2:000000020000000000000051           ← 新 timeline 写入

这是 PITR(时间点恢复)能精确回到任意时刻的基础。

6.3 自动切换需要外部组件

PG 内核故意不实现自动 Failover!原因是「自动判断主库死活」涉及网络分区、脑裂、仲裁等分布式难题,应该交给专业组件:

组件厂商 / 出品选举依赖特点
PatroniZalandoetcd / Consul / ZooKeeper最流行,K8s 生态首选
repmgr2ndQuadrant(已被 EnterpriseDB 收购)自带或外部老牌、轻量
pg_auto_failoverMicrosoft (CitusData)自带 monitor 节点配置简单,依赖 monitor
StolonSorintlabetcd / ConsulGo 写的,K8s 友好
PgPool-IINTT OSS Center自带 watchdog兼连接池 + 负载均衡

📌 与 MySQL 的区别

  • MySQL 8 的 InnoDB Cluster + MySQL Router + MySQL Shell 是官方一体化方案;
  • PG 走「内核只做基础设施 + 社区做编排」的路线,所以 Patroni 这类外部组件不可少。

6.4 Patroni 工作原理(简版)

                ┌──────────────────────────┐
                │   etcd (分布式 KV)       │
                │   /patroni/cluster/leader│
                └─────────┬────────────────┘
                          │(leader key 带 TTL,不续约就失效)
       ┌──────────────────┼──────────────────┐
       ↓                  ↓                  ↓
   ┌─────────┐        ┌─────────┐        ┌─────────┐
   │Patroni A│        │Patroni B│        │Patroni C│
   └────┬────┘        └────┬────┘        └────┬────┘
        │                  │                  │
   ┌────▼─────┐       ┌────▼─────┐       ┌────▼─────┐
   │ PG 主库  │       │ PG 从库  │       │ PG 从库  │
   └──────────┘       └──────────┘       └──────────┘

流程:

  1. 启动时所有节点抢 etcd 里的 leader key,谁先抢到就是主;
  2. 主节点定期续约 leader key(TTL 默认 30s);
  3. 主节点挂了 → leader key 过期 → 其他节点发现 → 重新选举;
  4. 新选出来的节点(一般是 LSN 最大的)自动 pg_promote(),并修改其他从库的 primary_conninfo 指向自己。

7. 逻辑复制(PG 10+)

7.1 与物理复制的根本区别

维度物理复制逻辑复制
复制单元WAL 字节流(页面级)行级变更事件(INSERT/UPDATE/DELETE)
主从版本要求必须完全一致可跨大版本(10 → 17)
主从表结构要求必须完全一致仅订阅表结构必须一致
从库可写✅ 订阅端可独立写其他表
双向复制✅(PG 16+ 内置)
DDL 复制❌(PG 16/17 部分支持)
复制粒度整个集群按表 / 按行(带 WHERE)

7.2 逻辑解码 Logical Decoding

逻辑解码:把 WAL 里以「页面变更」记录的内容,解析回「这是对 users 表的 INSERT,值是 (1, 'Alice')」这种逻辑事件。

PG 内置的输出插件 pgoutput 输出 PG 自己的逻辑复制协议;社区插件 wal2json 输出 JSON。

7.3 配置三步走

主库(发布者,Publisher):

ini
# postgresql.conf
wal_level = logical          # 必须升到 logical(默认 replica 不行)
max_wal_senders = 10
max_replication_slots = 10
sql
-- 创建发布
CREATE PUBLICATION ch15_pub_orders FOR TABLE ch15_orders, ch15_order_items;
-- 也可以全表
CREATE PUBLICATION ch15_pub_all FOR ALL TABLES;
-- 也可以带过滤条件(PG 15+)
CREATE PUBLICATION ch15_pub_paid FOR TABLE ch15_orders WHERE (status = 'paid');

从库(订阅者,Subscriber):

sql
-- 1. 先在从库上把表结构建好(DDL 不会自动复制!)
CREATE TABLE ch15_orders (...);
CREATE TABLE ch15_order_items (...);

-- 2. 创建订阅
CREATE SUBSCRIPTION ch15_sub_orders
    CONNECTION 'host=主库 port=5432 dbname=learn_pg user=repl_user password=s3cret'
    PUBLICATION ch15_pub_orders;

自动发生的事情

  1. 从库自动在主库创建一个逻辑复制槽(名字默认与订阅同名);
  2. 主库做一次初始拷贝(COPY 全表数据到从库);
  3. 之后每次主库的 INSERT/UPDATE/DELETE 都会通过逻辑解码推给从库。

7.4 验证

主库:

sql
SELECT * FROM pg_publication;
SELECT * FROM pg_replication_slots WHERE slot_type = 'logical';
SELECT * FROM pg_stat_replication;

从库:

sql
SELECT * FROM pg_subscription;
SELECT * FROM pg_stat_subscription;

7.5 限制与坑

  1. DDL 不复制:你在主库 ALTER TABLE ch15_orders ADD COLUMN xxx,从库不会自动跟进,需要手动同步。
  2. TRUNCATE 复制需要 PG 11+(早期版本完全不复制)。
  3. 大事务延迟:主库一个 1 亿行的 UPDATE,逻辑解码需要全量读完再发,从库迟迟拿不到(PG 14+ 有 streaming 选项缓解)。
  4. 主键依赖:UPDATE/DELETE 需要主键(或 REPLICA IDENTITY FULL),否则报错。
  5. 序列不复制SERIAL 列的 next value 不会同步。

7.6 应用:跨版本平滑升级

PG 14 → PG 17 完全不停机升级流程:

  1. 起一台 PG 17 实例(空库);
  2. 从 PG 14 上 pg_dump --schema-only 把表结构同步到 PG 17;
  3. 在 PG 14 上 CREATE PUBLICATION p_all FOR ALL TABLES;
  4. 在 PG 17 上 CREATE SUBSCRIPTION s_all CONNECTION '...' PUBLICATION p_all;
  5. 等延迟追平(监控 pg_stat_subscription);
  6. 业务代码切换连接串到 PG 17;
  7. 删除 PG 14 上的发布,下线老库。

7.7 应用:CDC 数据管道

逻辑复制 + 第三方消费工具,可以把 PG 变成一个 CDC(Change Data Capture)源,把数据变更实时打到 Kafka / ES / Hudi:

  • Debezium:开源 CDC 平台,PG connector 用 pgoutput / wal2json;
  • pg_recvlogical:PG 自带的命令行消费工具。

pg_recvlogical 用法:

bash
# 创建 slot
pg_recvlogical -d learn_pg --slot=ch15_cdc_slot --create-slot -P pgoutput

# 持续消费(输出到 stdout / 文件)
pg_recvlogical -d learn_pg --slot=ch15_cdc_slot --start -o publication_names=ch15_pub_orders -o proto_version=1 -f -

8. 读写分离方案

应用怎么把读流量打到从库、写流量打到主库?三种主流方案:

8.1 应用层路由

最简单粗暴,代码里两个连接池:

python
write_pool = create_pool('host=master ...')
read_pool  = create_pool('host=replica1,replica2 ...')

def write(sql): return write_pool.execute(sql)
def read(sql):  return read_pool.execute(sql)

优点:简单,可控;缺点:业务代码改造,从库挂了要自己处理。

8.2 PgBouncer + 多端口

PgBouncer 是轻量连接池,本身不做读写分离,但可以通过端口区分

ini
# pgbouncer.ini
[databases]
learn_pg_w = host=master  port=5432 dbname=learn_pg
learn_pg_r = host=replica port=5432 dbname=learn_pg

应用按需连不同的库名。

8.3 PgPool-II(透明读写分离)

PgPool-II 解析 SQL:SELECT 自动路由到从库,INSERT/UPDATE/DELETE 路由到主库,业务无感知。代价是引入一层代理,会增加延迟、可能成为瓶颈。


9. 复制延迟监控

主库视角:

sql
SELECT
    client_addr,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), sent_lsn))    AS pending_send,
    pg_size_pretty(pg_wal_lsn_diff(sent_lsn, write_lsn))               AS pending_write,
    pg_size_pretty(pg_wal_lsn_diff(write_lsn, flush_lsn))              AS pending_flush,
    pg_size_pretty(pg_wal_lsn_diff(flush_lsn, replay_lsn))             AS pending_replay,
    pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn))  AS total_lag,
    write_lag, flush_lag, replay_lag
FROM pg_stat_replication;

从库视角:

sql
SELECT
    pg_is_in_recovery(),
    pg_last_wal_receive_lsn(),
    pg_last_wal_replay_lsn(),
    EXTRACT(EPOCH FROM (now() - pg_last_xact_replay_timestamp())) AS lag_seconds;

LSN(Log Sequence Number):WAL 里每个字节的全局唯一位置,写作 0/16BFB30(16 进制 segment/offset)。pg_wal_lsn_diff(a, b) = a 字节数 - b 字节数


10. 与 MySQL 复制全方位对比

维度MySQLPostgreSQL
复制日志binlog(独立于 redo log)WAL(同时承担 redo log + 复制)
默认复制类型逻辑(binlog row 模式)物理(流复制)
半同步rpl_semi_sync_* 插件synchronous_commit = remote_write
多源复制5.7+ 支持物理复制不支持,逻辑可订阅多个 publication
GTID5.6+ 支持物理复制无对应;逻辑复制有 origin
自动 FailoverInnoDB Cluster + MySQL RouterPatroni / repmgr 等外部组件
并行复制binlog 多线程回放物理复制单进程 redo(可调 recovery_min_apply_delay
跨版本复制binlog 兼容性较强物理不可,逻辑可
DDL 复制binlog 自动复制物理:是;逻辑:否(PG 16/17 部分支持)
主从读写一致InnoDB Cluster + Group Replicationquorum 复制 + remote_apply

11. 综合实战:用 docker-compose 起一主一从

随教程附带的 code/02_setup_streaming_replication.sh 是完整脚本。这里给出关键 docker-compose 片段:

yaml
version: '3.8'
services:
  pg_master:
    image: postgres:16
    container_name: pg_master
    environment:
      POSTGRES_PASSWORD: s3cret
      POSTGRES_DB: learn_pg
    command: >
      postgres
      -c wal_level=logical
      -c max_wal_senders=10
      -c max_replication_slots=10
      -c hot_standby=on
      -c synchronous_commit=on
    ports:
      - "5432:5432"
    volumes:
      - ./master_data:/var/lib/postgresql/data
      - ./pg_hba_extra.conf:/etc/pg_hba_extra.conf

  pg_replica:
    image: postgres:16
    container_name: pg_replica
    depends_on:
      - pg_master
    ports:
      - "5433:5432"
    entrypoint: ["/bin/bash", "/setup_replica.sh"]
    volumes:
      - ./replica_data:/var/lib/postgresql/data
      - ./setup_replica.sh:/setup_replica.sh

setup_replica.sh 在容器启动时执行 pg_basebackup -R -h pg_master -U repl_user,然后 postgres -D /var/lib/postgresql/data,整套链路就跑起来了。


12. 小结

把这一章浓缩成 10 句话:

  1. PG 复制分物理(流复制 / 文件归档)和逻辑(pub/sub)两大体系。
  2. 流复制三大角色:walsender → 网络 → walreceiver → startup redo
  3. 配置流复制 = 主库改 wal_level/max_wal_senders/pg_hba.conf + 从库 pg_basebackup -R + 启动。
  4. 复制槽防止主库回收从库还需要的 WAL,但要监控 max_slot_wal_keep_size 防爆盘。
  5. 同步级别synchronous_commit 配合 synchronous_standby_names 调,支持 quorum。
  6. 从库默认可读(Hot Standby),但不可写。
  7. 故障切换 = 手动 pg_promote() / 自动 Patroni(PG 内核不带自动 Failover)。
  8. 逻辑复制基于逻辑解码,跨版本 / 跨表结构 / 可双向,但 DDL 不复制。
  9. 读写分离三种方案:应用层路由 / PgBouncer / PgPool-II。
  10. 与 MySQL 比,PG 物理复制对应「字节级镜像」,MySQL 复制天生就是「逻辑」。


🎮 配套演示

用浏览器打开 ./15_replication/demo.html,跟着可视化动画再走一遍本章核心概念。

配套代码在 ./15_replication/code/,每个脚本都可以独立 bash xxx.sh(或 python xxx.py)运行,先跑 init.sql 准备数据。


13. 面试高频题(答案附带详细要点)

Q1:PG 流复制原理?walsender、walreceiver、startup 进程各做什么?

考察点:流复制的核心进程协作。

标准答案

PG 流复制由三个进程串成一条流水线:

  1. walsender:跑在主库上,每个连过来的从库对应一个独立的 walsender 进程。它从主库 pg_wal/ 目录读取 WAL 字节流,通过长连接 TCP 推给从库。心跳 / 流控 / 同步反馈都在它手里。

  2. walreceiver:跑在从库上,只有一个,负责接收 walsender 推来的 WAL,调用 write/fsync 把字节落盘到从库本地 pg_wal/。它会回三种 ACK:write_lsn(已写到 OS 缓存)、flush_lsn(已 fsync)、replay_lsn(已 redo 完)。

  3. startup 进程:从库 redo 的核心,启动后就一直跑(不像主库 startup 只在恢复阶段跑)。它从 pg_wal/ 读 WAL,按 LSN 顺序把页面变更应用到 shared_buffers + 数据文件,使从库数据持续追上主库。

整个链路上,walsender 和 walreceiver 之间用 PG 内部的「复制协议」通信(基于 PG 协议 v3 的 CopyData 消息),保持心跳;从库通过 hot_standby = on 把 redo 后的状态对外只读暴露。

加分项:能说清「写完 WAL ≠ 应用 WAL」,因此延迟分 write/flush/replay 三个层次,监控 pg_stat_replication 时要分别看;能提到 hot_standby_feedback 防止从库长查询被主库 VACUUM 掉。

与 MySQL 对比:MySQL 主从是 IO Thread + SQL Thread 两个线程模型,IO Thread ≈ walreceiver,SQL Thread ≈ startup;但 MySQL 复制走 binlog(逻辑日志),PG 走 WAL(物理日志),所以 PG 复制是「字节复印」,MySQL 是「事件回放」。


Q2:什么是复制槽?为什么需要?有什么风险?怎么防范?

考察点:复制槽的作用、运维风险。

标准答案

作用:复制槽(Replication Slot)是主库上的「书签」,记录某个从库(或逻辑订阅者)当前消费到的 WAL 位置(restart_lsn)。主库回收 WAL 文件时会保证 restart_lsn 之前的 WAL 不会被删

为什么需要:早期的 wal_keep_segments(PG 13+ 改名 wal_keep_size)只能粗暴指定保留多少 WAL 段,从库断开太久仍可能落后到无法追上,需要重做 base backup。复制槽精确按需保留,从库再连回来一定能续上。

风险:从库长期不消费(机器宕机、网络断开但没拆 slot),主库 WAL 越堆越多,pg_wal/ 目录把磁盘撑爆,主库崩溃。这是 PG 运维事故 Top 1 之一。

防范

  1. PG 13+ 设置 max_slot_wal_keep_size(如 64GB),超过就强制释放 slot 而不再保留 WAL,主库牺牲那个从库以保命;
  2. 监控 pg_replication_slots,看 active 是否为 false 以及 restart_lsnpg_current_wal_lsn() 的字节数,超阈值告警;
  3. 弃用从库时一定要 pg_drop_replication_slot('xxx'),不要只删机器;
  4. 物理槽用得保守(用 wal_keep_size 兜底),逻辑槽必须用(否则订阅断开就丢)。

加分项:能区分物理槽(pg_create_physical_replication_slot)和逻辑槽(pg_create_logical_replication_slot 必须指定 plugin),并解释为什么逻辑槽必须存在(消费状态不能由订阅端单方面记录)。


Q3:synchronous_commit 五个值的区别?什么是 quorum 复制?

考察点:同步级别的取舍、新特性 quorum。

标准答案

synchronous_commit 控制主库 commit 时等到什么程度才算成功:

等什么数据安全
off啥都不等,连 WAL fsync 也不等单机崩溃可能丢几百 ms
local只等本机 WAL fsync单机安全,从库可能落后
on(默认)本机 fsync + 同步从库 远程写到 OS 缓存(=remote_write)主库挂,同步从库一定有这条 WAL
remote_flush本机 fsync + 同步从库 fsync 落盘强一致
remote_apply本机 fsync + 同步从库 redo 完毕切换无延迟读

同步从库」由 synchronous_standby_names 决定,分两种语法:

  • FIRST k (s1, s2, ...)优先级模式,列表前 k 个可用从库参与同步;
  • ANY k (s1, s2, s3, ...)Quorum 模式(PG 10+),任意 k 个 ACK 即视为成功。

Quorum 复制的优势:3 台同步从库要求 ANY 2,则任意一台从库挂掉仍能 commit,不卡主库;同步级别能跨机房做「至少 2 个机房 ACK 才算成功」的灾备。

取舍建议

  • 普通业务:synchronous_commit = on + 异步从库做读扩展;
  • 金融 / 不容丢数据:remote_apply + ANY 2 (s1, s2, s3) 跨 AZ 部署;
  • 日志型业务:local 即可,不阻塞写。

加分项:能说出「synchronous_commit 是会话级参数」,可以在事务里临时改值(重要事务用 remote_apply,普通事务 local);能解释「主库 commit 等不到 quorum 时会卡住直到从库恢复」,所以同步从库挂多了反而拖死主库。


Q4:物理复制和逻辑复制的区别?各适合什么场景?

考察点:两种复制的对比、选型能力。

标准答案

根本区别

维度物理复制逻辑复制
复制内容WAL 字节流(物理页面变更)逻辑事件(INSERT/UPDATE/DELETE 行)
主从必须同大版本、同表结构、同 WAL 格式不同版本、不同表结构均可
从库能否写不能可以独立写其他表
复制粒度整个集群按表,可带 WHERE 过滤(PG 15+)
DDL自动复制不复制(业务自己同步)
序列复制不复制
性能高(字节复印)低一些(解码 + 应用)

场景选型

  • 物理复制适合:主备灾备、读写分离、热备升级(小版本升级)。要求从库和主库一摸一样。
  • 逻辑复制适合
    • 跨大版本升级(PG 12 → PG 17 平滑切换);
    • CDC 数据管道(变更打到 Kafka / ES);
    • 多源汇聚(多个 OLTP 库的部分表汇到一个数据仓库);
    • 双向复制(PG 16+ 内置);
    • 子集复制(只把某些表 / 某些行同步给报表库)。

底层实现差异:物理复制就是 walreceiver + startup 直接 redo WAL;逻辑复制走「逻辑解码」流水线:walsender 把 WAL 喂给 reorder buffer,按事务边界缓存,commit 时调用 output plugin(pgoutput / wal2json)转换成逻辑事件再发出去。

加分项:能讲清「为什么物理复制要求主从同表结构」(因为物理复制传的是页面字节,包括元组在页内的偏移;表结构变了字节布局就乱);能讲清「逻辑解码不传中间状态」,所以一个大事务在主库提交后才会一次性发给订阅端,造成「批量延迟」,PG 14+ 的 streaming = on 能边写边发缓解。


Q5:PG 没有内置自动 Failover,怎么做高可用?Patroni 怎么工作?

考察点:高可用方案、Patroni 选举原理。

标准答案

PG 内核有意不实现自动 Failover,因为「判断主库是否真的挂了」涉及网络分区、脑裂、仲裁,是分布式难题,社区把这部分留给外部组件做。主流方案:

方案特点
Patroni最流行,依赖 etcd/Consul/ZK,K8s 生态首选
repmgr老牌,轻量,配置较繁琐
pg_auto_failover微软出品,自带 monitor 节点
StolonSorintlab 出品,K8s 友好
PgPool-II兼连接池 + 负载均衡,Watchdog 做选举

Patroni 工作流程

  1. 节点启动:每台 PG 旁起一个 Patroni 守护进程,所有 Patroni 都连同一个 etcd 集群;
  2. leader 选举:所有节点抢 etcd 上的 /patroni/<cluster>/leader key(CAS 写入),写入成功的成为主,并给 key 一个 TTL(默认 30s);
  3. 续约:主节点定期续约 leader key(loop_wait 默认 10s);从节点不断检查 leader key 是否还在;
  4. 故障检测:leader key TTL 到了没续约 → 从节点感知到 → 进入选举流程;
  5. 选举新主:所有候选者比较 LSN(pg_last_wal_replay_lsn()),LSN 最大的获胜,调用 pg_promote() 提升自己;
  6. 重定向其他从库:新主把自己的连接信息写回 etcd,其他从库读到后修改 primary_conninfo,rewind 到新 timeline 后挂到新主。

关键风险

  • 脑裂:网络分区导致主库与 etcd 失联,但客户端还能连主库,此时新主已选出,旧主仍接受写。Patroni 通过「rewind + 旧主自动 demote」缓解;
  • fencing:彻底防脑裂需要 STONITH(关电 / 关网卡),Patroni 自带 pre_promote 钩子可以集成。

加分项:能讲到 Patroni 通过 HAProxy + 自带 /master /replica HTTP 探针接口实现客户端负载均衡;能说出「Patroni 推荐 etcd 而不是 Consul」(CAS 性能更好);能解释「failover 后旧主用 pg_rewind 工具 reattach 比重做 base backup 快很多」。


Q6:什么是 pg_rewind?什么场景下用?

考察点:高可用运维细节。

标准答案

pg_rewind 是 PG 自带工具,作用是把一个旧主库快速变成新主库的从库,避免重新做 base backup。

场景

主库 A → 从库 B 流复制中。某天 A 短暂挂了,B 被 promote 成了新主,业务切到 B 继续写。一段时间后 A 修好,想让 A 重新加入集群当 B 的从库。

问题:A 和 B 在「分叉点」之后各自有不同的 WAL(A 上有「自己提交但 B 没收到」的旧事务,B 上有「自己 promote 后接受的新事务」)。直接把 A 当 B 的从库启动会因为「数据文件冲突」失败。

传统方案:删 A 的 data 目录,对 B 做 pg_basebackup 重新拷一份。1TB 数据要拷半天。

pg_rewind 方案

  1. 找到 A 和 B 的「时间线分叉点」(比较各自 timeline history file);
  2. 通过逻辑解码 WAL 找出 A 在分叉点之后修改了哪些数据块;
  3. 只把这些数据块从 B 重新拷过来覆盖 A,其他数据原地复用;
  4. A 启动时从分叉点开始 redo B 的 WAL,追上之后挂到 B 下面。

前提条件

  • 两边 wal_log_hints = on 或开了 data checksums(否则 hint bit 修改不写 WAL,rewind 找不全差异);
  • A 必须能 clean shutdown(如果 crash 残留可以用 --source-pgdata + 手动恢复)。

加分项:能说出 pg_rewind 速度的本质是「按差异块拷贝」而不是「按全量拷贝」;TB 级数据 1 分钟内能完成;Patroni 在 failover 后默认会调用 pg_rewind 恢复旧主;rewind 不行的话才回退到 pg_basebackup


Q7:怎么用 PG 做 CDC(变更数据捕获)?输出到 Kafka 怎么做?

考察点:逻辑复制的工程化应用。

标准答案

PG 的 CDC 基础设施是逻辑解码 + 逻辑复制槽,落地有两条路线:

路线 1:Debezium(业界主流)

PG (logical slot, pgoutput)
   ↓ logical decoding
Debezium PG Connector (Kafka Connect)

Kafka Topic(一表一 topic)

下游消费者(ES / ClickHouse / Hudi)

配置要点:

  • 主库 wal_level = logicalmax_replication_slots ≥ 1
  • 创建账号:CREATE ROLE debezium WITH REPLICATION LOGIN PASSWORD ...
  • Debezium connector 配置 slot.namepublication.autocreate.mode = filteredplugin.name = pgoutput
  • 全表初始 snapshot + 之后增量。

路线 2:自研pg_recvlogical

bash
pg_recvlogical -d learn_pg --slot=ch15_cdc --create-slot -P pgoutput
pg_recvlogical -d learn_pg --slot=ch15_cdc --start \
    -o publication_names=ch15_pub_orders -o proto_version=1 \
    | python parse_and_send_to_kafka.py

或用 wal2json plugin 输出 JSON:

bash
pg_recvlogical -d learn_pg --slot=ch15_cdc --create-slot -P wal2json
pg_recvlogical -d learn_pg --slot=ch15_cdc --start -f -
# 输出形如:
# {"change":[{"kind":"insert","schema":"public","table":"ch15_orders","columnvalues":[1, 100, "paid"]}]}

关键运维点

  1. slot 一定要监控:CDC 消费方挂了,slot 还在,主库 WAL 会涨;
  2. 大事务影响:1 亿行的 UPDATE 要等 commit 才发,下游延迟剧增,PG 14+ 的 streaming = on 能 commit 前流式发;
  3. 断点续传:消费方记录已处理的 LSN,重启后 pg_recvlogical --startpos=<LSN>
  4. DDL:业务自己保证下游 schema 兼容。

加分项:能解释「逻辑复制槽 + 复制协议」组合形成天然的 exactly-once 投递(消费方 ACK 后 PG 才推进 slot 位置);能讲 Debezium 的 schema registry 集成;能讲 wal2json vs pgoutput 的取舍(pgoutput 二进制紧凑,wal2json 可读性好但解析慢)。


14. 实操彩蛋:5 行命令搭起一主一从

bash
# 起主库
docker run -d --name pg_master \
    -e POSTGRES_PASSWORD=s3cret \
    -e POSTGRES_DB=learn_pg \
    -p 5432:5432 \
    postgres:16 \
    -c wal_level=logical -c max_wal_senders=10 -c hot_standby=on

# 进主库建账号
docker exec -it pg_master psql -U postgres -c \
    "CREATE ROLE repl_user WITH REPLICATION LOGIN PASSWORD 'r3pl';"

# 把 pg_hba.conf 加 replication 规则(追加一行后 reload)
docker exec pg_master bash -c \
    'echo "host replication repl_user 0.0.0.0/0 md5" >> /var/lib/postgresql/data/pg_hba.conf'
docker exec pg_master psql -U postgres -c "SELECT pg_reload_conf();"

# 起从库(pg_basebackup 一键克隆)
docker run -d --name pg_replica --link pg_master \
    -p 5433:5432 \
    -e POSTGRES_PASSWORD=s3cret postgres:16 \
    bash -c 'rm -rf /var/lib/postgresql/data/* && \
             PGPASSWORD=r3pl pg_basebackup -h pg_master -U repl_user -D /var/lib/postgresql/data -R -P -Xs && \
             postgres -D /var/lib/postgresql/data'

# 验证
docker exec pg_master psql -U postgres -c "SELECT * FROM pg_stat_replication;"

详细完整脚本见同目录 code/02_setup_streaming_replication.sh


🔗 延伸阅读

🎬 可视化演示

演示加载缓慢或样式异常?点此在新标签页打开 ↗

💻 示例代码

python
#!/usr/bin/env python3
"""
01_show_replication_status.py
=============================
从主库查询当前所有从库的复制状态、复制槽状态、复制延迟。

依赖:psycopg[binary] >= 3.1
    pip install "psycopg[binary]"

连接:默认 host=127.0.0.1 port=5432 dbname=learn_pg user=postgres
"""
import os
import sys
import psycopg

CONN_INFO = os.environ.get(
    "PG_DSN",
    "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres",
)


def hr(title: str) -> None:
    print("\n" + "=" * 70)
    print(f"  {title}")
    print("=" * 70)


def query_and_print(cur, sql: str) -> None:
    cur.execute(sql)
    cols = [d.name for d in cur.description]
    rows = cur.fetchall()
    if not rows:
        print("  (无记录)")
        return
    widths = [max(len(c), max((len(str(r[i])) for r in rows), default=0)) for i, c in enumerate(cols)]
    fmt = "  " + "  ".join(f"{{:<{w}}}" for w in widths)
    print(fmt.format(*cols))
    print("  " + "-" * (sum(widths) + 2 * (len(cols) - 1)))
    for r in rows:
        print(fmt.format(*[str(v) for v in r]))


def main() -> int:
    try:
        conn = psycopg.connect(CONN_INFO)
    except psycopg.Error as e:
        print(f"[FATAL] 无法连接主库: {e}", file=sys.stderr)
        return 1

    with conn, conn.cursor() as cur:
        hr("1. 当前 PG 角色")
        cur.execute("SELECT pg_is_in_recovery() AS in_recovery, current_setting('wal_level') AS wal_level;")
        in_recovery, wal_level = cur.fetchone()
        print(f"  pg_is_in_recovery = {in_recovery}   {'(从库)' if in_recovery else '(主库)'}")
        print(f"  wal_level         = {wal_level}")

        hr("2. pg_stat_replication(连上来的从库)")
        query_and_print(cur, """
            SELECT
                application_name,
                client_addr,
                state,
                sync_state,
                pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), sent_lsn))    AS pending_send,
                pg_size_pretty(pg_wal_lsn_diff(sent_lsn, write_lsn))               AS pending_write,
                pg_size_pretty(pg_wal_lsn_diff(write_lsn, flush_lsn))              AS pending_flush,
                pg_size_pretty(pg_wal_lsn_diff(flush_lsn, replay_lsn))             AS pending_replay,
                pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn))  AS total_lag,
                COALESCE(write_lag::text,  '-') AS write_lag,
                COALESCE(flush_lag::text,  '-') AS flush_lag,
                COALESCE(replay_lag::text, '-') AS replay_lag
            FROM pg_stat_replication
            ORDER BY application_name;
        """)

        hr("3. pg_replication_slots(所有复制槽)")
        query_and_print(cur, """
            SELECT
                slot_name,
                slot_type,
                database,
                active,
                COALESCE(active_pid::text, '-') AS active_pid,
                pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained_wal
            FROM pg_replication_slots
            ORDER BY slot_name;
        """)

        hr("4. WAL 写入位置(主库视角)")
        cur.execute("""
            SELECT
                pg_current_wal_lsn() AS current_lsn,
                pg_walfile_name(pg_current_wal_lsn()) AS current_walfile;
        """)
        lsn, walfile = cur.fetchone()
        print(f"  current_lsn      = {lsn}")
        print(f"  current_walfile  = {walfile}")

        hr("5. PUBLICATIONS")
        query_and_print(cur, """
            SELECT pubname, puballtables, pubinsert, pubupdate, pubdelete, pubtruncate
            FROM pg_publication;
        """)

    print("\n[DONE] 查询完成。如要持续监控,可加上 watch 1:")
    print("  watch -n 1 \"python3 01_show_replication_status.py | head -50\"")
    return 0


if __name__ == "__main__":
    sys.exit(main())
bash
#!/usr/bin/env bash
# =====================================================================
# 02_setup_streaming_replication.sh
# ---------------------------------------------------------------------
# 用 docker 起一主一从 PG 16 流复制集群。
#
# 主库:pg_master,端口 5432
# 从库:pg_replica,端口 5433
#
# 默认账号:postgres / s3cret
# 复制账号:repl_user / r3pl_pwd
#
# 用法:
#   ./02_setup_streaming_replication.sh up        # 起集群
#   ./02_setup_streaming_replication.sh status    # 查复制状态
#   ./02_setup_streaming_replication.sh write     # 主库写入测试数据
#   ./02_setup_streaming_replication.sh read      # 从库验证可读
#   ./02_setup_streaming_replication.sh failover  # 演示 promote 从库
#   ./02_setup_streaming_replication.sh down      # 销毁
# =====================================================================

set -euo pipefail

PG_IMAGE="postgres:16"
NETWORK="pg_repl_net"
MASTER="pg_master"
REPLICA="pg_replica"
PGPASS="s3cret"
REPL_USER="repl_user"
REPL_PASS="r3pl_pwd"

cmd_up() {
    docker network inspect "$NETWORK" >/dev/null 2>&1 || docker network create "$NETWORK"

    echo "[1/5] 启动主库 $MASTER..."
    docker run -d --name "$MASTER" --network "$NETWORK" \
        -e POSTGRES_PASSWORD="$PGPASS" \
        -e POSTGRES_DB=learn_pg \
        -p 5432:5432 \
        "$PG_IMAGE" \
        -c wal_level=logical \
        -c max_wal_senders=10 \
        -c max_replication_slots=10 \
        -c hot_standby=on \
        -c synchronous_commit=on \
        -c wal_keep_size=128MB

    echo "[2/5] 等待主库就绪..."
    for i in {1..30}; do
        if docker exec "$MASTER" pg_isready -U postgres >/dev/null 2>&1; then break; fi
        sleep 1
    done

    echo "[3/5] 创建复制账号 + 修改 pg_hba.conf..."
    docker exec "$MASTER" psql -U postgres -c \
        "CREATE ROLE $REPL_USER WITH REPLICATION LOGIN PASSWORD '$REPL_PASS';"
    docker exec "$MASTER" bash -c \
        "echo 'host replication $REPL_USER 0.0.0.0/0 md5' >> /var/lib/postgresql/data/pg_hba.conf"
    docker exec "$MASTER" psql -U postgres -c "SELECT pg_reload_conf();"

    echo "[4/5] 在主库创建物理复制槽 + PUBLICATION..."
    docker exec "$MASTER" psql -U postgres -d learn_pg -c \
        "SELECT pg_create_physical_replication_slot('ch15_standby1_slot');"
    docker exec "$MASTER" psql -U postgres -d learn_pg -c \
        "CREATE TABLE IF NOT EXISTS ch15_demo (id BIGSERIAL PRIMARY KEY, msg TEXT, ts TIMESTAMPTZ DEFAULT now());"

    echo "[5/5] 启动从库 $REPLICA(pg_basebackup 自动克隆主库)..."
    docker run -d --name "$REPLICA" --network "$NETWORK" \
        -e POSTGRES_PASSWORD="$PGPASS" \
        -e PGPASSWORD="$REPL_PASS" \
        -p 5433:5432 \
        --entrypoint /bin/bash \
        "$PG_IMAGE" \
        -c "
            set -e
            rm -rf /var/lib/postgresql/data/*
            echo '正在做基础备份...'
            pg_basebackup -h $MASTER -p 5432 -U $REPL_USER \
                -D /var/lib/postgresql/data \
                -Fp -Xs -P -R -S ch15_standby1_slot
            echo '基础备份完成,启动从库...'
            chown -R postgres:postgres /var/lib/postgresql/data
            chmod 0700 /var/lib/postgresql/data
            su postgres -c 'postgres -D /var/lib/postgresql/data'
        "

    sleep 5
    echo
    echo "✅ 集群已起。下面查复制状态:"
    cmd_status
}

cmd_status() {
    echo "--- 主库 pg_stat_replication ---"
    docker exec "$MASTER" psql -U postgres -d learn_pg -c "
        SELECT application_name, client_addr, state, sync_state,
               pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn) AS lag_bytes
        FROM pg_stat_replication;
    "
    echo "--- 主库 pg_replication_slots ---"
    docker exec "$MASTER" psql -U postgres -d learn_pg -c "
        SELECT slot_name, slot_type, active, restart_lsn FROM pg_replication_slots;
    "
    echo "--- 从库角色 ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_is_in_recovery() AS in_recovery, pg_last_wal_replay_lsn() AS replayed;" || true
}

cmd_write() {
    echo "--- 主库写入 5 行数据 ---"
    docker exec "$MASTER" psql -U postgres -d learn_pg -c \
        "INSERT INTO ch15_demo(msg) SELECT 'hello-' || g FROM generate_series(1,5) g RETURNING *;"
}

cmd_read() {
    echo "--- 从库查询数据(应能看到主库写入的行) ---"
    sleep 1
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT * FROM ch15_demo ORDER BY id DESC LIMIT 10;"
    echo "--- 试试在从库写(应失败:read-only) ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "INSERT INTO ch15_demo(msg) VALUES('从库写入')" || true
}

cmd_failover() {
    echo "--- Step 1: 停掉主库 ---"
    docker stop "$MASTER"

    echo "--- Step 2: 在从库 pg_promote() ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_promote(wait => true, wait_seconds => 30);"

    echo "--- Step 3: 验证从库变成主库 ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_is_in_recovery() AS still_replica;"

    echo "--- Step 4: 在新主库写入 ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "INSERT INTO ch15_demo(msg) VALUES('写入新主库') RETURNING *;"
}

cmd_down() {
    docker rm -f "$MASTER" "$REPLICA" 2>/dev/null || true
    docker network rm "$NETWORK" 2>/dev/null || true
    echo "✅ 已清理。"
}

case "${1:-}" in
    up)        cmd_up ;;
    status)    cmd_status ;;
    write)     cmd_write ;;
    read)      cmd_read ;;
    failover)  cmd_failover ;;
    down)      cmd_down ;;
    *)
        echo "Usage: $0 {up|status|write|read|failover|down}"
        exit 1
        ;;
esac
python
#!/usr/bin/env python3
"""
03_logical_replication.py
=========================
演示 PG 逻辑复制:在主库建 PUBLICATION,在从库建 SUBSCRIPTION,
然后在主库写数据,验证从库能同步看到。

前置条件:
    - 主库(端口 5432)和从库(端口 5433,独立实例非物理 standby)都在跑
    - 两边 wal_level >= logical
    - 两边都有 learn_pg 数据库
    - repl_user 账号已创建
    - 两边都有 ch15_logical_demo 表(脚本会自动建)

依赖:psycopg[binary]>=3.1
    pip install "psycopg[binary]"

用法:
    python3 03_logical_replication.py setup    # 建表 + 发布 + 订阅
    python3 03_logical_replication.py write    # 主库写入测试
    python3 03_logical_replication.py verify   # 验证从库同步
    python3 03_logical_replication.py status   # 查 publication / subscription
    python3 03_logical_replication.py teardown # 清理
"""
import os
import sys
import time

import psycopg

MASTER_DSN = os.environ.get(
    "MASTER_DSN",
    "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres",
)
REPLICA_DSN = os.environ.get(
    "REPLICA_DSN",
    "host=127.0.0.1 port=5433 dbname=learn_pg user=postgres password=postgres",
)
# subscription 端连接主库用的 DSN(注意密码改成 repl_user 的)
SUB_CONN = os.environ.get(
    "SUB_CONN",
    "host=127.0.0.1 port=5432 dbname=learn_pg user=repl_user password=r3pl_pwd",
)


def exec_sql(dsn: str, sql: str, fetch: bool = False):
    with psycopg.connect(dsn, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute(sql)
        if fetch and cur.description:
            return cur.fetchall()
    return None


def setup() -> None:
    print("[1/4] 主库建表 + 发布")
    exec_sql(MASTER_DSN, """
        CREATE TABLE IF NOT EXISTS ch15_logical_demo (
            id     BIGSERIAL PRIMARY KEY,
            payload TEXT NOT NULL,
            ts     TIMESTAMPTZ DEFAULT now()
        );
        DROP PUBLICATION IF EXISTS ch15_pub_logical_demo;
        CREATE PUBLICATION ch15_pub_logical_demo FOR TABLE ch15_logical_demo;
    """)

    print("[2/4] 从库建表(DDL 不会自动复制,必须手工建)")
    exec_sql(REPLICA_DSN, """
        CREATE TABLE IF NOT EXISTS ch15_logical_demo (
            id     BIGSERIAL PRIMARY KEY,
            payload TEXT NOT NULL,
            ts     TIMESTAMPTZ DEFAULT now()
        );
        TRUNCATE ch15_logical_demo;
    """)

    print("[3/4] 从库建订阅")
    exec_sql(REPLICA_DSN, "DROP SUBSCRIPTION IF EXISTS ch15_sub_logical_demo;")
    exec_sql(REPLICA_DSN, f"""
        CREATE SUBSCRIPTION ch15_sub_logical_demo
        CONNECTION '{SUB_CONN}'
        PUBLICATION ch15_pub_logical_demo;
    """)

    print("[4/4] 等 5 秒让初始 sync 完成...")
    time.sleep(5)
    print("✅ 设置完成。")


def write() -> None:
    print("--- 主库插入 10 行数据 ---")
    rows = exec_sql(
        MASTER_DSN,
        """
        INSERT INTO ch15_logical_demo(payload)
        SELECT 'logical-' || g FROM generate_series(1,10) g
        RETURNING id, payload;
        """,
        fetch=True,
    )
    for r in rows or []:
        print(f"  inserted: {r}")


def verify() -> None:
    print("--- 主库当前数据量 ---")
    m = exec_sql(MASTER_DSN, "SELECT count(*), max(id) FROM ch15_logical_demo;", fetch=True)
    print(f"  master: count={m[0][0]} max_id={m[0][1]}")

    print("--- 等 2 秒让逻辑复制追上 ---")
    time.sleep(2)

    print("--- 从库当前数据量 ---")
    r = exec_sql(REPLICA_DSN, "SELECT count(*), max(id) FROM ch15_logical_demo;", fetch=True)
    print(f"  replica: count={r[0][0]} max_id={r[0][1]}")

    if m[0][0] == r[0][0]:
        print("✅ 主从一致")
    else:
        print("⚠️  主从不一致,请检查 pg_stat_subscription / pg_stat_replication")


def status() -> None:
    print("=== 主库 pg_publication ===")
    for row in exec_sql(MASTER_DSN, "SELECT pubname, puballtables FROM pg_publication;", fetch=True) or []:
        print(f"  {row}")

    print("\n=== 主库 pg_replication_slots(逻辑槽) ===")
    for row in exec_sql(
        MASTER_DSN,
        "SELECT slot_name, slot_type, plugin, active FROM pg_replication_slots WHERE slot_type='logical';",
        fetch=True,
    ) or []:
        print(f"  {row}")

    print("\n=== 主库 pg_stat_replication ===")
    for row in exec_sql(
        MASTER_DSN,
        "SELECT application_name, client_addr, state, sync_state FROM pg_stat_replication;",
        fetch=True,
    ) or []:
        print(f"  {row}")

    print("\n=== 从库 pg_subscription ===")
    for row in exec_sql(
        REPLICA_DSN,
        "SELECT subname, subenabled, subslotname FROM pg_subscription;",
        fetch=True,
    ) or []:
        print(f"  {row}")

    print("\n=== 从库 pg_stat_subscription ===")
    for row in exec_sql(
        REPLICA_DSN,
        "SELECT subname, pid, received_lsn, latest_end_lsn FROM pg_stat_subscription;",
        fetch=True,
    ) or []:
        print(f"  {row}")


def teardown() -> None:
    print("[1/3] 从库删订阅")
    try:
        exec_sql(REPLICA_DSN, "ALTER SUBSCRIPTION ch15_sub_logical_demo DISABLE;")
        exec_sql(REPLICA_DSN, "ALTER SUBSCRIPTION ch15_sub_logical_demo SET (slot_name = NONE);")
        exec_sql(REPLICA_DSN, "DROP SUBSCRIPTION ch15_sub_logical_demo;")
    except psycopg.Error as e:
        print(f"  warn: {e}")

    print("[2/3] 主库删发布 + 槽")
    try:
        exec_sql(MASTER_DSN, "DROP PUBLICATION IF EXISTS ch15_pub_logical_demo;")
        exec_sql(MASTER_DSN, "SELECT pg_drop_replication_slot('ch15_sub_logical_demo') WHERE EXISTS (SELECT 1 FROM pg_replication_slots WHERE slot_name='ch15_sub_logical_demo');")
    except psycopg.Error as e:
        print(f"  warn: {e}")

    print("[3/3] 删表")
    exec_sql(MASTER_DSN, "DROP TABLE IF EXISTS ch15_logical_demo;")
    exec_sql(REPLICA_DSN, "DROP TABLE IF EXISTS ch15_logical_demo;")
    print("✅ 已清理。")


COMMANDS = {
    "setup": setup,
    "write": write,
    "verify": verify,
    "status": status,
    "teardown": teardown,
}

if __name__ == "__main__":
    cmd = sys.argv[1] if len(sys.argv) > 1 else ""
    if cmd not in COMMANDS:
        print(f"Usage: {sys.argv[0]} {{{'|'.join(COMMANDS.keys())}}}", file=sys.stderr)
        sys.exit(1)
    COMMANDS[cmd]()
bash
#!/usr/bin/env bash
# =====================================================================
# 04_promote_failover.sh
# ---------------------------------------------------------------------
# 演示「主库挂了 → 从库提升为主库 → 业务继续写入」的 Failover 流程。
#
# 前置:已运行 02_setup_streaming_replication.sh up 起好集群。
# 用法:
#   ./04_promote_failover.sh demo       # 跑完整剧本
#   ./04_promote_failover.sh promote    # 只做提升
#   ./04_promote_failover.sh check      # 检查从库当前角色
#   ./04_promote_failover.sh rewind     # 用 pg_rewind 让旧主重新入集群
# =====================================================================

set -euo pipefail

MASTER="pg_master"
REPLICA="pg_replica"

step() {
    echo
    echo "========================================================"
    echo "  $*"
    echo "========================================================"
}

cmd_check() {
    echo "--- 主库角色 ---"
    docker exec "$MASTER" psql -U postgres -d learn_pg \
        -c "SELECT pg_is_in_recovery() AS in_recovery;" 2>/dev/null || echo "(主库不可达)"
    echo "--- 从库角色 ---"
    docker exec "$REPLICA" psql -U postgres -d learn_pg \
        -c "SELECT pg_is_in_recovery() AS in_recovery;"
}

cmd_promote() {
    step "在从库上调用 pg_promote()"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_promote(wait => true, wait_seconds => 30);"
    sleep 2
    cmd_check
}

cmd_demo() {
    step "Step 1: 主库写一条标记数据"
    docker exec "$MASTER" psql -U postgres -d learn_pg -c \
        "INSERT INTO ch15_demo(msg) VALUES('failover-marker-A') RETURNING *;" || true

    sleep 1

    step "Step 2: 验证从库已经收到(流复制)"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT * FROM ch15_demo WHERE msg LIKE 'failover-marker%' ORDER BY id;"

    step "Step 3: 模拟主库故障 —— docker stop pg_master"
    docker stop "$MASTER"

    step "Step 4: 此时业务尝试连主库(应失败)"
    docker exec "$REPLICA" psql -h "$MASTER" -U postgres -d learn_pg \
        -c "SELECT 1;" 2>&1 | head -3 || echo "(预期失败)"

    step "Step 5: 把从库 promote 成新主库"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_promote(wait => true, wait_seconds => 30);"
    sleep 2

    step "Step 6: 检查从库是否已经变主"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT pg_is_in_recovery() AS still_replica;"

    step "Step 7: 在新主库写入"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "INSERT INTO ch15_demo(msg) VALUES('failover-marker-B-on-new-master') RETURNING *;"

    step "Step 8: 看 timeline 切换"
    docker exec "$REPLICA" bash -c "ls /var/lib/postgresql/data/pg_wal/ | head -10"

    echo
    echo "✅ Failover 演示完毕。新主在 pg_replica(端口 5433)。"
    echo "   要让旧主重新加入,运行: ./04_promote_failover.sh rewind"
}

cmd_rewind() {
    step "用 pg_rewind 让旧主 pg_master 变成新主 pg_replica 的从库"
    echo "[1/4] 先启动旧主(启动后会以 standalone 模式起来)"
    docker start "$MASTER"
    sleep 3

    echo "[2/4] 立即 stop(避免它写入分叉数据)"
    docker exec "$MASTER" su postgres -c "pg_ctl -D /var/lib/postgresql/data stop -m fast" || true
    sleep 2

    echo "[3/4] 在旧主容器内执行 pg_rewind(从新主拉取差异块)"
    docker exec -e PGPASSWORD=postgres "$MASTER" su postgres -c \
        "pg_rewind --target-pgdata=/var/lib/postgresql/data \
                   --source-server='host=pg_replica port=5432 user=postgres dbname=postgres password=postgres' \
                   --progress" || {
            echo "⚠️  pg_rewind 失败。常见原因:旧主未开 wal_log_hints / data_checksums。"
            echo "    回退方案:删 data 目录,重新 pg_basebackup。"
            return 1
        }

    echo "[4/4] 写 standby.signal + primary_conninfo,启动旧主作为从库"
    docker exec "$MASTER" bash -c "
        touch /var/lib/postgresql/data/standby.signal
        cat >> /var/lib/postgresql/data/postgresql.auto.conf <<EOF
primary_conninfo = 'host=pg_replica port=5432 user=repl_user password=r3pl_pwd application_name=oldmaster'
EOF
        chown postgres:postgres /var/lib/postgresql/data/standby.signal
    "
    docker exec -d "$MASTER" su postgres -c "postgres -D /var/lib/postgresql/data"

    sleep 3
    echo
    echo "✅ rewind 完成。新主视角:"
    docker exec "$REPLICA" psql -U postgres -d learn_pg -c \
        "SELECT application_name, client_addr, state FROM pg_stat_replication;"
}

case "${1:-}" in
    demo)    cmd_demo ;;
    promote) cmd_promote ;;
    check)   cmd_check ;;
    rewind)  cmd_rewind ;;
    *)
        echo "Usage: $0 {demo|promote|check|rewind}"
        exit 1
        ;;
esac
python
#!/usr/bin/env python3
"""
05_replication_lag.py
=====================
持续监控复制延迟,输出主库各从库的 pending_send / pending_write /
pending_flush / pending_replay 字节数与时间延迟。

用法:
    python3 05_replication_lag.py                # 监控主库 5432
    python3 05_replication_lag.py --interval 0.5 # 0.5s 一次
    python3 05_replication_lag.py --once         # 只跑一次
    PG_DSN='host=... port=...' python3 05_replication_lag.py

依赖:psycopg[binary]>=3.1
"""
import argparse
import os
import sys
import time
from datetime import datetime

import psycopg

DEFAULT_DSN = "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres"


SQL = """
SELECT
    application_name,
    client_addr::text,
    state,
    sync_state,
    pg_wal_lsn_diff(pg_current_wal_lsn(), sent_lsn)    AS pending_send_bytes,
    pg_wal_lsn_diff(sent_lsn, write_lsn)               AS pending_write_bytes,
    pg_wal_lsn_diff(write_lsn, flush_lsn)              AS pending_flush_bytes,
    pg_wal_lsn_diff(flush_lsn, replay_lsn)             AS pending_replay_bytes,
    pg_wal_lsn_diff(pg_current_wal_lsn(), replay_lsn)  AS total_lag_bytes,
    EXTRACT(EPOCH FROM write_lag)::float  AS write_lag_s,
    EXTRACT(EPOCH FROM flush_lag)::float  AS flush_lag_s,
    EXTRACT(EPOCH FROM replay_lag)::float AS replay_lag_s
FROM pg_stat_replication
ORDER BY application_name;
"""


def human_bytes(n) -> str:
    if n is None:
        return "-"
    n = float(n)
    for unit in ["B", "KB", "MB", "GB", "TB"]:
        if abs(n) < 1024.0:
            return f"{n:6.1f}{unit}"
        n /= 1024.0
    return f"{n:6.1f}PB"


def human_secs(s) -> str:
    if s is None:
        return "-"
    if s < 1:
        return f"{int(s*1000)}ms"
    if s < 60:
        return f"{s:.1f}s"
    return f"{s/60:.1f}m"


def render(rows) -> None:
    ts = datetime.now().strftime("%H:%M:%S")
    print(f"\n[{ts}]  {len(rows)} 个从库连接")
    if not rows:
        print("  (无从库连接)")
        return
    print(f"  {'app':<20} {'addr':<16} {'state':<10} {'sync':<10}"
          f" {'send':>9} {'write':>9} {'flush':>9} {'replay':>9} {'total':>9}"
          f" | {'wlag':>6} {'flag':>6} {'rlag':>6}")
    print("  " + "-" * 130)
    for r in rows:
        (app, addr, state, sync,
         ps, pw, pf, pr, total,
         wlag, flag, rlag) = r
        print(f"  {(app or '-'):<20} {(addr or '-'):<16} {(state or '-'):<10} {(sync or '-'):<10}"
              f" {human_bytes(ps):>9} {human_bytes(pw):>9} {human_bytes(pf):>9}"
              f" {human_bytes(pr):>9} {human_bytes(total):>9}"
              f" | {human_secs(wlag):>6} {human_secs(flag):>6} {human_secs(rlag):>6}")


def main() -> int:
    parser = argparse.ArgumentParser()
    parser.add_argument("--interval", type=float, default=2.0, help="刷新间隔秒数")
    parser.add_argument("--once", action="store_true", help="只跑一次")
    parser.add_argument("--dsn", default=os.environ.get("PG_DSN", DEFAULT_DSN))
    args = parser.parse_args()

    try:
        conn = psycopg.connect(args.dsn)
    except psycopg.Error as e:
        print(f"[FATAL] 连接失败: {e}", file=sys.stderr)
        return 1

    try:
        while True:
            with conn.cursor() as cur:
                cur.execute(SQL)
                rows = cur.fetchall()
            conn.commit()  # 释放快照
            render(rows)
            if args.once:
                break
            time.sleep(args.interval)
    except KeyboardInterrupt:
        print("\n[bye]")
    finally:
        conn.close()
    return 0


if __name__ == "__main__":
    sys.exit(main())
markdown
# 第 15 章 配套代码 · 复制与高可用

## 准备工作

1.`psql -h 127.0.0.1 -U postgres -d learn_pg -f ../init.sql` 在主库创建 `ch15_orders / ch15_order_items` 等业务表 + `ch15_pub_orders` 发布 + `ch15_standby1_slot` 物理槽
2. 安装依赖:`pip install "psycopg[binary]>=3.1"`
3. 02 / 04 是 docker 化的「一主一从」搭建脚本,需要本机能跑 `docker`,并且 5432 / 5433 端口空闲
4. (可选)通过环境变量覆盖默认连接:
   - `export PG_DSN="host=... port=... dbname=... user=..."`(01、05)
   - `export MASTER_DSN=... REPLICA_DSN=... SUB_CONN=...`(03 用)

> ⚠️ **本章很多脚本需要两个 PG 实例**(一主一从)。流复制 / 逻辑复制都做不了「单实例自演」。`02_setup_streaming_replication.sh up` 会自动起两个 docker 容器,03 / 04 / 05 在此基础上演示。

## 脚本一览(推荐运行顺序)

| 脚本 | 一句话说明 | 关键 PG 特性 |
|------|------------|--------------|
| `01_show_replication_status.py` | 单纯读 `pg_stat_replication / pg_replication_slots / pg_publication`,看主库上挂了哪些从库 | `pg_stat_replication``pg_walfile_name` |
| `02_setup_streaming_replication.sh` | docker 一键起一主一从;子命令 `up/status/write/read/failover/down` | `pg_basebackup -R -S``hot_standby` |
| `03_logical_replication.py` | 主从两个独立实例之间的逻辑复制;`setup/write/verify/status/teardown` | `CREATE PUBLICATION / SUBSCRIPTION` |
| `04_promote_failover.sh` | 演示「主库挂了 → `pg_promote` 提升 → `pg_rewind` 让旧主回来」 | `pg_promote``pg_rewind` |
| `05_replication_lag.py` | 持续监控复制延迟(字节 + 时间双维度) | `pg_wal_lsn_diff``replay_lag` |

## 运行示例

```bash
bash 02_setup_streaming_replication.sh up
bash 02_setup_streaming_replication.sh status
bash 02_setup_streaming_replication.sh write
bash 02_setup_streaming_replication.sh read

python3 05_replication_lag.py --interval 1
bash 04_promote_failover.sh demo
bash 02_setup_streaming_replication.sh down

预期输出

05_replication_lag.py 输出每 2 秒刷新一次:

[14:23:01]  1 个从库连接
  app                  addr             state      sync       send      write     flush     replay    total | wlag   flag   rlag
  ----------------------------------------------------------------------------------------------------
  walreceiver          172.18.0.3       streaming  async      0.0B      0.0B      0.0B      0.0B      0.0B  | 12ms   18ms   24ms

常见报错与依赖

  • connection refused → PG 未启动,或 docker 网络没建好
  • relation "ch15_orders" does not exist → 没跑 init.sql
  • WARNING: out of replication slots → 调高 max_replication_slots
  • pg_basebackup: FATAL: number of requested standby connections exceeds max_wal_senders → 调高 max_wal_senders
  • pg_rewind 报「源端缺少 WAL」→ 旧主未开 wal_log_hintsdata_checksums,回退方案是删 data 目录重新 pg_basebackup
  • 03 需要 5433 端口启一个独立 PG 实例(不是物理 standby),自行用 docker 起:docker run -d -p 5433:5432 -e POSTGRES_PASSWORD=postgres postgres:16 -c wal_level=logical
  • 04 / pg_promote 后旧主再回来必须先 pg_rewind 或重新基础备份,否则会出现 timeline 分叉

01_show_replication_status.py ↗ · 02_setup_streaming_replication.sh ↗ · 03_logical_replication.py ↗ · 04_promote_failover.sh ↗ · 05_replication_lag.py ↗ · README.md ↗