Skip to content

第 16 章 分区与分库分表:让 PG 装下亿级数据

目标读者:单表数据量已经过亿,查询变慢、备份变难,但又不想立刻上分布式数据库的同学。

学完你会:能用 PG 声明式分区把一张大表拆成 N 个小表 + 让查询自动只扫到相关分区 + 用 pg_partman 自动维护时序分区 + 看懂 Citus / FDW 的「分库分表」思路


0. 导读:为什么要分区 / 分表 / 分库?

一张表数据少时怎么写都行;一旦行数飙到 1 亿、几亿,问题接踵而至:

  • 查询变慢:扫描全表慢,索引也变成几十 GB;
  • VACUUM 变长:单次 vacuum 全表要几小时;
  • DDL 阻塞ALTER TABLE 加个字段要锁表一整夜;
  • 备份变重pg_dump 一次几百 GB;
  • 删历史数据DELETE WHERE created_at < ... 慢且产生大量 dead tuple;
  • 磁盘满:一台机器装不下了。

解决思路有三层:

            数据规模 ──→ 越来越大
   单表 ──→ 分区表 ──→ 分库分表 ──→ 分布式数据库
   (PG)    (PG 内置)   (Citus/FDW)   (Spanner/CockroachDB)

1. 分区 vs 分表 vs 分库:先把概念分清

名词物理位置逻辑位置应用感知谁来路由
分区(Partition)同一台机器、同一个 PG 实例同一个「逻辑表」,PG 透明拆成多个 partition不感知(select 一张表名)PG 自动
分表(Sharding @ 单库)同一个 PG 实例多张物理表 ch16_orders_2025_01ch16_orders_2025_02应用要拼表名应用代码
分库(Sharding @ 多机)多台机器、多个 PG 实例多张物理表分散在不同机器上不感知(中间件路由)中间件 / 协调节点

📌 与 MySQL 的区别

  • MySQL 8 的 PARTITION BY ≈ PG 的「分区」;
  • MySQL 「分库分表」通常靠 ShardingSphere、MyCAT、Vitess;
  • PG 的「分库分表」官方推荐路线是 Citus(已被微软收购,开源版仍可用)。

PG 10 之前用「继承表 + 触发器」实现分区(pg_pathman 扩展),PG 10 起内置「声明式分区(Declarative Partitioning)」,从此分区成为一等公民。


2. PG 10+ 声明式分区:三种类型

声明式分区的三大要素:

  1. 分区父表CREATE TABLE ... PARTITION BY ...不存数据,只是个虚表;
  2. 分区键(partition key):父表上指定的「按哪一列拆」;
  3. 分区子表CREATE TABLE ... PARTITION OF parent FOR VALUES ...真正存数据

2.1 RANGE 分区:按值区间(最常用)

典型场景:按时间分区(每月一个分区,方便归档)。

sql
-- 父表
CREATE TABLE ch16_orders (
    id          BIGSERIAL,
    user_id     BIGINT NOT NULL,
    amount      NUMERIC(12,2) NOT NULL,
    status      TEXT,
    created_at  TIMESTAMPTZ NOT NULL,
    PRIMARY KEY (id, created_at)        -- 主键必须包含分区键!
) PARTITION BY RANGE (created_at);

-- 子分区(左闭右开)
CREATE TABLE ch16_orders_2025_01 PARTITION OF ch16_orders
    FOR VALUES FROM ('2025-01-01') TO ('2025-02-01');

CREATE TABLE ch16_orders_2025_02 PARTITION OF ch16_orders
    FOR VALUES FROM ('2025-02-01') TO ('2025-03-01');

CREATE TABLE ch16_orders_2025_03 PARTITION OF ch16_orders
    FOR VALUES FROM ('2025-03-01') TO ('2025-04-01');

-- 兜底分区(任何不匹配的值落到这里)
CREATE TABLE ch16_orders_default PARTITION OF ch16_orders DEFAULT;

插入和查询完全像普通表

sql
INSERT INTO ch16_orders (user_id, amount, status, created_at)
VALUES (1, 99.99, 'paid', '2025-01-15');
-- PG 自动路由到 ch16_orders_2025_01

SELECT * FROM ch16_orders WHERE created_at >= '2025-02-01' AND created_at < '2025-03-01';
-- 只扫 ch16_orders_2025_02 一个分区

左闭右开FROM ('2025-01-01') TO ('2025-02-01') 包含 1 月 1 日,不含 2 月 1 日,这样相邻分区不会重叠。

2.2 LIST 分区:按枚举值

典型场景:按地区、按渠道、按业务线。

sql
CREATE TABLE users (
    id      BIGSERIAL,
    region  TEXT NOT NULL,
    name    TEXT,
    PRIMARY KEY (id, region)
) PARTITION BY LIST (region);

CREATE TABLE users_cn   PARTITION OF users FOR VALUES IN ('CN');
CREATE TABLE users_us   PARTITION OF users FOR VALUES IN ('US', 'CA');
CREATE TABLE users_eu   PARTITION OF users FOR VALUES IN ('DE', 'FR', 'UK', 'NL');
CREATE TABLE users_other PARTITION OF users DEFAULT;

LIST 分区一个值只能落到一个分区(不像 RANGE 可以连续)。

2.3 HASH 分区:按哈希均匀打散(PG 11+)

典型场景:值分布不均匀,想强制均匀分布;或要按 user_id 这种高基数列分散读写压力。

sql
CREATE TABLE ch16_events (
    id          BIGSERIAL,
    user_id     BIGINT NOT NULL,
    event_type  TEXT,
    payload     JSONB,
    PRIMARY KEY (id, user_id)
) PARTITION BY HASH (user_id);

-- 4 个分区,模数 4
CREATE TABLE ch16_events_p0 PARTITION OF ch16_events FOR VALUES WITH (MODULUS 4, REMAINDER 0);
CREATE TABLE ch16_events_p1 PARTITION OF ch16_events FOR VALUES WITH (MODULUS 4, REMAINDER 1);
CREATE TABLE ch16_events_p2 PARTITION OF ch16_events FOR VALUES WITH (MODULUS 4, REMAINDER 2);
CREATE TABLE ch16_events_p3 PARTITION OF ch16_events FOR VALUES WITH (MODULUS 4, REMAINDER 3);

原理:PG 用内部 hash 函数 hash_user_id 算出哈希值,再 % 4 决定落到 p0~p3 中的哪个。

⚠️ HASH 分区的坑

  • 不支持「按时间归档」,因为不是按值连续分区;
  • 调整分区数(从 4 改成 8)需要重建所有分区;
  • WHERE user_id = 100 能裁剪到一个分区;WHERE user_id BETWEEN 100 AND 200 裁剪不了(因为哈希后不连续)。

2.4 多级分区:分区下再分区

sql
-- 一级 RANGE 按月
CREATE TABLE logs (
    id BIGSERIAL,
    region TEXT NOT NULL,
    created_at TIMESTAMPTZ NOT NULL,
    msg TEXT,
    PRIMARY KEY (id, region, created_at)
) PARTITION BY RANGE (created_at);

CREATE TABLE logs_2025_01 PARTITION OF logs
    FOR VALUES FROM ('2025-01-01') TO ('2025-02-01')
    PARTITION BY LIST (region);                 -- 二级 LIST 按地区

CREATE TABLE logs_2025_01_cn PARTITION OF logs_2025_01 FOR VALUES IN ('CN');
CREATE TABLE logs_2025_01_us PARTITION OF logs_2025_01 FOR VALUES IN ('US');

适合:先按时间筛 → 再按地区筛 的查询模式。


3. 默认分区 DEFAULT:兜底未匹配的值

sql
CREATE TABLE ch16_orders_default PARTITION OF ch16_orders DEFAULT;

作用:任何不匹配的分区键值都进 default 分区,避免插入失败

⚠️ 代价

  • 添加新分区时,PG 会扫描整个 default 分区确认没有冲突的行,加 ACCESS EXCLUSIVE 锁,default 越大越慢;
  • default 分区无法被 partition pruning 利用(每个查询都要扫它)。

最佳实践:把 default 分区当「告警用」,定期清空 / 检查里面有没有数据,正常情况 default 应该是空的。


4. 分区裁剪 Partition Pruning

分区裁剪(Partition Pruning):查询计划器根据 WHERE 条件,自动跳过不可能命中的分区,只扫相关分区。

PG 默认开启:

sql
SHOW enable_partition_pruning;   -- on

4.1 计划期裁剪(Planning-Time Pruning)

WHERE 是字面量常量时,规划器在生成计划时直接剪掉无关分区:

sql
EXPLAIN
SELECT * FROM ch16_orders
WHERE created_at >= '2025-02-01' AND created_at < '2025-03-01';

输出(PG 16):

 Append  (cost=0.00..XX rows=YY width=ZZ)
   ->  Seq Scan on ch16_orders_2025_02 orders_1
         Filter: ...

只看到 ch16_orders_2025_02,其他分区根本没出现在计划里。

4.2 执行期裁剪(Execution-Time Pruning,PG 11+)

WHERE 是参数化的(PREPARE / 函数参数 / Subquery),规划器没法在编译期知道值,但执行期可以裁:

sql
PREPARE q (date) AS
    SELECT * FROM ch16_orders WHERE created_at >= $1 AND created_at < $1 + INTERVAL '1 month';

EXPLAIN (ANALYZE)
EXECUTE q ('2025-02-01');

输出会出现 Subplans Removed: N,表示执行期剪掉了 N 个分区:

 Append  (...)
   Subplans Removed: 5
   ->  Seq Scan on ch16_orders_2025_02 orders_1

4.3 怎样的 WHERE 才能触发裁剪?

WHERE是否裁剪
created_at >= '2025-02-01'
created_at = '2025-02-15'
created_at IN ('2025-02-01', '2025-03-01')✅(PG 11+)
created_at + INTERVAL '1 day' >= '...'❌(左侧加表达式破坏)
extract(month FROM created_at) = 2❌(同上)
created_at >= now() - INTERVAL '7 day'✅(PG 12+)

规则:WHERE 里的分区键必须直接出现在比较操作符的一侧,不能套表达式。

4.4 验证裁剪效果

