PostgreSQL 逻辑复制集群搭建指北:基于 Docker 的 publisher + subscriber 一键拉起


难度 中等
编写人 编写内容 编写时间
growdu 初稿,使用 docker.m.daocloud.io/library/postgres:16docker compose 起 publisher / subscriber,跑通 INSERT / UPDATE / DELETE 的实时同步 2026-09-01

本文是「PostgreSQL 逻辑复制系列」的第 N 篇,重点不在原理拆解,而在「用一条 docker compose up -d 把 publisher/subscriber 拉起来」。同系列前文:

PostgreSQL 的逻辑复制(Logical Replication)自 10.0 起成为内核能力,基于 发布(Publication)/ 订阅(Subscription) 模型,把发布端表的 WAL 解码为逻辑变更(pgoutput 插件),由订阅端 apply worker 回放。

背景

它与流复制的区别在于:

  • 粒度细到表、甚至 DML 操作类型(INSERT/UPDATE/DELETE/TRUNCATE 可独立订阅);
  • 可以跨大版本(如 14 → 16);
  • 既能整库级同步,也能构建多主、级联、选择性分发表;
  • 默认是 单向,反向回环需要用 originREPLICA IDENTITY FULL 等手段处理。

下面这套配置用 docker.m.daocloud.io/library/postgres:16 镜像在一台主机上起两个 PG 实例,做一个能跑通的最小化逻辑复制集群,并把 publisher / subscriber 都做成由 docker compose 一键拉起的服务。

总体方案

+------------------+    logical WAL (pgoutput)    +------------------+
|  pg-publisher    |  --------------------------->|  pg-subscriber   |
|  hostname:       |     via replication slot     |  hostname:       |
|  pg-publisher    |     sub_slot_app             |  pg-subscriber   |
|  port 5432       |                              |  port 5432       |
|  DB: appdb       |                              |  DB: appdb       |
|  PUBLICATION:    |                              |  SUBSCRIPTION:   |
|  pub_app         |                              |  sub_app         |
+------------------+                              +------------------+

要点:

维度 publisher subscriber
wal_level 必须 logical 同上(作为上游保留能力)
max_replication_slots ≥ 订阅端数量 不强依赖
max_wal_senders ≥ 订阅端 × 工作进程 同上
max_logical_replication_workers ≥ 同时复制表数 apply worker + tablesync worker
表结构 必须存在并包含主键 必须预先存在copy_data=true 时尤其重要)
复制账号 REPLICATION 属性 用同一账号连接

docker-compose.yml

工作目录建议放在 /tmp/pg-logical-replication,文件如下:

services:
  publisher:
    image: docker.m.daocloud.io/library/postgres:16
    container_name: pg-publisher
    hostname: pg-publisher
    environment:
      POSTGRES_USER: repuser
      POSTGRES_PASSWORD: reppass@2026
      POSTGRES_DB: appdb
    command:
      - "postgres"
      - "-c"
      - "wal_level=logical"
      - "-c"
      - "max_replication_slots=10"
      - "-c"
      - "max_wal_senders=10"
      - "-c"
      - "max_logical_replication_workers=10"
      - "-c"
      - "max_worker_processes=20"
      - "-c"
      - "shared_preload_libraries="
      - "-c"
      - "listen_addresses=*"
    ports:
      - "5433:5432"
    volumes:
      - publisher_data:/var/lib/postgresql/data
      - ./init-publisher.sh:/docker-entrypoint-initdb.d/01-init.sh:ro
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U repuser -d appdb"]
      interval: 5s
      timeout: 5s
      retries: 10

  subscriber:
    image: docker.m.daocloud.io/library/postgres:16
    container_name: pg-subscriber
    hostname: pg-subscriber
    environment:
      POSTGRES_USER: repuser
      POSTGRES_PASSWORD: reppass@2026
      POSTGRES_DB: appdb
      PUBLISHER_HOST: pg-publisher
    command:
      - "postgres"
      - "-c"
      - "wal_level=logical"
      - "-c"
      - "max_replication_slots=10"
      - "-c"
      - "max_logical_replication_workers=10"
      - "-c"
      - "max_worker_processes=20"
      - "-c"
      - "shared_preload_libraries="
      - "-c"
      - "listen_addresses=*"
    ports:
      - "5434:5432"
    volumes:
      - subscriber_data:/var/lib/postgresql/data
      - ./init-subscriber.sh:/docker-entrypoint-initdb.d/01-init.sh:ro
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U repuser -d appdb"]
      interval: 5s
      timeout: 5s
      retries: 10
    depends_on:
      publisher:
        condition: service_healthy

