主题
第 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_01、ch16_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+ 声明式分区:三种类型
声明式分区的三大要素:
- 分区父表:
CREATE TABLE ... PARTITION BY ...,不存数据,只是个虚表; - 分区键(partition key):父表上指定的「按哪一列拆」;
- 分区子表:
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; -- on4.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_14.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 |
| 100 | 5ms |
| 1000 | 50ms |
| 10000 | 500ms+(规划比执行还慢!) |
经验值:单父表分区数 ≤ 几百,最多 1000。
7.2 分区键选择决定能否裁剪
关键:最常用的查询条件 = 分区键。
| 业务 | 分区键 |
|---|---|
| 按月查订单 | created_at(RANGE) |
| 按用户读自己数据 | user_id(HASH) |
| 多租户 SaaS | tenant_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_102008、orders_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_fdw | MySQL |
oracle_fdw | Oracle |
file_fdw(内置) | 本地 CSV / TSV |
mongo_fdw | MongoDB |
parquet_fdw | Parquet 文件 |
redis_fdw | Redis |
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 全方位对比
| 维度 | MySQL | PostgreSQL |
|---|---|---|
| 内置分区类型 | RANGE / LIST / HASH / KEY | RANGE / LIST / HASH(PG 11+) |
| 多级分区 | 仅 SUBPARTITION(最多 2 层) | 任意层级 |
| 默认分区 | 部分版本支持 | 原生 DEFAULT 分区 |
| 分区裁剪 | 支持 | 支持,且支持执行期裁剪(PG 11+) |
| 分区索引级联 | 自动 | PG 11+ 自动 |
| 主键限制 | 不能用任意列 | 主键必须包含所有分区键 |
| 跨分区 JOIN 优化 | 较弱 | partitionwise_join(PG 11+) |
| Sharding 中间件 | ShardingSphere / Vitess / MyCAT | Citus / FDW |
| MPP 数据仓库 | TiDB / OceanBase 同 wire | Greenplum |
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 句话总结:
- 分区 = 单实例内拆,PG 自动路由;分库 = 多实例分散,需要中间件。
- PG 10+ 声明式分区有 RANGE / LIST / HASH 三种,多级可嵌套。
- 分区键必须包含在主键里,这是 PG 的硬约束。
- DEFAULT 分区做兜底,但要避免它装实际业务数据。
- 分区裁剪是分区性能的核心,WHERE 必须直接对分区键操作。
- PG 11+ 父表索引会级联到所有子分区,无需手动到每个分区建索引。
- DETACH CONCURRENTLY(PG 14+)实现不锁表的分区下线。
- pg_partman 自动维护时序分区,配合 pg_cron 做调度。
- 分区数控制在 几十到几百,过多反而拖慢规划器。
- 真正的分库分表用 Citus(推荐)/ FDW(轻量)/ Greenplum(MPP)。
🎮 配套演示
用浏览器打开
./16_partition/demo.html,跟着可视化动画再走一遍本章核心概念。配套代码在
./16_partition/code/,每个脚本都可以独立python xxx.py运行,先跑init.sql准备数据。
13. 面试高频题
Q1:PG 10+ 声明式分区与「继承表 + 触发器」相比,优势是什么?
考察点:分区机制演进、声明式分区原理。
标准答案:
PG 10 之前实现分区的「老办法」是:
- 父表是普通表(也存数据);
- 子表
INHERITS (parent)继承父表结构; - 写一个 INSERT 触发器,根据值决定 redirect 到哪个子表;
- 给子表加
CHECK约束,让 constraint exclusion 做裁剪。
老办法的痛点:
- 触发器开销:每次 INSERT 都跑一遍触发器函数,1000 行/秒就是 1000 次 plpgsql 调用,CPU 占用高;
- CHECK 约束维护:子表加错约束就插不进,运维事故频发;
- 约束剪枝粒度粗:依赖
constraint_exclusion = on,规划期才裁,且必须严格表达式匹配; - INSERT 路由不在内核:每行单独 trigger,无法批量优化;
- 无法保证父表只读:父表自己可能也有数据(必须自己约束不写)。
声明式分区(PG 10+)的优势:
- 内核原生路由:
INSERT INTO parent直接由 executor 内核根据分区键算出目标子表,无触发器开销,单条 INSERT 性能提升数倍; - 父表自动只读:父表是
relkind='p'(partitioned),完全不存数据,保证逻辑表的纯净; - 更强的剪枝:PG 11+ 支持执行期剪枝(
Subplans Removed),prepared statement、子查询、运行时参数都能剪; - 索引级联:PG 11+ 父表建索引自动下发所有子分区;
- DETACH CONCURRENTLY:PG 14+ 不锁表下线分区;
- Partitionwise JOIN/Aggregate:跨分区操作下推优化。
加分项:能讲出「老办法依然在用」的场景——pg_pathman 这类扩展用 hook 拦截 planner,比早期声明式分区在「分区数极多时」性能还好;不过 PG 13+ 后差距已经很小,新项目无脑用声明式即可。
Q2:分区裁剪是什么?计划期裁剪和执行期裁剪有什么区别?怎么用 EXPLAIN 验证?
考察点:分区性能核心机制。
标准答案:
分区裁剪(Partition Pruning):根据 WHERE 条件自动跳过不可能命中的分区,只扫相关分区,是分区表性能的关键。
两种时机:
计划期裁剪(Planning-Time Pruning):WHERE 是字面量常量,规划器编译时就知道值,直接在 Plan 里去掉不相关分区。
EXPLAIN输出里你根本看不到被剪掉的分区。执行期裁剪(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_idhash 分布; - 没有自然区间的列:UUID、订单号 hash 等。
典型场景对照:
| 业务 | 分区类型 | 分区键 |
|---|---|---|
| 电商订单(看最近一月) | RANGE | created_at(月) |
| 微信日志(按天清理) | RANGE | log_date(日) |
| SaaS 租户表 | LIST 或 HASH | tenant_id |
| IM 消息表(写均匀) | HASH | user_id |
| 全国门店报表 | LIST | province_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_orders和ch16_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_fdw、oracle_fdw、sqlserver_fdw:其他 RDBMS;mongo_fdw、redis_fdw、elasticsearch_fdw:NoSQL;file_fdw(内置)、csv_fdw、parquet_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):
- WHERE 下推:
WHERE user_id = 100会原样发到远程库执行,只返回符合条件的行; - 投影下推:只 SELECT 需要的列,不传无关列;
- JOIN 下推(PG 9.6+):如果 JOIN 两边都在同一 foreign server 上,PG 可以把整个 JOIN 下推到远程执行,本地只接收结果;
- 聚合下推(PG 10+):
SELECT count(*) FROM remote_orders直接在远程 count; - 排序 / 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 社区最流行的分区自动维护扩展,主要做三件事:
- 预建未来分区:避免「写入时分区不存在」;
- 删除过期分区:按保留策略自动 detach + drop;
- 管理 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 detach | detach 后归档到冷存储 |
infinite_time_partitions | 持续创建未来分区 | true |
经典翻车场景:
p_premake太小,维护 cron 挂了:cron 失败几天,写入时间超过预建范围 → 落入 default 分区 → default 越来越大;- 没装 pg_cron 也没设置 BGW:以为「装了 partman 就自动跑了」,其实需要手动调
run_maintenance(); - default 分区有数据后建新分区超慢:partman 添加新分区时要扫 default 找冲突行,几十 GB default 加一次分区要锁几十分钟;
retention_keep_table = false误删数据:保留策略写错(比如把 30 days 写成 30 months),第二天醒来发现两年的分区被 drop 了;强烈建议先retention_keep_table = true跑一段时间,确认没事再改 false;- 大事务 + 分区 detach 死锁:detach 拿 ACCESS EXCLUSIVE 锁,正在跑的长 SELECT 会挡住;PG 14+ 可以用
partman.partition_data_proc+DETACH CONCURRENTLY; - 分区数失控:
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。
🔗 延伸阅读
- 第 9 章 物理存储与堆表结构:理解每个分区子表底层就是独立的堆文件,知道为什么分区能加速 VACUUM、缓解膨胀。
- 第 6 章 索引与查询优化:分区索引是 PG 11+ 的杀手锏,配合本章的级联建索引一起看。
- 第 17 章 性能调优:分区裁剪与计划缓存(
prepare/plan_cache_mode)的互动,分区数对规划器开销的影响。
🎬 可视化演示
演示加载缓慢或样式异常?点此在新标签页打开 ↗
💻 示例代码
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 分区)等基础测试表。
安装 Python 客户端:
bashpip install "psycopg[binary]>=3.1"(可选)通过环境变量覆盖默认连接:
bashexport PG_DSN="host=127.0.0.1 port=5432 dbname=learn_pg user=postgres password=postgres"05_postgres_fdw.py 需要可达的「远端」库(脚本里默认用本地的
learn_pg_archive当远端,跑前先createdb learn_pg_archive)。
脚本一览(推荐运行顺序)
| 脚本 | 一句话说明 | 关键 PG 特性 |
|---|---|---|
01_range_partition_demo.py | 按月 RANGE 建 ch16_demo_orders,演示路由、裁剪、DETACH CONCURRENTLY | PARTITION BY RANGE、default 分区、pg_inherits |
02_hash_partition_demo.py | 8 个 HASH 分区灌 10 万行,验证分布均匀性与裁剪行为 | PARTITION BY HASH、点查裁剪、范围查询无法裁剪 |
03_partition_pruning.py | 用 EXPLAIN 对比裁剪成功/失败 6 种 SQL 写法 | enable_partition_pruning、partitionwise_join、prepare 计划缓存 |
04_partman_simulation.py | 没装 pg_partman 也能跑,纯 Python 模拟自动建分区 + 滚动归档 | ch16_my_partman_config、ch16_ts_logs 时序表 |
05_postgres_fdw.py | 用 postgres_fdw 把冷数据归档到「远端」,演示 Remote SQL 下推 | CREATE SERVER、CREATE 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 ↗