sql
-- 1000 万行 + 12 个月分区的实测
EXPLAIN (ANALYZE, BUFFERS)
SELECT count(*) FROM ch16_orders WHERE created_at >= '2025-06-01' AND created_at < '2025-07-01';

-- 没裁剪:扫 1000 万行,3.5s
-- 裁剪后:只扫 1 个分区 80 万行,280ms  → 12 倍提速

5. 分区表上的索引(PG 11+ 索引级联)

PG 11 起,在父表上建索引会自动级联到所有子分区:

sql
CREATE INDEX ON ch16_orders (created_at);
CREATE INDEX ON ch16_orders (user_id);

PG 自动在每个分区上建同名索引,未来新建的分区也会自动建。

老版本(PG 10)的痛:必须在每个分区上手动建索引,加新分区时也要记得建。

5.1 唯一索引 / 主键的限制

主键 / 唯一索引必须包含所有分区键。原因:PG 没法跨分区查找冲突,只能在每个分区内保证唯一,所以必须把分区键放进主键,分区内唯一才能推出全表唯一。

sql
-- ❌ 不行:主键不包含分区键 created_at
CREATE TABLE ch16_orders (
    id BIGSERIAL PRIMARY KEY,
    created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);
-- ERROR: unique constraint on partitioned table must include all partitioning columns

-- ✅ 行:主键 = (id, created_at)
CREATE TABLE ch16_orders (
    id BIGSERIAL,
    created_at TIMESTAMPTZ NOT NULL,
    PRIMARY KEY (id, created_at)
) PARTITION BY RANGE (created_at);

5.2 外键

PG 11+ 支持「分区表作为外键的引用方」(即 REFERENCES partitioned_table); PG 12+ 支持「分区表作为外键的来源方」(即 FOREIGN KEY ... REFERENCES other_table)。


6. 分区维护

6.1 添加新分区(最常见运维操作)

sql
-- 月底之前,预先建好下个月的分区
CREATE TABLE ch16_orders_2025_04 PARTITION OF ch16_orders
    FOR VALUES FROM ('2025-04-01') TO ('2025-05-01');

没有提前建会怎样?

  • 有 default 分区:数据进 default,但 default 一旦有数据,未来添加新分区会很慢(要锁 default 扫一遍)。
  • 没有 default 分区:插入直接报错 no partition of relation "ch16_orders" found for row

结论必须用脚本提前建未来分区!自己写 cron 或用 pg_partman(下一节)。

6.2 删除老分区

sql
-- 方法 1:直接 DROP(最快)
DROP TABLE ch16_orders_2024_01;

-- 方法 2:先 DETACH 再处理(推荐:可以先归档再删)
ALTER TABLE ch16_orders DETACH PARTITION ch16_orders_2024_01;
-- 此时 ch16_orders_2024_01 变成普通表,可以归档到冷存储

PG 14+ 的并发 DETACH(不锁表!)

sql
ALTER TABLE ch16_orders DETACH PARTITION ch16_orders_2024_01 CONCURRENTLY;

CONCURRENTLY 模式只取 SHARE UPDATE EXCLUSIVE 锁,不会阻塞 SELECT/INSERT,巨大业务库的运维利器。

6.3 ATTACH:把已有表挂上去

把一张「独立表」转成「子分区」:

sql
-- 1. 先建一张独立表(甚至可以从其他系统导入)
CREATE TABLE ch16_orders_2024_12 (LIKE ch16_orders INCLUDING ALL);
INSERT INTO ch16_orders_2024_12 SELECT * FROM external_source;

-- 2. 加 CHECK 约束(让 PG 不必扫描全表验证)
ALTER TABLE ch16_orders_2024_12 ADD CONSTRAINT chk
    CHECK (created_at >= '2024-12-01' AND created_at < '2025-01-01') NOT VALID;
ALTER TABLE ch16_orders_2024_12 VALIDATE CONSTRAINT chk;

-- 3. ATTACH(这一步只检查约束,不扫数据,秒级完成)
ALTER TABLE ch16_orders ATTACH PARTITION ch16_orders_2024_12
    FOR VALUES FROM ('2024-12-01') TO ('2025-01-01');

诀窍:先加 NOT VALID 的 CHECK 再 VALIDATE,可以让最后的 ATTACH 跳过全表扫描。

6.4 pg_partman 扩展:自动化分区维护

pg_partman 是社区最流行的分区管理工具:

  • 自动建未来分区:每次后台任务跑都会预建 N 个未来分区;
  • 自动删过期分区:按保留策略自动 detach + drop;
  • 配合 BGW 或 pg_cron 做调度

简化用法:

sql
-- 1. 装扩展
CREATE EXTENSION pg_partman;

-- 2. 把现有分区表交给 pg_partman 管理
SELECT partman.create_parent(
    p_parent_table => 'public.ch16_orders',
    p_control      => 'created_at',
    p_type         => 'range',
    p_interval     => 'monthly',
    p_premake      => 4                -- 预建 4 个未来月分区
);

-- 3. 设保留策略:保留最近 12 个月,更老的自动 detach
UPDATE partman.part_config
SET retention = '12 months', retention_keep_table = false
WHERE parent_table = 'public.ch16_orders';

-- 4. 定时任务(用 pg_cron)
SELECT cron.schedule('partman-ch16_orders', '0 4 * * *',
                     $$ SELECT partman.run_maintenance('public.ch16_orders'); $$);

之后你就完全不用关心分区了。


7. 分区表的性能注意点

7.1 分区数不要太多

每张分区表都是一张物理表,规划器要为每个查询遍历所有分区做剪枝判断:

分区数单 SELECT 规划耗时(粗略)
12<1ms
1005ms
100050ms
10000500ms+(规划比执行还慢!)

经验值:单父表分区数 ≤ 几百,最多 1000

7.2 分区键选择决定能否裁剪

关键:最常用的查询条件 = 分区键

业务分区键
按月查订单created_at(RANGE)
按用户读自己数据user_id(HASH)
多租户 SaaStenant_id(LIST 或 HASH)
全国连锁门店region(LIST)

7.3 跨分区 JOIN 性能:partitionwise_join

如果两张分区表用相同分区键 JOIN,PG 11+ 可以「按分区对齐」做 JOIN,每个分区内单独 hash join,大幅提升性能:

sql
SET enable_partitionwise_join = on;        -- 默认 off!
SELECT * FROM ch16_orders o JOIN ch16_order_items oi
ON o.id = oi.order_id AND o.created_at = oi.created_at;

注意 enable_partitionwise_join 默认是 off,要手动开。开启后 work_mem 占用会增加(每个分区 JOIN 都要分 work_mem)。

7.4 分区聚合:partitionwise_aggregate

类似地:

sql
SET enable_partitionwise_aggregate = on;
SELECT region, sum(amount) FROM ch16_orders GROUP BY region;

每个分区先聚合再汇总,少传中间数据。


8. 真正的分库分表:Sharding 方案

PG 的内置分区只能在单实例内拆表,单机硬件上限挡在那。要跨多机,需要 sharding 方案。

8.1 Citus(最主流)

原理

        ┌─────────────────┐
        │  Coordinator    │  ← 应用连这里
        └────────┬────────┘

        ┌────────┼────────┬─────────┐
        ↓        ↓        ↓         ↓
   ┌────────┐┌────────┐┌────────┐┌────────┐
   │Worker 1││Worker 2││Worker 3││Worker 4│
   │ shard 1││ shard 2││ shard 3││ shard 4│
   │ shard 5││ shard 6││ shard 7││ shard 8│
   └────────┘└────────┘└────────┘└────────┘
  • Coordinator(协调节点):保存元数据(哪张表分了多少 shard、shard 在哪台 worker),改写并分发 SQL;
  • Worker(工作节点):每个 worker 上是普通 PG 实例,存若干 shard;
  • Shard:实际是 worker 上的物理表 orders_102008orders_102009,按 hash(distribution_column) 分散;
  • 应用感知:完全不感知,连 coordinator 像连普通 PG。

Citus 表分两类:

sql
-- 1. 分布式表(distributed table):按列 hash 分到 worker
SELECT create_distributed_table('ch16_orders', 'user_id');

-- 2. 引用表(reference table):每个 worker 全量复制(如字典表)
SELECT create_reference_table('ch16_countries');

Citus 的强项

  • 分布式 OLAP:聚合查询自动 push down 到各 worker 并行算;
  • 分布式事务(基于 2PC);
  • 跨 shard JOIN(如果 join 列与分布列对齐,下推到 worker;否则 coordinator 拉数据)。

Citus 的弱项

  • 跨 shard 的 UNIQUE 约束实现复杂;
  • 修改分布列(重新分片)很重;
  • 自增主键不能简单用 SERIAL(要用 BIGSERIAL + 自增起点错开,或用 UUID)。

8.2 其他方案

方案现状说明
pg_shard❌ 已废弃,被 Citus 合并早期 sharding 扩展
Postgres-XL🧊 维护减少老牌 share-nothing 方案,社区不活跃
Greenplum✅ 活跃基于 PG 的 MPP 数据仓库(Pivotal/VMware → Broadcom)
YugabyteDB / CockroachDB✅ 活跃完全重写的分布式 SQL,不是 PG 但兼容 wire 协议

9. 外部表 FDW:联邦查询 / 假分片

Foreign Data Wrapper(FDW):让 PG 能够像查本地表一样查其他数据源(其他 PG、MySQL、文件、Mongo、ES、Parquet ...)。

常用 FDW 列表:

扩展数据源
postgres_fdw(内置)远程 PG
mysql_fdwMySQL
oracle_fdwOracle
file_fdw(内置)本地 CSV / TSV
mongo_fdwMongoDB
parquet_fdwParquet 文件
redis_fdwRedis

9.1 postgres_fdw:跨库访问

sql
-- 1. 装扩展
CREATE EXTENSION postgres_fdw;

-- 2. 定义远程服务器
CREATE SERVER pg_remote
    FOREIGN DATA WRAPPER postgres_fdw
    OPTIONS (host '192.168.1.20', port '5432', dbname 'archive_pg');

-- 3. 凭证映射
CREATE USER MAPPING FOR postgres
    SERVER pg_remote
    OPTIONS (user 'postgres', password 'remote_pwd');