volumes:
  publisher_data:
  subscriber_data:

几个工程上的细节:

  • daocloud.io 镜像源直连,无需额外 registry 代理;
  • listen_addresses=* 让同一 compose 网络里的容器能远程连入;
  • healthcheck 使用 pg_isready,让 depends_on.condition: service_healthy 真正生效;
  • 两个容器分别映射到宿主机 5433 / 5434,避免本地有 PG 时的端口冲突。

publisher 初始化脚本

init-publisher.sh 由 docker-entrypoint 在首次启动自动执行,做 5 件事:

  1. repuserREPLICATION 属性;
  2. appdb 的 DML 权限授予 repuser
  3. 创建测试用的 app_user / app_order 表与种子数据;
  4. pg_hba.conf 上为 repuser 放行 scram-sha-256(包含 replication 段);
  5. 创建 pub_app PUBLICATION(FOR TABLE 显式列表,也可以改成 FOR ALL TABLES)。
#!/usr/bin/env bash
set -euo pipefail

echo "[publisher] init roles, schema, publication"

# 1) 复制账号 & 库内权限
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
ALTER ROLE repuser WITH REPLICATION PASSWORD 'reppass@2026';
GRANT ALL PRIVILEGES ON DATABASE appdb TO repuser;
GRANT ALL ON ALL TABLES    IN SCHEMA public TO repuser;
GRANT ALL ON ALL SEQUENCES IN SCHEMA public TO repuser;
ALTER DEFAULT PRIVILEGES IN SCHEMA public GRANT ALL ON TABLES    TO repuser;
ALTER DEFAULT PRIVILEGES IN SCHEMA public GRANT ALL ON SEQUENCES TO repuser;
SQL

# 2) 表结构 (必须先于 CREATE PUBLICATION)
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
CREATE TABLE IF NOT EXISTS app_user (
    id          BIGSERIAL    PRIMARY KEY,
    username    VARCHAR(64)  NOT NULL UNIQUE,
    email       VARCHAR(128) NOT NULL,
    created_at  TIMESTAMPTZ  NOT NULL DEFAULT now()
);

CREATE TABLE IF NOT EXISTS app_order (
    id           BIGSERIAL    PRIMARY KEY,
    user_id      BIGINT       NOT NULL REFERENCES app_user(id),
    amount       NUMERIC(12,2) NOT NULL,
    remark       TEXT,
    created_at   TIMESTAMPTZ  NOT NULL DEFAULT now()
);
SQL

# 3) 种子数据
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
INSERT INTO app_user(username, email) VALUES
    ('alice', 'alice@example.com'),
    ('bob',   'bob@example.com'),
    ('carol', 'carol@example.com')
ON CONFLICT DO NOTHING;
SQL

# 4) 让远端复制流量走 scram-sha-256
echo "host    all             repuser        0.0.0.0/0               scram-sha-256" >> /var/lib/postgresql/data/pg_hba.conf
echo "host    replication     repuser        0.0.0.0/0               scram-sha-256" >> /var/lib/postgresql/data/pg_hba.conf
psql -v ON_ERROR_STOP=1 -U repuser -d appdb -c "SELECT pg_reload_conf();"

# 5) 创建发布
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
DROP PUBLICATION IF EXISTS pub_app;
CREATE PUBLICATION pub_app FOR TABLE app_user, app_order;
SQL

echo "[publisher] publication:"
psql -U repuser -d appdb -c "SELECT pubname, puballtables FROM pg_publication;"
echo "[publisher] tables:"
psql -U repuser -d appdb -c "\dt"

注意:docker-entrypoint.sh 会保证 SQL 脚本只在 数据目录为空 时跑一次。如果你想重做一次发布,删卷即可:docker compose down -v

subscriber 初始化脚本

subscriber 端要做两件事:

  1. 预先创建空表(列、主键、外键、默认值都要齐全,因为初始 copy_data 通过 INSERT … SELECT 把 publisher 的内容灌进来);
  2. CREATE SUBSCRIPTION 指向 publisher,开始拉取逻辑变更。
#!/usr/bin/env bash
set -euo pipefail

PUBLISHER_HOST="${PUBLISHER_HOST:-pg-publisher}"
echo "[subscriber] waiting for publisher at ${PUBLISHER_HOST}:5432 ..."
for i in $(seq 1 60); do
    if pg_isready -h "${PUBLISHER_HOST}" -p 5432 -U repuser -d appdb >/dev/null 2>&1; then
        echo "[subscriber] publisher is up after ${i} tries"
        break
    fi
    sleep 2