-- 4. 定义外部表
CREATE FOREIGN TABLE orders_2024 (
    id BIGINT,
    user_id BIGINT,
    amount NUMERIC,
    created_at TIMESTAMPTZ
) SERVER pg_remote OPTIONS (schema_name 'public', table_name 'ch16_orders');

-- 5. 像本地表一样用
SELECT * FROM orders_2024 WHERE user_id = 123;
-- PG 自动把 WHERE 条件下推到远程库执行

9.2 「假分片」:FDW + 分区表 = 跨库分区

把一张大分区表的不同分区放到不同 PG 实例:

sql
CREATE TABLE ch16_orders (
    id BIGINT, created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);

-- 2024 分区放远程归档库(用 FDW)
CREATE FOREIGN TABLE orders_2024 PARTITION OF ch16_orders
    FOR VALUES FROM ('2024-01-01') TO ('2025-01-01')
    SERVER pg_remote OPTIONS (table_name 'orders_2024');

-- 2025 分区在本地
CREATE TABLE orders_2025 PARTITION OF ch16_orders
    FOR VALUES FROM ('2025-01-01') TO ('2026-01-01');

查询自动按时间路由到对应的库。性能不如 Citus(需要跨网),但零成本


10. 与 MySQL 全方位对比

维度MySQLPostgreSQL
内置分区类型RANGE / LIST / HASH / KEYRANGE / LIST / HASH(PG 11+)
多级分区仅 SUBPARTITION(最多 2 层)任意层级
默认分区部分版本支持原生 DEFAULT 分区
分区裁剪支持支持,且支持执行期裁剪(PG 11+)
分区索引级联自动PG 11+ 自动
主键限制不能用任意列主键必须包含所有分区键
跨分区 JOIN 优化较弱partitionwise_join(PG 11+)
Sharding 中间件ShardingSphere / Vitess / MyCATCitus / FDW
MPP 数据仓库TiDB / OceanBase 同 wireGreenplum

11. 实战脚本(同目录 code/

  • 01_range_partition_demo.py:建表 + 插入跨月数据 + EXPLAIN 看分区裁剪;
  • 02_hash_partition_demo.py:HASH 分区均匀性验证(统计每个分区行数);
  • 03_partition_pruning.py:EXPLAIN 对比裁剪 / 不裁剪扫描的分区数;
  • 04_partman_simulation.py:Python 模拟 pg_partman 自动建/删分区;
  • 05_postgres_fdw.py:用 postgres_fdw 跨库联邦查询。

12. 小结

10 句话总结:

  1. 分区 = 单实例内拆,PG 自动路由;分库 = 多实例分散,需要中间件。
  2. PG 10+ 声明式分区有 RANGE / LIST / HASH 三种,多级可嵌套。
  3. 分区键必须包含在主键里,这是 PG 的硬约束。
  4. DEFAULT 分区做兜底,但要避免它装实际业务数据。
  5. 分区裁剪是分区性能的核心,WHERE 必须直接对分区键操作。
  6. PG 11+ 父表索引会级联到所有子分区,无需手动到每个分区建索引。
  7. DETACH CONCURRENTLY(PG 14+)实现不锁表的分区下线。
  8. pg_partman 自动维护时序分区,配合 pg_cron 做调度。
  9. 分区数控制在 几十到几百,过多反而拖慢规划器。
  10. 真正的分库分表用 Citus(推荐)/ FDW(轻量)/ Greenplum(MPP)。

🎮 配套演示

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

配套代码在 ./16_partition/code/,每个脚本都可以独立 python xxx.py 运行,先跑 init.sql 准备数据。


13. 面试高频题

Q1:PG 10+ 声明式分区与「继承表 + 触发器」相比,优势是什么?

考察点:分区机制演进、声明式分区原理。

标准答案

PG 10 之前实现分区的「老办法」是:

  1. 父表是普通表(也存数据);
  2. 子表 INHERITS (parent) 继承父表结构;
  3. 写一个 INSERT 触发器,根据值决定 redirect 到哪个子表;
  4. 给子表加 CHECK 约束,让 constraint exclusion 做裁剪。

老办法的痛点

  • 触发器开销:每次 INSERT 都跑一遍触发器函数,1000 行/秒就是 1000 次 plpgsql 调用,CPU 占用高;
  • CHECK 约束维护:子表加错约束就插不进,运维事故频发;
  • 约束剪枝粒度粗:依赖 constraint_exclusion = on,规划期才裁,且必须严格表达式匹配;
  • INSERT 路由不在内核:每行单独 trigger,无法批量优化;
  • 无法保证父表只读:父表自己可能也有数据(必须自己约束不写)。

声明式分区(PG 10+)的优势

  1. 内核原生路由INSERT INTO parent 直接由 executor 内核根据分区键算出目标子表,无触发器开销,单条 INSERT 性能提升数倍;
  2. 父表自动只读:父表是 relkind='p'(partitioned),完全不存数据,保证逻辑表的纯净
  3. 更强的剪枝:PG 11+ 支持执行期剪枝(Subplans Removed),prepared statement、子查询、运行时参数都能剪;
  4. 索引级联:PG 11+ 父表建索引自动下发所有子分区;
  5. DETACH CONCURRENTLY:PG 14+ 不锁表下线分区;
  6. Partitionwise JOIN/Aggregate:跨分区操作下推优化。

加分项:能讲出「老办法依然在用」的场景——pg_pathman 这类扩展用 hook 拦截 planner,比早期声明式分区在「分区数极多时」性能还好;不过 PG 13+ 后差距已经很小,新项目无脑用声明式即可。


Q2:分区裁剪是什么?计划期裁剪和执行期裁剪有什么区别?怎么用 EXPLAIN 验证?

考察点:分区性能核心机制。

标准答案

分区裁剪(Partition Pruning):根据 WHERE 条件自动跳过不可能命中的分区,只扫相关分区,是分区表性能的关键。

两种时机

  1. 计划期裁剪(Planning-Time Pruning):WHERE 是字面量常量,规划器编译时就知道值,直接在 Plan 里去掉不相关分区。EXPLAIN 输出里你根本看不到被剪掉的分区

  2. 执行期裁剪(Execution-Time Pruning,PG 11+):WHERE 用了参数(PREPARE / 函数参数 / 子查询结果 / now()),规划器编译时不知道值,但执行期可以裁。EXPLAIN ANALYZE 输出里有 Subplans Removed: N,表示执行期剪了 N 个分区。

EXPLAIN 验证

sql
-- 计划期裁剪
EXPLAIN SELECT * FROM ch16_orders WHERE created_at >= '2025-02-01' AND created_at < '2025-03-01';
-- 输出:Append → Seq Scan on ch16_orders_2025_02   (只见 2 月分区)

-- 执行期裁剪
PREPARE q (date) AS SELECT * FROM ch16_orders WHERE created_at >= $1;
EXPLAIN (ANALYZE) EXECUTE q ('2025-06-01');
-- 输出:Append (... Subplans Removed: 5) → Seq Scan on ch16_orders_2025_06 + 后续分区

触发条件:WHERE 里分区键必须直接出现在比较操作符的一侧,不能套表达式。

WHERE 形式是否触发裁剪
created_at >= '2025-02-01'
created_at + INTERVAL '1d' >= '...'❌(左侧加了表达式)
extract(month FROM created_at) = 2
created_at >= now() - INTERVAL '7d'✅(PG 12+ stable function)

关闭裁剪做对比

sql
SET enable_partition_pruning = off;
EXPLAIN ANALYZE SELECT count(*) FROM ch16_orders WHERE created_at >= '2025-02-01' ...
-- 看到 N 个分区都被扫

加分项:能解释为什么 partition pruning 失效要用 query rewrite:把 extract(year FROM created_at) = 2025 改写成 created_at >= '2025-01-01' AND created_at < '2026-01-01',是常见调优手法;能讲到 enable_partition_pruning vs 老的 constraint_exclusion(前者专为声明式分区,后者给老的继承式分区,PG 17 仍并存)。


Q3:什么时候用 RANGE / LIST / HASH 分区?给业务场景举例。

考察点:分区类型选型。

标准答案

RANGE 分区——按值的连续区间拆,适合:

  • 时序数据(最经典):日志、订单、监控指标,按 created_at 月/日分区;可以快速归档(DROP/DETACH 老分区);
  • 数值范围:按价格段、年龄段分;
  • 能利用「范围查询」自动裁剪WHERE ts BETWEEN ... AND ...)。

LIST 分区——按枚举值拆,适合:

  • 多租户:按 tenant_id 分(如果租户数不多且固定);
  • 按地区/渠道region IN ('CN', 'US', 'EU')
  • 按状态status IN ('active', 'archived'),用于把冷数据集中存放;
  • 已知一组离散值,且每个值数据量有差异时用得多。

HASH 分区——按哈希取模拆,适合:

  • 均匀分散写压力user_id 高基数,用 hash 让多分区均匀写入,避免最新分区写热点;
  • 分布式 sharding 的雏形:对接 Citus 时按 user_id hash 分布;
  • 没有自然区间的列:UUID、订单号 hash 等。

典型场景对照

业务分区类型分区键
电商订单(看最近一月)RANGEcreated_at(月)
微信日志(按天清理)RANGElog_date(日)
SaaS 租户表LIST 或 HASHtenant_id
IM 消息表(写均匀)HASHuser_id
全国门店报表LISTprovince_code

误用案例

  • 用 RANGE 分按 user_id:用户增长不均匀,老分区可能空、新分区爆满;不如 HASH。
  • 用 HASH 分按 created_at:完全没法利用时间范围裁剪 + 不能归档老数据;坑爹。
  • 用 LIST 分按 country_code:200 多个国家全建分区,规划器变慢;改 HASH 或合并到地理大区。

加分项:能补充「多级分区」场景——一级 RANGE 按月、二级 HASH 按 user_id,既能按时间归档又能均匀分散写压力;能讲到 PG 不支持 KEY 分区(MySQL 有),KEY 在 PG 里就是 HASH。


Q4:分区表上的主键 / 唯一约束有什么限制?为什么?

考察点:分区表约束实现原理。

标准答案

限制:分区表上的主键(PRIMARY KEY)和唯一约束(UNIQUE)必须包含所有分区键列

sql
-- ❌ 报错
CREATE TABLE ch16_orders (
    id BIGSERIAL PRIMARY KEY,         -- 主键只有 id
    created_at TIMESTAMPTZ NOT NULL
) PARTITION BY RANGE (created_at);
-- ERROR: unique constraint on partitioned table must include all partitioning columns

-- ✅ 行
CREATE TABLE ch16_orders (
    id BIGSERIAL,
    created_at TIMESTAMPTZ NOT NULL,
    PRIMARY KEY (id, created_at)      -- 主键 = (id, created_at)
) PARTITION BY RANGE (created_at);

原因

PG 的唯一索引是每个分区独立维护的(每个分区一棵 B-Tree)。要想保证「全表唯一」,必须能从「分区键」反推到「具体分区」,然后只查那个分区的索引就能确定唯一性。

如果主键不包含分区键:

  • 插入 id=100 时,PG 不知道这一行要去哪个分区;
  • 即使 id=100 已经存在于某个分区,也得查所有分区的索引才能发现冲突;
  • 这等于全表 N 次索引查找,性能崩溃。

包含分区键后

  • 插入 (id=100, created_at='2025-02-15'),PG 直接路由到 ch16_orders_2025_02
  • ch16_orders_2025_02 内查 (id=100, created_at='2025-02-15') 是否存在 → O(log n);
  • 全表唯一 = 分区键定位分区 + 分区内 (主键 - 分区键) 部分的「全局唯一性是用户责任」。

实务影响

  • 业务上 id 可能依然要全局唯一(如订单号),就要:
    • 用 UUID / Snowflake 等天然全局唯一 ID 生成方案;
    • 主键写 (id, created_at) 但应用代码保证 id 不重复;
  • 外键约束的限制类似:FK 引用列也要满足这条规则;
  • PG 11+ 才支持「分区表作为外键引用方」,PG 12+ 才支持「分区表作为外键来源方」。

加分项:能讲到 PG 14+ 引入的「全局索引(global index)」一直是社区呼声,但截至 PG 17 仍未实现;能讲为什么 Citus 解决了类似问题——它在 coordinator 上做唯一性校验,但代价是分布式事务开销;能讲序列(SEQUENCE)在分区表上的注意点——BIGSERIAL 默认共用一个序列,可能成为瓶颈,可以改成每分区独立序列 + ID 错位生成。


Q5:怎么用 PG 做「分库分表」?Citus 的工作原理?

考察点:跨实例 sharding 方案、Citus 架构。

标准答案

PG 内置分区只能在单实例内拆表,受限于单机硬件。要跨多机分散数据/算力,主流方案是 Citus。

Citus 架构

       App ──→ Coordinator ──→  Worker 1 (shard 1, 5)
                              ├→ Worker 2 (shard 2, 6)
                              ├→ Worker 3 (shard 3, 7)
                              └→ Worker 4 (shard 4, 8)
  • Coordinator(协调节点):是个 PG 实例,安装了 citus 扩展。它保存元数据(哪张表分了多少 shard、shard 在哪台 worker),接收应用 SQL 后改写并分发到 worker;
  • Worker(工作节点):每个 worker 是普通 PG 实例(也装 citus),存放若干 shard 物理表,名字形如 orders_102008(102008 是 shard ID);
  • Distribution Column:分布键,hash(distribution_column) % shard_count 决定行落到哪个 shard。

两种表类型

sql
-- 1. 分布式表:按 user_id hash 分散到所有 worker
SELECT create_distributed_table('ch16_orders', 'user_id');
-- 默认 32 shard

-- 2. 引用表:每个 worker 全量复制(适合小字典表)
SELECT create_reference_table('ch16_countries');

典型查询

  • SELECT * FROM ch16_orders WHERE user_id = 100:coordinator 根据 hash 计算只发给一个 worker,O(1)
  • SELECT count(*) FROM ch16_orders:coordinator 把 SQL 改写下推到所有 worker 并行算,再合并结果(map-reduce 风格);
  • 分布式 JOIN:如果 JOIN 列与分布列对齐(如 ch16_ordersch16_order_items 都按 user_id 分),下推到 worker 本地 JOIN;否则需要 coordinator 拉一张表过来 JOIN(性能差);
  • 分布式事务:基于 2PC(PostgreSQL 的 prepared transaction),跨 shard 一致提交。

优缺点

优点

  • 应用感知度低,连 coordinator 像普通 PG;
  • 分布式 OLAP 极强,聚合并行加速明显;
  • 是 PG 官方扩展,跟新版 PG 兼容(被微软收购后仍开源);
  • 支持引用表,适合事实表 + 维度表场景。

缺点

  • 跨 shard 唯一约束实现复杂(用 reference table 或应用层保证);
  • 修改分布列需重分片(很重);
  • 序列(SERIAL)有限制,推荐 UUID / Snowflake;
  • coordinator 是单点(虽然可以做 HA)。

其他 sharding 方案

  • FDW + 分区:把不同分区放到不同 PG 实例(postgres_fdw),轻量但性能差;
  • YugabyteDB / CockroachDB:完全重写的分布式 SQL,PG wire 兼容,但不是真 PG;
  • Greenplum:基于 PG 改的 MPP 数据仓库,OLAP 极强,OLTP 不适合;
  • 应用层 sharding:业务自己路由到不同 PG(最灵活但最重)。

加分项:能讲 Citus 的「co-location group」概念——多张表用相同分布列时归为一组,跨表 JOIN 全部 push down;能讲 Citus 11+ 引入「无 coordinator 模式」(每个节点都能接连接);能对比 Vitess(MySQL)和 Citus 的设计差异(Vitess 用 vtgate 做 SQL 路由,Citus 在 PG 内核做)。


Q6:什么是 FDW?postgres_fdw 怎么实现「跨库 JOIN」?性能怎样?

考察点:联邦查询、FDW 下推机制。

标准答案

Foreign Data Wrapper(FDW):PG 的扩展机制,把外部数据源「包装」成一张 PG 内的「外部表」(FOREIGN TABLE),让 SQL 能像查本地表一样查它们。

支持的数据源(社区维护):

  • postgres_fdw(内置):另一个 PG 实例;
  • mysql_fdworacle_fdwsqlserver_fdw:其他 RDBMS;
  • mongo_fdwredis_fdwelasticsearch_fdw:NoSQL;
  • file_fdw(内置)、csv_fdwparquet_fdw:文件;
  • hdfs_fdw:大数据生态。

配置 4 步走

sql
CREATE EXTENSION postgres_fdw;
CREATE SERVER s FOREIGN DATA WRAPPER postgres_fdw OPTIONS (host '...', dbname '...');
CREATE USER MAPPING FOR postgres SERVER s OPTIONS (user '...', password '...');
CREATE FOREIGN TABLE remote_orders (...) SERVER s OPTIONS (table_name 'ch16_orders');

-- 用法
SELECT * FROM remote_orders WHERE user_id = 100;
SELECT u.name, o.amount FROM users u JOIN remote_orders o ON u.id = o.user_id;

跨库 JOIN 的实现

PG 的 FDW 框架会做下推优化(pushdown)

  1. WHERE 下推WHERE user_id = 100 会原样发到远程库执行,只返回符合条件的行;
  2. 投影下推:只 SELECT 需要的列,不传无关列;
  3. JOIN 下推(PG 9.6+):如果 JOIN 两边都在同一 foreign server 上,PG 可以把整个 JOIN 下推到远程执行,本地只接收结果;
  4. 聚合下推(PG 10+):SELECT count(*) FROM remote_orders 直接在远程 count;
  5. 排序 / LIMIT 下推:能下推就下推,减少传输。

性能特点

  • 好的情况:WHERE 选择性好 + 都能下推 → 性能接近本地查询,只是多了网络一跳;
  • 差的情况:JOIN 跨多个 server(无法下推),PG 只能从远程拉两张全表回本地 hash join,百万行级别就崩盘
  • 差的情况:远程表统计信息没拉过来(要手动 ANALYZE remote_table),规划器估错行数。

用 EXPLAIN VERBOSE 看下推

sql
EXPLAIN VERBOSE SELECT * FROM remote_orders WHERE user_id = 100;
-- 输出会有 "Remote SQL: SELECT id, user_id, amount FROM public.ch16_orders WHERE user_id = 100"

「假分片」用法:把分区表的不同分区放到不同 PG 实例:

sql
CREATE FOREIGN TABLE orders_2024 PARTITION OF ch16_orders
    FOR VALUES FROM ('2024-01-01') TO ('2025-01-01')
    SERVER pg_archive;

冷数据放归档库,热数据本地,应用零改造。

加分项:能讲出 use_remote_estimate = true 的作用——让 FDW 用远程库的统计信息估行数(默认 false 用本地默认值不准);能讲到 fetch_size 调优——大查询调大 fetch_size 减少 round trip;能对比 Citus 和 FDW:Citus 是为分布式重写过的执行引擎,下推 / 并行做得彻底;FDW 是「能下推就下推」的尽力而为。


Q7:用 pg_partman 做时序分区表自动维护,关键参数有哪些?怎么搞挂的?

考察点:分区表运维细节。

标准答案

pg_partman 是 PG 社区最流行的分区自动维护扩展,主要做三件事:

  1. 预建未来分区:避免「写入时分区不存在」;
  2. 删除过期分区:按保留策略自动 detach + drop;
  3. 管理 default 分区:把误入 default 的数据迁回正确分区。

核心配置

sql
SELECT partman.create_parent(
    p_parent_table => 'public.ch16_orders',
    p_control      => 'created_at',     -- 分区键列
    p_type         => 'range',
    p_interval     => 'monthly',        -- 分区粒度:daily/weekly/monthly/yearly/数字
    p_premake      => 4,                -- 预建未来 4 个分区
    p_start_partition => '2025-01-01'   -- 起始分区
);

-- 配置保留策略(保留最近 24 个月)
UPDATE partman.part_config SET
    retention = '24 months',
    retention_keep_table = false,        -- false = drop 表;true = 只 detach
    retention_keep_index = false,
    infinite_time_partitions = true      -- 即使没数据也持续创建未来分区
WHERE parent_table = 'public.ch16_orders';