done

# 本实例的本地权限
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
GRANT ALL PRIVILEGES ON DATABASE appdb TO repuser;
SQL

# 关键: subscriber 端必须预先存在同名表结构, 才能承接初始 copy_data 同步
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<'SQL'
CREATE TABLE IF NOT EXISTS app_user (
    id          BIGSERIAL    PRIMARY KEY,
    username    VARCHAR(64)  NOT NULL UNIQUE,
    email       VARCHAR(128) NOT NULL,
    created_at  TIMESTAMPTZ  NOT NULL DEFAULT now()
);

CREATE TABLE IF NOT EXISTS app_order (
    id           BIGSERIAL    PRIMARY KEY,
    user_id      BIGINT       NOT NULL REFERENCES app_user(id),
    amount       NUMERIC(12,2) NOT NULL,
    remark       TEXT,
    created_at   TIMESTAMPTZ  NOT NULL DEFAULT now()
);
SQL

export PGPASSWORD='reppass@2026'

# 创建逻辑复制订阅 (核心步骤)
psql -v ON_ERROR_STOP=1 -U repuser -d appdb <<SQL
DROP SUBSCRIPTION IF EXISTS sub_app;
CREATE SUBSCRIPTION sub_app
    CONNECTION 'host=${PUBLISHER_HOST} port=5432 user=repuser password=reppass@2026 dbname=appdb'
    PUBLICATION pub_app
    WITH (copy_data = true, create_slot = true, enabled = true, slot_name = sub_slot_app);
SQL

echo "[subscriber] subscription:"
psql -U repuser -d appdb -c "SELECT subname, subenabled, subslotname, subpublications FROM pg_subscription;"

CREATE SUBSCRIPTION 末尾的几个开关很重要:

  • copy_data = true:先做一次全量初始同步;
  • create_slot = true:自动在 publisher 上建一个名字为 sub_slot_app 的 replication slot;
  • enabled = true:开启 apply worker;
  • slot_name = …:固定复制槽名,方便后续运维诊断。

部署与验证

启动

cd /tmp/pg-logical-replication
docker compose up -d
docker compose ps

启动后大约 5–10 秒两端 healthcheck 全绿。

元数据核对

Publisher 端:

docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "SELECT * FROM pg_publication;"
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "SELECT slot_name, plugin, slot_type, database, active, restart_lsn FROM pg_replication_slots;"
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "SELECT pid, usename, application_name, client_addr, state, sync_state, write_lsn, replay_lsn FROM pg_stat_replication;"

我在部署中输出的结果(节选):

  oid  | pubname | pubowner | puballtables | pubinsert | pubupdate | pubdelete | pubtruncate | pubviaroot
-------+---------+----------+--------------+-----------+-----------+-----------+-------------+------------
 16412 | pub_app |       10 | f            | t         | t         | t         | t           | f
  slot_name   |  plugin  | slot_type | database | active | restart_lsn
--------------+----------+-----------+----------+--------+-------------
 sub_slot_app | pgoutput | logical   | appdb    | t      | 0/1979E08
 pid | usename | application_name | client_addr |   state   | sync_state | write_lsn | replay_lsn
-----+---------+------------------+-------------+-----------+------------+-----------+------------
  96 | repuser | sub_app          | 172.21.0.3  | streaming | async      | 0/1979E40 | 0/1979E40

Subscriber 端:

docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT * FROM pg_subscription;"
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT pid, relid::regclass AS table, received_lsn, last_msg_send_time, last_msg_receipt_time FROM pg_stat_subscription;"
 oid  | subdbid | subskiplsn | subname | subowner | subenabled | subtwophasestate | subpublications
-------+---------+------------+---------+----------+------------+------------------+-----------------
 16410 |   16384 | 0/0        | sub_app |       10 | t          | d                | {pub_app}
 pid | table | received_lsn |      last_msg_send_time       |     last_msg_receipt_time
-----+-------+--------------+-------------------------------+-------------------------------
  82 |       | 0/1979E40    | 2026-08-31 08:59:26.799086+00 | 2026-08-31 08:59:26.799119+00

pg_stat_subscriptionlast_msg_send_timelast_msg_receipt_time 在持续更新,received_lsn 与 publisher 端 replay_lsn 同步推进,说明 apply worker 工作正常。

端到端 DML 验证

# 1) INSERT
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "INSERT INTO app_user(username, email) VALUES('dave', 'dave@example.com');"
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT id, username, email FROM app_user WHERE username='dave';"