-- 调度(用 pg_cron 每天凌晨 4 点维护)
SELECT cron.schedule('partman-ch16_orders', '0 4 * * *',
                     $$ SELECT partman.run_maintenance('public.ch16_orders'); $$);

关键参数

参数作用推荐值
p_premake预建多少个未来分区≥ 4(防止维护 cron 挂了仍有缓冲)
p_interval分区粒度看业务数据量:100w/月用 monthly,1 亿/天用 daily
retention数据保留多久业务定(如日志 30 天,订单 24 个月)
retention_keep_table删表 or detachdetach 后归档到冷存储
infinite_time_partitions持续创建未来分区true

经典翻车场景

  1. p_premake 太小,维护 cron 挂了:cron 失败几天,写入时间超过预建范围 → 落入 default 分区 → default 越来越大;
  2. 没装 pg_cron 也没设置 BGW:以为「装了 partman 就自动跑了」,其实需要手动调 run_maintenance()
  3. default 分区有数据后建新分区超慢:partman 添加新分区时要扫 default 找冲突行,几十 GB default 加一次分区要锁几十分钟;
  4. retention_keep_table = false 误删数据:保留策略写错(比如把 30 days 写成 30 months),第二天醒来发现两年的分区被 drop 了;强烈建议先 retention_keep_table = true 跑一段时间,确认没事再改 false;
  5. 大事务 + 分区 detach 死锁:detach 拿 ACCESS EXCLUSIVE 锁,正在跑的长 SELECT 会挡住;PG 14+ 可以用 partman.partition_data_proc + DETACH CONCURRENTLY
  6. 分区数失控p_premake = 100 + infinite_time_partitions = true → 分区数飙到上千,规划器变慢。

监控

sql
-- 查 default 分区的行数
SELECT count(*) FROM ch16_orders_default;

-- 查最近一次维护时间
SELECT * FROM partman.part_config WHERE parent_table = 'public.ch16_orders';

加分项:能讲到 partman 也支持 epoch 时间戳分区、支持子分区(多级 partition);能比较 partman 和「自己写 cron + plpgsql」的取舍——partman 帮你处理了 default 分区数据迁移、ATTACH 优化等坑;能讲 timescaleDB 的 hypertable 是另一种思路——它在分区基础上再封装一层 chunk 自动管理,用户体验更好但失去通用性。


14. 实操彩蛋:60 秒看懂分区裁剪

sql
-- 1. 建一张 12 月分区表
CREATE TABLE quick_demo (
    id BIGSERIAL,
    ts TIMESTAMPTZ NOT NULL,
    val INT,
    PRIMARY KEY (id, ts)
) PARTITION BY RANGE (ts);

DO $$
DECLARE m INT;
BEGIN
    FOR m IN 1..12 LOOP
        EXECUTE format(
            'CREATE TABLE quick_demo_%s PARTITION OF quick_demo FOR VALUES FROM (%L) TO (%L)',
            lpad(m::text, 2, '0'),
            format('2025-%s-01', lpad(m::text, 2, '0'))::date,
            (format('2025-%s-01', lpad(m::text, 2, '0'))::date + INTERVAL '1 month')::date
        );
    END LOOP;
END$$;

-- 2. 灌点数据
INSERT INTO quick_demo (ts, val)
SELECT '2025-01-01'::timestamptz + (random() * INTERVAL '364 day'),
       (random() * 1000)::INT
FROM generate_series(1, 100000);

-- 3. 见证奇迹
EXPLAIN ANALYZE SELECT count(*) FROM quick_demo
WHERE ts >= '2025-06-01' AND ts < '2025-07-01';
-- 只扫 quick_demo_06,用时 ~5ms

EXPLAIN ANALYZE SELECT count(*) FROM quick_demo
WHERE extract(month FROM ts) = 6;
-- 扫所有 12 个分区,用时 ~80ms(裁剪失败!)

亲手敲一遍,分区裁剪的威力立刻 GET。完整版本见 code/03_partition_pruning.py


🔗 延伸阅读


🎬 可视化演示

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

💻 示例代码

python
#!/usr/bin/env python3
"""
01_range_partition_demo.py
==========================
演示 RANGE 分区表:
  1. 从零建表 + 6 个月分区 + default 分区
  2. 灌入跨月数据
  3. 验证数据自动路由到对应分区
  4. EXPLAIN 看分区裁剪
  5. 演示「插入 default 分区」和「添加新分区」

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

import psycopg

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


def hr(s):
    print("\n" + "=" * 68)
    print(f"  {s}")
    print("=" * 68)


def run(cur, sql, fetch=True):
    cur.execute(sql)
    if fetch and cur.description:
        rows = cur.fetchall()
        for r in rows:
            print(" ", r)
        return rows


def main():
    with psycopg.connect(DSN, autocommit=True) as conn, conn.cursor() as cur:
        hr("Step 1: 清理 + 建分区父表")
        run(cur, "DROP TABLE IF EXISTS ch16_demo_orders CASCADE;", fetch=False)
        run(cur, """
            CREATE TABLE ch16_demo_orders (
                id BIGSERIAL,
                user_id BIGINT NOT NULL,
                amount NUMERIC(12,2),
                created_at TIMESTAMPTZ NOT NULL,
                PRIMARY KEY (id, created_at)
            ) PARTITION BY RANGE (created_at);
        """, fetch=False)
        print("  ✓ 父表 ch16_demo_orders 已建(PARTITION BY RANGE created_at)")

        hr("Step 2: 建 3 个月分区 + default 分区")
        for ym, frm, to in [("2025_01", "2025-01-01", "2025-02-01"),
                            ("2025_02", "2025-02-01", "2025-03-01"),
                            ("2025_03", "2025-03-01", "2025-04-01")]:
            run(cur, f"""
                CREATE TABLE ch16_demo_orders_{ym} PARTITION OF ch16_demo_orders
                    FOR VALUES FROM ('{frm}') TO ('{to}');
            """, fetch=False)
        run(cur, "CREATE TABLE ch16_demo_orders_default PARTITION OF ch16_demo_orders DEFAULT;",
            fetch=False)
        print("  ✓ 3 个月分区 + default 分区建好")

        hr("Step 3: 父表索引(PG 11+ 自动级联到所有子分区)")
        run(cur, "CREATE INDEX ON ch16_demo_orders (user_id);", fetch=False)
        print("  ✓ 在父表建索引,自动下发到所有子分区")
        run(cur, """
            SELECT indexrelid::regclass AS index_name
            FROM pg_index
            WHERE indrelid IN (
                SELECT inhrelid FROM pg_inherits WHERE inhparent='ch16_demo_orders'::regclass
            )
            ORDER BY index_name;
        """)

        hr("Step 4: 插入跨月数据(PG 自动路由到对应分区)")
        run(cur, """
            INSERT INTO ch16_demo_orders (user_id, amount, created_at) VALUES
                (1, 100.00, '2025-01-15 10:00:00'),
                (2, 200.00, '2025-01-20 11:00:00'),
                (3, 300.00, '2025-02-05 12:00:00'),
                (4, 400.00, '2025-02-25 13:00:00'),
                (5, 500.00, '2025-03-10 14:00:00'),
                (6, 600.00, '2025-09-01 15:00:00')   -- 落入 default!
            RETURNING id, created_at;
        """)

        hr("Step 5: 各分区行数")
        run(cur, """
            SELECT
                relname AS partition,
                pg_stat_get_live_tuples(c.oid)::INT AS rows
            FROM pg_class c
            JOIN pg_inherits i ON i.inhrelid = c.oid
            WHERE i.inhparent = 'ch16_demo_orders'::regclass
            ORDER BY relname;
        """)
        # 强制 ANALYZE 更新统计信息
        run(cur, "ANALYZE ch16_demo_orders;", fetch=False)

        hr("Step 6: EXPLAIN 演示分区裁剪")
        print("\n--- 6.1 WHERE 用字面量常量(计划期裁剪) ---")
        run(cur, """
            EXPLAIN SELECT * FROM ch16_demo_orders
            WHERE created_at >= '2025-02-01' AND created_at < '2025-03-01';
        """)

        print("\n--- 6.2 WHERE 不能用:表达式套在分区键上 ---")
        run(cur, """
            EXPLAIN SELECT * FROM ch16_demo_orders
            WHERE EXTRACT(MONTH FROM created_at) = 2;
        """)

        print("\n--- 6.3 关闭分区裁剪做对比 ---")
        run(cur, "SET enable_partition_pruning = off;", fetch=False)
        run(cur, """
            EXPLAIN SELECT * FROM ch16_demo_orders
            WHERE created_at >= '2025-02-01' AND created_at < '2025-03-01';
        """)
        run(cur, "RESET enable_partition_pruning;", fetch=False)

        hr("Step 7: 添加 4 月分区,把 default 里的数据搬走")
        # 演示常见运维操作:default 里有了未匹配数据,需要先建对应的分区
        run(cur, """
            CREATE TABLE ch16_demo_orders_2025_04 PARTITION OF ch16_demo_orders
                FOR VALUES FROM ('2025-04-01') TO ('2025-05-01');
        """, fetch=False)
        print("  ✓ 4 月分区已建,但 default 里 2025-09-01 的行还需要单独处理")
        # 实际上 default 里的 09 月行不在 04 月范围,只能新建 09 月分区或保留
        run(cur, """
            SELECT created_at FROM ch16_demo_orders_default;
        """)

        hr("Step 8: DETACH CONCURRENTLY(PG 14+ 不锁表下线分区)")
        # 注意:DETACH CONCURRENTLY 不能在事务里执行
        try:
            run(cur, "ALTER TABLE ch16_demo_orders DETACH PARTITION ch16_demo_orders_2025_01 CONCURRENTLY;",
                fetch=False)
            print("  ✓ 1 月分区已下线(变成普通独立表 ch16_demo_orders_2025_01)")
        except psycopg.Error as e:
            print(f"  ⚠️  DETACH 失败(可能 PG 版本 < 14):{e}")

        hr("最终各分区状态")
        run(cur, """
            SELECT
                relname,
                pg_stat_get_live_tuples(c.oid)::INT AS rows,
                pg_size_pretty(pg_relation_size(c.oid)) AS size
            FROM pg_class c
            JOIN pg_inherits i ON i.inhrelid = c.oid
            WHERE i.inhparent = 'ch16_demo_orders'::regclass
            ORDER BY relname;
        """)

        print("\n[DONE] 完整演示结束。可以打开 demo.html 看可视化效果。")
        return 0


if __name__ == "__main__":
    sys.exit(main() or 0)
python
#!/usr/bin/env python3
"""
02_hash_partition_demo.py
=========================
HASH 分区均匀性验证:
  1. 建 8 个 HASH 分区,按 user_id 分
  2. 灌入 100,000 个不同 user_id 的事件
  3. 统计每个分区行数,验证哈希分布是否均匀
  4. 演示「WHERE user_id = X」能裁剪到 1 个分区
  5. 演示「WHERE user_id BETWEEN A AND B」无法裁剪(哈希后不连续)

依赖:psycopg[binary]>=3.1
"""
import os
import sys
import statistics

import psycopg

DSN = os.environ.get(
    "PG_DSN",
    "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres",
)
N_PARTITIONS = 8
N_ROWS = 100_000


def hr(s):
    print("\n" + "=" * 68)
    print(f"  {s}")
    print("=" * 68)


def main():
    with psycopg.connect(DSN, autocommit=True) as conn, conn.cursor() as cur:
        hr("Step 1: 建 HASH 分区表(8 个分区)")
        cur.execute("DROP TABLE IF EXISTS ch16_hash_demo CASCADE;")
        cur.execute("""
            CREATE TABLE ch16_hash_demo (
                id BIGSERIAL,
                user_id BIGINT NOT NULL,
                payload TEXT,
                PRIMARY KEY (id, user_id)
            ) PARTITION BY HASH (user_id);
        """)
        for i in range(N_PARTITIONS):
            cur.execute(f"""
                CREATE TABLE ch16_hash_demo_p{i} PARTITION OF ch16_hash_demo
                    FOR VALUES WITH (MODULUS {N_PARTITIONS}, REMAINDER {i});
            """)
        print(f"  ✓ 已建 {N_PARTITIONS} 个 HASH 分区 ch16_hash_demo_p0 ~ p{N_PARTITIONS-1}")

        hr(f"Step 2: 灌入 {N_ROWS:,} 行(user_id 从 1 到 {N_ROWS})")
        # 用单条多行 INSERT 提升插入速度
        cur.execute(f"""
            INSERT INTO ch16_hash_demo (user_id, payload)
            SELECT g, 'event-' || g
            FROM generate_series(1, {N_ROWS}) g;
        """)
        cur.execute("ANALYZE ch16_hash_demo;")
        print(f"  ✓ 灌入完成")

        hr("Step 3: 统计各分区行数")
        cur.execute(f"""
            SELECT relname, pg_stat_get_live_tuples(c.oid)::INT AS rows,
                   pg_size_pretty(pg_relation_size(c.oid)) AS size
            FROM pg_class c
            JOIN pg_inherits i ON i.inhrelid = c.oid
            WHERE i.inhparent = 'ch16_hash_demo'::regclass
            ORDER BY relname;
        """)
        rows = cur.fetchall()
        counts = [r[1] for r in rows]
        ideal = N_ROWS / N_PARTITIONS
        print(f"  分区名             行数        大小        与理想值偏差")
        print(f"  --------------------------------------------------------")
        for name, c, sz in rows:
            diff_pct = (c - ideal) / ideal * 100
            bar = "█" * int(c / ideal * 20)
            print(f"  {name:<18} {c:<10,} {sz:<10}  {diff_pct:+.2f}%  {bar}")
        print(f"  --------------------------------------------------------")
        print(f"  理想值: {ideal:,.0f} 行/分区")
        print(f"  实际标准差: {statistics.stdev(counts):.2f}")
        print(f"  最大偏差: {(max(counts)-min(counts))/ideal*100:.2f}%")
        print(f"  → 哈希分布均匀(PG 内置哈希函数质量很好)")

        hr("Step 4: WHERE user_id = X 触发分区裁剪(点查命中 1 个分区)")
        cur.execute("EXPLAIN SELECT * FROM ch16_hash_demo WHERE user_id = 12345;")
        for r in cur.fetchall():
            print(" ", r[0])

        hr("Step 5: WHERE user_id BETWEEN ... 不能裁剪(哈希后不连续)")
        cur.execute("EXPLAIN SELECT * FROM ch16_hash_demo WHERE user_id BETWEEN 100 AND 200;")
        for r in cur.fetchall():
            print(" ", r[0])
        print("  ↑ 注意所有分区都被扫到(HASH 分区的硬伤)")

        hr("Step 6: WHERE user_id IN (a, b, c) 部分裁剪(PG 11+)")
        cur.execute("EXPLAIN SELECT * FROM ch16_hash_demo WHERE user_id IN (1, 100, 1000);")
        for r in cur.fetchall():
            print(" ", r[0])
        print("  ↑ 各 user_id 算 hash 后落到 ≤3 个分区,PG 11+ 能裁剪")

        hr("结论")
        print("  ✓ HASH 分区适合:均匀分散写入热点 / 点查(=)")
        print("  ✗ HASH 分区不适合:范围查询、按时间归档")
        print("  → 时序数据用 RANGE,高基数随机访问用 HASH")

        return 0


if __name__ == "__main__":
    sys.exit(main() or 0)
python
#!/usr/bin/env python3
"""
03_partition_pruning.py
=======================
对比「裁剪 vs 不裁剪」的性能差异,并展示什么样的 WHERE 条件能 / 不能触发裁剪。

要求:先跑过 init.sql(已有 ch16_orders 分区表 + 6 万行数据)。
"""
import os
import sys
import time

import psycopg

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


def hr(s):
    print("\n" + "=" * 70)
    print(f"  {s}")
    print("=" * 70)


def explain(cur, sql, label):
    print(f"\n[{label}]")
    print(f"SQL: {sql}")
    cur.execute(f"EXPLAIN (ANALYZE, BUFFERS) {sql}")
    rows = [r[0] for r in cur.fetchall()]
    # 统计扫描的分区数
    scanned = sum(1 for line in rows if "Scan on ch16_orders_2025" in line)
    pruned = next((line for line in rows if "Subplans Removed" in line), None)
    timing = next((line for line in rows if "Execution Time" in line), "")
    print("  --- EXPLAIN ---")
    for line in rows:
        print("  ", line)
    print(f"  >>> 扫描了 {scanned} 个分区"
          + (f"   {pruned.strip()}" if pruned else "")
          + f"   {timing.strip()}")


def main():
    with psycopg.connect(DSN, autocommit=True) as conn, conn.cursor() as cur:
        # 检查表是否存在
        cur.execute("""
            SELECT count(*) FROM pg_class WHERE relname = 'ch16_orders' AND relkind = 'p';
        """)
        if cur.fetchone()[0] == 0:
            print("[FATAL] 表 ch16_orders 不存在或不是分区表,请先跑 init.sql")
            return 1

        cur.execute("ANALYZE ch16_orders;")

        hr("场景 A:完美裁剪(WHERE 直接对分区键比较)")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE created_at >= '2025-03-01' AND created_at < '2025-04-01'",
                "A.1 范围常量")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE created_at = '2025-04-15 12:00:00'",
                "A.2 点查常量")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE created_at IN ('2025-02-15','2025-04-15','2025-05-15')",
                "A.3 IN 多值")

        hr("场景 B:失败裁剪(WHERE 套了表达式)")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE EXTRACT(MONTH FROM created_at) = 3",
                "B.1 函数包裹分区键")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE created_at + INTERVAL '1 day' >= '2025-03-01'",
                "B.2 加表达式")
        explain(cur,
                "SELECT count(*) FROM ch16_orders WHERE created_at::date = '2025-03-15'",
                "B.3 cast 改变类型")

        hr("场景 C:执行期裁剪(PREPARE 参数化查询)")
        cur.execute("DEALLOCATE ALL;")
        cur.execute("PREPARE q (timestamptz, timestamptz) AS "
                    "SELECT count(*) FROM ch16_orders WHERE created_at >= $1 AND created_at < $2;")
        # 跑几次让规划器生成 generic plan(PG 用启发式 5 次自定义后切 generic)
        for _ in range(7):
            cur.execute("EXECUTE q ('2025-04-01'::timestamptz, '2025-05-01'::timestamptz);")

        explain(cur,
                "EXECUTE q ('2025-04-01'::timestamptz, '2025-05-01'::timestamptz)",
                "C.1 PREPARE + EXECUTE(看 Subplans Removed: N)")
        cur.execute("DEALLOCATE q;")

        hr("场景 D:开 vs 关 enable_partition_pruning 性能对比")
        sql = ("SELECT count(*) FROM ch16_orders "
               "WHERE created_at >= '2025-03-01' AND created_at < '2025-04-01'")

        for setting in ("on", "off"):
            cur.execute(f"SET enable_partition_pruning = {setting};")
            # 跑 5 次取最小耗时
            times = []
            for _ in range(5):
                t0 = time.perf_counter()
                cur.execute(sql)
                cur.fetchall()
                times.append((time.perf_counter() - t0) * 1000)
            print(f"  enable_partition_pruning={setting:>3}  最快 {min(times):.2f}ms"
                  f"  平均 {sum(times)/len(times):.2f}ms")
        cur.execute("RESET enable_partition_pruning;")

        hr("场景 E:partitionwise_join 跨分区表 JOIN 优化(PG 11+)")
        # 简单演示:构造另一张同分区键的表然后 JOIN
        cur.execute("DROP TABLE IF EXISTS ch16_payments CASCADE;")
        cur.execute("""
            CREATE TABLE ch16_payments (
                id BIGSERIAL,
                order_id BIGINT,
                paid_at TIMESTAMPTZ NOT NULL,
                PRIMARY KEY (id, paid_at)
            ) PARTITION BY RANGE (paid_at);
        """)
        for ym, frm, to in [("2025_03", "2025-03-01", "2025-04-01"),
                            ("2025_04", "2025-04-01", "2025-05-01")]:
            cur.execute(f"CREATE TABLE ch16_payments_{ym} PARTITION OF ch16_payments "
                        f"FOR VALUES FROM ('{frm}') TO ('{to}');")
        cur.execute("""
            INSERT INTO ch16_payments (order_id, paid_at)
            SELECT id, created_at FROM ch16_orders
            WHERE created_at >= '2025-03-01' AND created_at < '2025-05-01'
            LIMIT 5000;
        """)
        cur.execute("ANALYZE ch16_payments;")

        sql_join = ("SELECT count(*) FROM ch16_orders o JOIN ch16_payments p "
                    "ON o.id = p.order_id AND o.created_at = p.paid_at "
                    "WHERE o.created_at >= '2025-03-01' AND o.created_at < '2025-05-01'")

        for setting in ("off", "on"):
            cur.execute(f"SET enable_partitionwise_join = {setting};")
            cur.execute(f"EXPLAIN (ANALYZE) {sql_join}")
            timing = next((r[0] for r in cur.fetchall() if "Execution Time" in r[0]), "")
            print(f"  enable_partitionwise_join={setting}  {timing.strip()}")
        cur.execute("RESET enable_partitionwise_join;")
        cur.execute("DROP TABLE ch16_payments CASCADE;")

        hr("总结")
        print("  ✓ WHERE 直接对分区键比较 → 完美裁剪")
        print("  ✓ PG 11+ PREPARE 后参数查询 → 执行期裁剪(Subplans Removed)")
        print("  ✗ 函数 / cast / 表达式包裹分区键 → 失败")
        print("  💡 partitionwise_join 默认 OFF,做大宽表 JOIN 记得开")

    return 0


if __name__ == "__main__":
    sys.exit(main() or 0)
python
#!/usr/bin/env python3
"""
04_partman_simulation.py
========================
用纯 Python + SQL 模拟 pg_partman 的核心功能:
  - 自动预建未来 N 个月分区
  - 自动 detach + drop 超过保留期的老分区
  - 维护一张配置表记录每张分区表的策略