# 2) UPDATE
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "UPDATE app_user SET email='alice.new@example.com' WHERE username='alice';"
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT email FROM app_user WHERE username='alice';"

# 3) DELETE
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "DELETE FROM app_user WHERE username='carol';"
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT count(*) FROM app_user WHERE username='carol';"

# 4) INSERT app_order
docker exec -e PGPASSWORD=reppass@2026 pg-publisher \
    psql -U repuser -d appdb -c "INSERT INTO app_order(user_id, amount, remark) VALUES(2, 299.50, 'order-dave');"
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT id, user_id, amount, remark FROM app_order;"

# 5) 复制延迟
docker exec -e PGPASSWORD=reppass@2026 pg-subscriber \
    psql -U repuser -d appdb -c "SELECT now() - last_msg_receipt_time AS apply_lag FROM pg_stat_subscription;"

部署实测里,两端 app_user / app_order 在几秒后内容完全一致:

[publisher]
 id | username |         email
----+----------+-----------------------
  1 | alice    | alice.new@example.com
  2 | bob      | bob@example.com
  4 | frank    | frank@example.com
  5 | dave     | dave@example.com
(4 rows)

 id | user_id | amount |   remark
----+---------+--------+------------
  1 |       2 | 299.50 | order-dave

[subscriber]
 id | username |         email
----+----------+-----------------------
  1 | alice    | alice.new@example.com
  2 | bob      | bob@example.com
  4 | frank    | frank@example.com
  5 | dave     | dave@example.com
(4 rows)

 id | user_id | amount |   remark
----+---------+--------+------------
  1 |       2 | 299.50 | order-dave

    apply_lag
-----------------
 00:00:02.307543

注意:不要在 subscriber 端直接执行 DML。复制过来的变更与本地 BIGSERIAL 会冲突:

ERROR:  duplicate key value violates unique constraint "app_user_pkey"
DETAIL:  Key (id)=(1) already exists.

如果非要写,需要把 subscriber 端的 identity / sequence 与 publisher 解耦,或建立双向逻辑复制 + origin 过滤。

常用排查命令

-- publisher: 槽是否活跃,是否堆积
SELECT slot_name, plugin, slot_type, database, active, restart_lsn
  FROM pg_replication_slots;

-- publisher: walsender 是否在流式发送
SELECT pid, usename, application_name, client_addr, state, sync_state,
       write_lsn, replay_lsn, (now() - backend_start) AS dur
  FROM pg_stat_replication;

-- subscriber: 订阅是否启用、远端 LSN
SELECT subname, subenabled, subslotname, subpublications
  FROM pg_subscription;

SELECT pid, relid::regclass AS table,
       received_lsn, last_msg_send_time, last_msg_receipt_time,
       (now() - last_msg_receipt_time) AS apply_lag
  FROM pg_stat_subscription;

-- subscriber: 单表同步状态
SELECT * FROM pg_subscription_rel;

-- subscriber 主动同步进度诊断(最直接)
SELECT s.subname,
       sr.srsubstate,
       sr.srrelid::regclass
  FROM pg_subscription s
  JOIN pg_subscription_rel sr ON sr.srsubid = s.oid;

srsubstate 取值含义:

  • i = initialize(订阅尚未开始向该表推数据);
  • d = data being copied(初始 copy_data 进行中);
  • s = synchronized(已同步完成,正在等后续变更);
  • r = ready(复制状态最终一致,正在持续接收变更)。

常见坑

  1. subscriber 端没有同名表结构CREATE SUBSCRIPTION + copy_data=true 不会自动建表,只会尝试 INSERT … SELECT。会直接报 relation "app_user" does not exist。修复:在 subscriber 上预建表。
  2. **wal_level 不等于 logical**:物理复制(流复制)只需要 replica,发布订阅链路必须 logical,否则 CREATE PUBLICATIONwal_level is not logical
  3. max_replication_slots 不足:每多一个订阅都要占一个槽;扩展新订阅前要规划好。
  4. 删除 / disable 顺序:先 ALTER SUBSCRIPTION … DISABLE; 再删,否则远端 walsender 还会继续往拉。
  5. REPLICA IDENTITY 默认是 DEFAULT(=主键):对于无主键表,要么 ALTER TABLE … REPLICA IDENTITY FULL;,要么在 PUBLICATION 里 WHERE 过滤掉,但代价是 WAL 体积明显变大。
  6. DDL 不通过逻辑复制:表结构变更要单独在两端各执行。社区方案 pglogical 风格的双向 DDL 不能裸用纯 pgoutput