适合无法装 pg_partman 扩展(如云托管 PG)的场景,思路完全一致。

用法:
    python3 04_partman_simulation.py setup       # 初始化配置 + 创建 demo 分区表
    python3 04_partman_simulation.py maintain    # 跑一次维护(可加到 cron)
    python3 04_partman_simulation.py status      # 查看当前分区
    python3 04_partman_simulation.py teardown    # 清理
"""
import os
import sys
from datetime import date, timedelta

import psycopg

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

CONFIG_TABLE = "ch16_my_partman_config"
PARENT_TABLE = "ch16_ts_logs"
PREMAKE = 4              # 预建未来 4 个月
RETENTION_MONTHS = 12    # 保留 12 个月
INTERVAL = "month"


def add_months(d: date, n: int) -> date:
    """日期加 n 个月(n 可为负数),返回月初"""
    y = d.year
    m = d.month + n
    while m > 12:
        m -= 12; y += 1
    while m < 1:
        m += 12; y -= 1
    return date(y, m, 1)


def setup(cur):
    """初始化:建配置表 + 父表 + 当月分区"""
    print("[setup] 创建配置表...")
    cur.execute(f"""
        CREATE TABLE IF NOT EXISTS {CONFIG_TABLE} (
            parent_table TEXT PRIMARY KEY,
            partition_col TEXT NOT NULL,
            interval_kind TEXT NOT NULL,
            premake INT NOT NULL DEFAULT 4,
            retention_months INT NOT NULL DEFAULT 12,
            last_maintenance TIMESTAMPTZ
        );
    """)

    print(f"[setup] 创建父表 {PARENT_TABLE}(按月 RANGE 分区)...")
    cur.execute(f"DROP TABLE IF EXISTS {PARENT_TABLE} CASCADE;")
    cur.execute(f"""
        CREATE TABLE {PARENT_TABLE} (
            id BIGSERIAL,
            ts TIMESTAMPTZ NOT NULL,
            level TEXT,
            msg TEXT,
            PRIMARY KEY (id, ts)
        ) PARTITION BY RANGE (ts);
    """)

    # 建当月分区
    today = date.today()
    cur_start = today.replace(day=1)
    cur_end = add_months(cur_start, 1)
    pname = f"{PARENT_TABLE}_{cur_start.strftime('%Y_%m')}"
    cur.execute(f"""
        CREATE TABLE {pname} PARTITION OF {PARENT_TABLE}
            FOR VALUES FROM ('{cur_start}') TO ('{cur_end}');
    """)
    print(f"  ✓ 已建当月分区 {pname}")

    # 注册到配置表
    cur.execute(f"""
        INSERT INTO {CONFIG_TABLE}
            (parent_table, partition_col, interval_kind, premake, retention_months)
        VALUES ('{PARENT_TABLE}', 'ts', '{INTERVAL}', {PREMAKE}, {RETENTION_MONTHS})
        ON CONFLICT (parent_table) DO UPDATE
            SET premake = EXCLUDED.premake, retention_months = EXCLUDED.retention_months;
    """)
    print("  ✓ 配置已注册")


def list_partitions(cur, parent: str):
    cur.execute(f"""
        SELECT c.relname,
               pg_get_expr(c.relpartbound, c.oid) AS bound
        FROM pg_class c
        JOIN pg_inherits i ON i.inhrelid = c.oid
        WHERE i.inhparent = %s::regclass
        ORDER BY c.relname;
    """, (parent,))
    return cur.fetchall()


def maintain(cur):
    """逐一处理配置表中所有分区表"""
    cur.execute(f"SELECT parent_table, premake, retention_months FROM {CONFIG_TABLE};")
    for parent, premake, retention in cur.fetchall():
        print(f"\n[maintain] 处理 {parent}(premake={premake}, retention={retention}个月)")
        existing = {row[0] for row in list_partitions(cur, parent)}
        print(f"  当前分区数: {len(existing)}")

        # ========= 1. 预建未来分区 =========
        today = date.today().replace(day=1)
        for i in range(premake + 1):     # 包含当月,所以 +1
            start = add_months(today, i)
            end = add_months(start, 1)
            name = f"{parent}_{start.strftime('%Y_%m')}"
            if name in existing:
                continue
            print(f"  + 创建未来分区 {name} [{start} ~ {end})")
            cur.execute(f"""
                CREATE TABLE {name} PARTITION OF {parent}
                    FOR VALUES FROM ('{start}') TO ('{end}');
            """)

        # ========= 2. 清理过期分区 =========
        cutoff = add_months(today, -retention)
        print(f"  保留截止日期: {cutoff}(更早的将被 detach + drop)")
        for name, bound in list_partitions(cur, parent):
            # 解析 bound 字符串:FOR VALUES FROM ('2024-01-01 ...') TO ('2024-02-01 ...')
            import re
            m = re.search(r"FROM \('([\d-]+)", bound or "")
            if not m:
                continue
            try:
                p_start = date.fromisoformat(m.group(1))
            except ValueError:
                continue
            if p_start < cutoff:
                print(f"  - 删除过期分区 {name}(start={p_start})")
                # PG 14+ 可以 CONCURRENTLY 不锁表
                try:
                    cur.execute(f"ALTER TABLE {parent} DETACH PARTITION {name} CONCURRENTLY;")
                except psycopg.Error:
                    cur.execute(f"ALTER TABLE {parent} DETACH PARTITION {name};")
                cur.execute(f"DROP TABLE {name};")

        cur.execute(f"UPDATE {CONFIG_TABLE} SET last_maintenance = now() WHERE parent_table = '{parent}';")
        print(f"  ✓ {parent} 维护完成")


def status(cur):
    cur.execute(f"""
        SELECT parent_table, premake, retention_months, last_maintenance
        FROM {CONFIG_TABLE};
    """)
    print("\n=== 配置 ===")
    for r in cur.fetchall():
        print(f"  {r}")
    print(f"\n=== {PARENT_TABLE} 分区列表 ===")
    for name, bound in list_partitions(cur, PARENT_TABLE):
        print(f"  {name}  {bound}")


def teardown(cur):
    cur.execute(f"DROP TABLE IF EXISTS {PARENT_TABLE} CASCADE;")
    cur.execute(f"DROP TABLE IF EXISTS {CONFIG_TABLE};")
    print("✓ 已清理")


def main():
    cmd = sys.argv[1] if len(sys.argv) > 1 else "status"
    fns = {"setup": setup, "maintain": maintain, "status": status, "teardown": teardown}
    if cmd not in fns:
        print(f"Usage: {sys.argv[0]} {{{'|'.join(fns.keys())}}}", file=sys.stderr)
        return 1
    with psycopg.connect(DSN, autocommit=True) as conn, conn.cursor() as cur:
        fns[cmd](cur)
    return 0


if __name__ == "__main__":
    sys.exit(main() or 0)
python
#!/usr/bin/env python3
"""
05_postgres_fdw.py
==================
演示 postgres_fdw(PG 内置扩展)的「跨库联邦查询」能力。

场景:本地库 (learn_pg) 通过 FDW 访问远端库的表,模拟「假分片」。
为了演示方便,我们用「同一个 PG 实例的另一个数据库」当作远端:
   - learn_pg            → 本地库(建外部表)
   - learn_pg_remote     → 远程库(实际数据放这里)

用法:
    python3 05_postgres_fdw.py setup    # 创建远程库 + 远程数据 + FDW 配置
    python3 05_postgres_fdw.py query    # 跨库查询 + 看 EXPLAIN 下推
    python3 05_postgres_fdw.py join     # 跨库 JOIN
    python3 05_postgres_fdw.py teardown # 清理