publisher 端 DSN

部署完成后,发布端连接串如下:

# 容器内 / 同一 docker-compose 网络
postgresql://repuser:reppass%402026@pg-publisher:5432/appdb

# 宿主机(端口已映射为 5433)
postgresql://repuser:reppass%402026@127.0.0.1:5433/appdb

# key=value 风格
host=pg-publisher port=5432 user=repuser password=reppass@2026 dbname=appdb

注意密码里的 @ 需要百分号编码为 %40。一行自检:

PGPASSWORD='reppass@2026' psql \
  "postgresql://repuser:reppass@2026@127.0.0.1:5433/appdb" \
  -c "SELECT 1 AS ok, current_setting('wal_level') AS wal_level;"

预期输出:

 ok | wal_level
----+-----------
  1 | logical

wal_level = logical 说明拿到的是真正的发布端。

小结

  • docker.m.daocloud.io/library/postgres:16 + docker compose,可以把一个 publisher + subscriber 的逻辑复制集群拉起来不超过 10 秒,并跑通 INSERT / UPDATE / DELETE 的实时同步;
  • 关键参数都来自 wal_level=logicalmax_replication_slotsmax_wal_sendersmax_logical_replication_workers
  • subscriber 端 必须先有同名表,初始 copy_data 才能正确灌入;
  • 想要更进一步(双向、级联、selective replication)时,再叠加 origin / WHERE / 多 PUBLICATION 即可,框架本身已经比较稳固。

同系列前文


文章作者: growdu
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 growdu !
  目录
分类导航
随笔2 AI27 算法1 计算机基础13 博客搭建7 ChatGPT2 集群63 计算机通信1 数据库34 数据库深入80 DPDK26 Docker11 Elasticsearch4 编辑工具4 FAQ1 Go Web1 hometown2 编程语言16 网络9 OPC1 Linux38 openGauss4 页面12 PostgreSQL54 程序员自我修养1 协议11 成长之路1 stock1 存储5 工具20 VPP18 视频作品1 Vue13 Web1 代码示例11 数据库15 BenchmarkSQL1 PostgreSQL 源码修炼之路14
最热文章
1
13 逻辑复制深入
数据库深入🔥 1570
2
0 Postgresql存储、索引及系统优化、主备切换
PostgreSQL🔥 1495
3
一文读懂openguass dcf网络模块
集群🔥 1420
4
逻辑复制源码分析
数据库深入🔥 1327
5
PostgreSQL 分区表:从一行 `PARTITION BY` 到路由热路径的全链路拆解
数据库🔥 1094
6
applyparallelworker.c 之 LA 端源码深度解析:Leader Apply Worker 的指挥中枢
数据库深入🔥 1082
7
PostgreSQL Background Worker 全解:从 `RegisterBackgroundWorker` 到逻辑复制 4 类 worker 的全生命周期
数据库🔥 1078
8
PostgreSQL的后台进程walsender分析 - 关系型数据库 - 亿速云
PostgreSQL🔥 1033
9
PostgreSQL 逻辑复制的监控:六张视图 + 一组可执行 SQL,把 publisher/subscriber 的速率与健康度彻底看透
数据库🔥 1032
10
PostgreSQL 逻辑复制支持 DDL 之后:DDL 与 DML 的时序难题(重点:分区表)
数据库🔥 999
11
reorderbuffer.c 源码深度解析:PostgreSQL 逻辑复制的"事务重组引擎
数据库深入🔥 953
12
PostgreSQL 内核开发:读取一张表的 9 步标准流程与缓存全景
数据库🔥 938
13
从 `postgres` 二进制到生产级守护 —— PostgreSQL 最外层模块与启动全流程拆解
数据库🔥 936
14
支持逻辑复制同步 DDL 适配 SQL Server 方案
数据库深入🔥 934
15
PostgreSQL 逻辑复制的 ReorderBuffer 与事务机制:从一行 WAL 到一致性变更流的全链路绑定
数据库🔥 913
16
DDL同步架构(美化版)
数据库深入🔥 908
17
PostgreSQL Latch 机制详解:从一行 SetLatch 到 epoll 的内核之旅
数据库🔥 871
18
pgbench 源码全解:一个 C 文件如何撑起 PostgreSQL 官方压测工具
数据库🔥 860
19
PostgreSQL libpq 机制与缓冲区详解
数据库🔥 850
20
PostgreSQL 逻辑复制 spill 文件深度剖析:从 `xid-*.spill` 到 TPC-C 的增长方程
数据库🔥 845