"""
import os
import sys

import psycopg
from psycopg import sql

LOCAL_DSN = os.environ.get(
    "PG_LOCAL",
    "host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres",
)
ADMIN_DSN = os.environ.get(
    "PG_ADMIN",
    "host=127.0.0.1 port=5432 dbname=postgres user=postgres password=postgres",
)
REMOTE_DBNAME = "learn_pg_remote"


def hr(s):
    print("\n" + "=" * 68)
    print(f"  {s}")
    print("=" * 68)


def setup():
    hr("Step 1: 创建「远端」数据库")
    with psycopg.connect(ADMIN_DSN, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute(f"DROP DATABASE IF EXISTS {REMOTE_DBNAME};")
        cur.execute(f"CREATE DATABASE {REMOTE_DBNAME};")
        print(f"  ✓ 已创建数据库 {REMOTE_DBNAME}")

    hr("Step 2: 在远端建表 + 灌数据")
    remote_dsn = LOCAL_DSN.replace("dbname=learn_pg", f"dbname={REMOTE_DBNAME}")
    with psycopg.connect(remote_dsn, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute("""
            CREATE TABLE ch16_archive_orders (
                id BIGSERIAL PRIMARY KEY,
                user_id BIGINT NOT NULL,
                amount NUMERIC(12,2),
                status TEXT,
                created_at TIMESTAMPTZ NOT NULL
            );
        """)
        cur.execute("""
            INSERT INTO ch16_archive_orders (user_id, amount, status, created_at)
            SELECT
                (random()*100+1)::BIGINT,
                (random()*1000)::NUMERIC(12,2),
                (ARRAY['paid','done'])[1+(random()*1)::INT],
                '2024-01-01'::timestamptz + (random() * INTERVAL '364 day')
            FROM generate_series(1, 10000);
        """)
        cur.execute("CREATE INDEX ON ch16_archive_orders (user_id);")
        cur.execute("CREATE INDEX ON ch16_archive_orders (created_at);")
        cur.execute("ANALYZE ch16_archive_orders;")
        print("  ✓ 远端表 ch16_archive_orders 建好,灌入 10,000 行")

    hr("Step 3: 在本地库装 postgres_fdw + 配置远端")
    with psycopg.connect(LOCAL_DSN, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute("CREATE EXTENSION IF NOT EXISTS postgres_fdw;")
        print("  ✓ 装载 postgres_fdw 扩展")

        # 解析 LOCAL_DSN 拿 host/port,假设远端 PG 也在同一个地址
        import re
        host = (re.search(r"host=(\S+)", LOCAL_DSN) or [None, "127.0.0.1"])[1]
        port = (re.search(r"port=(\S+)", LOCAL_DSN) or [None, "5432"])[1]
        user = (re.search(r"user=(\S+)", LOCAL_DSN) or [None, "postgres"])[1]
        pwd  = (re.search(r"password=(\S+)", LOCAL_DSN) or [None, "postgres"])[1]

        cur.execute("DROP SERVER IF EXISTS pg_archive CASCADE;")
        cur.execute(f"""
            CREATE SERVER pg_archive
            FOREIGN DATA WRAPPER postgres_fdw
            OPTIONS (host '{host}', port '{port}', dbname '{REMOTE_DBNAME}',
                     use_remote_estimate 'true', fetch_size '1000');
        """)
        print("  ✓ 远端 server 已注册(开启 use_remote_estimate)")

        cur.execute(f"""
            CREATE USER MAPPING FOR CURRENT_USER
            SERVER pg_archive
            OPTIONS (user '{user}', password '{pwd}');
        """)
        print("  ✓ 用户映射已建")

        cur.execute("DROP FOREIGN TABLE IF EXISTS ch16_archive_orders;")
        cur.execute("""
            CREATE FOREIGN TABLE ch16_archive_orders (
                id BIGINT,
                user_id BIGINT,
                amount NUMERIC(12,2),
                status TEXT,
                created_at TIMESTAMPTZ
            ) SERVER pg_archive
              OPTIONS (schema_name 'public', table_name 'ch16_archive_orders');
        """)
        print("  ✓ 外部表 ch16_archive_orders 已建")

        # 拉远端统计信息(重要!)
        cur.execute("ANALYZE ch16_archive_orders;")
        print("  ✓ ANALYZE 已拉远端统计信息")


def query():
    hr("跨库 SELECT + EXPLAIN VERBOSE 看下推")
    with psycopg.connect(LOCAL_DSN, autocommit=True) as conn, conn.cursor() as cur:
        sql_q = "SELECT * FROM ch16_archive_orders WHERE user_id = 50 LIMIT 5"
        print(f"\nSQL: {sql_q}")
        cur.execute(sql_q)
        for r in cur.fetchall():
            print(" ", r)

        print("\n=== EXPLAIN VERBOSE ===")
        cur.execute(f"EXPLAIN (VERBOSE) {sql_q}")
        for r in cur.fetchall():
            print(" ", r[0])
        print("\n  ↑ 注意 'Remote SQL':WHERE / LIMIT 都被下推到远程库执行")

        print("\n=== EXPLAIN 聚合下推(PG 10+)===")
        cur.execute("EXPLAIN (VERBOSE) SELECT count(*) FROM ch16_archive_orders WHERE status = 'paid'")
        for r in cur.fetchall():
            print(" ", r[0])


def join_demo():
    hr("跨库 JOIN:本地 users JOIN 远端 ch16_archive_orders")
    with psycopg.connect(LOCAL_DSN, autocommit=True) as conn, conn.cursor() as cur:
        # 建一张本地小表
        cur.execute("DROP TABLE IF EXISTS ch16_local_users;")
        cur.execute("""
            CREATE TABLE ch16_local_users (
                id BIGINT PRIMARY KEY,
                name TEXT
            );
        """)
        cur.execute("""
            INSERT INTO ch16_local_users (id, name)
            SELECT g, 'user-' || g FROM generate_series(1, 100) g;
        """)
        cur.execute("ANALYZE ch16_local_users;")

        sql_q = """
            SELECT u.name, count(o.id) AS order_cnt, sum(o.amount) AS total
            FROM ch16_local_users u
            JOIN ch16_archive_orders o ON o.user_id = u.id
            WHERE u.id <= 5
            GROUP BY u.name
            ORDER BY u.name;
        """
        print(f"\nSQL:\n{sql_q}")
        cur.execute(sql_q)
        for r in cur.fetchall():
            print(" ", r)

        print("\n=== EXPLAIN VERBOSE ===")
        cur.execute(f"EXPLAIN (VERBOSE, ANALYZE) {sql_q}")
        for r in cur.fetchall():
            print(" ", r[0])
        print("\n  ↑ 跨 server 的 JOIN 不能整体下推,PG 会从远端拉数据回来本地 JOIN")
        print("    所以远端 WHERE 越精确越好(这里 u.id <= 5 限制只拉 5 个用户的订单)")


def teardown():
    hr("清理")
    with psycopg.connect(LOCAL_DSN, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute("DROP FOREIGN TABLE IF EXISTS ch16_archive_orders;")
        cur.execute("DROP SERVER IF EXISTS pg_archive CASCADE;")
        cur.execute("DROP TABLE IF EXISTS ch16_local_users;")
        cur.execute("DROP EXTENSION IF EXISTS postgres_fdw;")
    with psycopg.connect(ADMIN_DSN, autocommit=True) as conn, conn.cursor() as cur:
        cur.execute(f"DROP DATABASE IF EXISTS {REMOTE_DBNAME};")
    print("  ✓ 已清理")


def main():
    cmd = sys.argv[1] if len(sys.argv) > 1 else ""
    fns = {"setup": setup, "query": query, "join": join_demo, "teardown": teardown}
    if cmd not in fns:
        print(f"Usage: {sys.argv[0]} {{{'|'.join(fns.keys())}}}", file=sys.stderr)
        return 1
    fns[cmd]()
    return 0


if __name__ == "__main__":
    sys.exit(main() or 0)
markdown
# 第 16 章 配套代码

本章演示 **声明式分区**(RANGE / HASH / LIST)、**分区裁剪****pg_partman 自动维护****postgres_fdw 跨库查询**等核心特性。所有自建表都加了 `ch16_` 前缀,避免和其他章节的 `orders / events` 等同名表打架。

## 准备工作

1. 跑一次 init.sql 初始化父表与分区子表:

   ```bash
   psql -h 127.0.0.1 -U postgres -d learn_pg -f ../init.sql

会创建 ch16_orders(按月 RANGE 分区)、ch16_events(HASH 分区)、ch16_countries(LIST 分区)等基础测试表。

  1. 安装 Python 客户端:

    bash
    pip install "psycopg[binary]>=3.1"
  2. (可选)通过环境变量覆盖默认连接:

    bash
    export PG_DSN="host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres"
  3. 05_postgres_fdw.py 需要可达的「远端」库(脚本里默认用本地的 learn_pg_archive 当远端,跑前先 createdb learn_pg_archive)。

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

脚本一句话说明关键 PG 特性
01_range_partition_demo.py按月 RANGE 建 ch16_demo_orders,演示路由、裁剪、DETACH CONCURRENTLYPARTITION BY RANGE、default 分区、pg_inherits
02_hash_partition_demo.py8 个 HASH 分区灌 10 万行,验证分布均匀性与裁剪行为PARTITION BY HASH、点查裁剪、范围查询无法裁剪
03_partition_pruning.py用 EXPLAIN 对比裁剪成功/失败 6 种 SQL 写法enable_partition_pruningpartitionwise_join、prepare 计划缓存
04_partman_simulation.py没装 pg_partman 也能跑,纯 Python 模拟自动建分区 + 滚动归档ch16_my_partman_configch16_ts_logs 时序表
05_postgres_fdw.py用 postgres_fdw 把冷数据归档到「远端」,演示 Remote SQL 下推CREATE SERVERCREATE USER MAPPING、外部表 + 分区结合

预期输出(节选 03)

======================================================================
  场景 A:等值查询(裁剪到 1 个分区)
======================================================================
  Append  (cost=0.00..xx rows=xx width=xx)
    ->  Seq Scan on ch16_orders_2025_03  ...
  ✓ 只扫到 1 个分区

常见报错与依赖

  • connection refused → PG 未启动,或 PG_DSN 端口写错。
  • relation "ch16_orders" does not exist → 没跑 ../init.sql
  • permission denied to create extension "postgres_fdw" → 05 脚本需要超级用户,或预先 CREATE EXTENSION postgres_fdw;
  • cannot drop partition ... because other objects depend on it → 子分区上有外键 / 视图引用,DETACH 前先解依赖。
  • pg_partman 真要用得 CREATE EXTENSION pg_partman;,本章 04 是模拟实现,不依赖该扩展。

01_range_partition_demo.py ↗ · 02_hash_partition_demo.py ↗ · 03_partition_pruning.py ↗ · 04_partman_simulation.py ↗ · 05_postgres_fdw.py ↗ · README.md ↗