Skip to content

第 5 章 高级查询:从 JOIN 到窗口函数的完整地图

目标读者:会写 SELECT * FROM t WHERE id = 1,但一到多表连接、子查询嵌套、窗口分析就懵圈的同学。

学完你会:能独立写出「用户最近 N 笔订单」「销售排名」「组织架构树」「同比环比」「UPSERT」等真实业务 SQL,并能看懂 EXPLAIN 输出做性能调优。


0. 导读:为什么高级查询这么重要?

日常业务里,「查单表」只占开发工作量的不到 20%,剩下的 80% 都是:

  • 把多张表「拼」起来(JOIN)
  • 把一个问题「拆」成小问题递进回答(子查询、CTE)
  • 对结果集做「二次加工」(排序排名、同比、移动平均)
  • 做「写入 + 查询」的组合动作(INSERT ... RETURNING、UPSERT)

很多人 SQL 写 5 年还只停留在「循环拼 SQL + 在业务代码里做关联」,其实是把数据库当 KV 用,既浪费机器又慢。

PostgreSQL 在高级查询这一块,几乎是业界 SQL 标准的「范本」WITH RECURSIVE、窗口函数、LATERALDISTINCT ONRETURNING、可写 CTE,每一个都能让你的代码量砍一半。

下面我们按「由易到难」逐个拿下。


1. JOIN 全家桶

1.1 概念 + 生活类比

想象你有两本小本子:

  • 本子 A:用户名册users,谁的手机号是谁)
  • 本子 B:订单流水orders,哪个用户下了哪张单)

JOIN 就是「拿一张表的某一列去另一张表里对号入座」。对得上就拼到一起。

不同 JOIN 的区别,一张图说清:

A (用户)     B (订单)
 1  Alice     1  100  
 2  Bob       1  200
 3  Carol     2  300
                4  400  (订单属于用户4,而用户4不存在)
JOIN 类型结果口诀
INNER JOIN{1-100, 1-200, 2-300}两边都有才保留
LEFT JOIN{1-100, 1-200, 2-300, 3-NULL}左边全保,右边缺的补 NULL
RIGHT JOIN{1-100, 1-200, 2-300, NULL-400}右边全保,左边缺的补 NULL
FULL OUTER JOIN{1-100, 1-200, 2-300, 3-NULL, NULL-400}两边都保,缺的补 NULL
CROSS JOIN3×4=12 行笛卡儿积(不写 ON)

1.2 SQL 示例

sql
-- INNER JOIN(最常见)
SELECT u.name, o.amount
FROM ch5_users u
INNER JOIN ch5_orders o ON u.id = o.user_id;

-- LEFT JOIN:统计每个用户下单数(包括 0 单)
SELECT u.name, COUNT(o.id) AS order_cnt
FROM ch5_users u
LEFT JOIN ch5_orders o ON u.id = o.user_id
GROUP BY u.name;

-- FULL OUTER JOIN:找出「有用户无订单」+「有订单无用户」的脏数据
SELECT u.id AS uid, o.id AS oid
FROM ch5_users u
FULL OUTER JOIN ch5_orders o ON u.id = o.user_id
WHERE u.id IS NULL OR o.id IS NULL;

-- CROSS JOIN:生成「所有用户 x 所有月份」报表骨架
SELECT u.name, m.month
FROM ch5_users u
CROSS JOIN generate_series(1, 12) AS m(month);

1.3 ON / USING / NATURAL JOIN

sql
-- ON:最灵活,任意表达式
SELECT * FROM a JOIN b ON a.uid = b.user_id AND b.status = 'paid';

-- USING:两表同名列,写一次即可,结果里只保留一列
SELECT * FROM a JOIN b USING (id);   -- 要求 a.id 和 b.id 同名

-- NATURAL JOIN:自动找所有同名列做等值连接(慎用,改表结构会坏事)
SELECT * FROM a NATURAL JOIN b;

📌 与 MySQL 的区别

  • MySQL 8.0 之前不支持 FULL OUTER JOIN,需要 LEFT UNION RIGHT 模拟;PG 原生支持。
  • MySQL 的 JOIN / INNER JOIN / CROSS JOIN 在不带 ON 时行为一致,PG 严格区分:CROSS JOIN 专指笛卡儿积。

1.4 自连接(Self Join)

同一张表自己连自己,常见于「找推荐人」「比较上下级」:

sql
-- 员工表里找出 员工 + 他的直属经理
SELECT e.name AS emp, m.name AS mgr
FROM ch5_employees e
LEFT JOIN ch5_employees m ON e.manager_id = m.id;

1.5 多表 JOIN 与优化原则

sql
SELECT u.name, p.name, oi.qty
FROM ch5_orders o
JOIN ch5_order_items oi ON oi.order_id = o.id
JOIN ch5_products p     ON p.id = oi.product_id
JOIN ch5_users u        ON u.id = o.user_id
WHERE o.created_at >= '2025-01-01'
ORDER BY o.created_at DESC
LIMIT 100;

3 条黄金原则

  1. 小表带大表:PG 优化器一般自己选 Hash Join 的 build 端,但过滤条件尽早下推(把 WHERE o.created_at >= ... 写清楚,让 orders 先过滤,再去 JOIN)。
  2. JOIN 列必须有索引:否则只能走 Hash JoinMerge Join,一旦数据倾斜就 OOM 或拖垮磁盘。
  3. 返回列「用多少取多少」SELECT * 会导致 Index Only Scan 失效,宽行也会让排序/哈希的 work_mem 被撑爆。

2. 子查询(Subquery)

按「返回什么形状」分四类:

子查询类型返回典型用法
标量子查询1 行 1 列SELECT (SELECT MAX(amt) FROM ch5_orders) ...
行子查询1 行 N 列WHERE (a,b) = (SELECT x,y FROM ...)
列子查询N 行 1 列WHERE uid IN (SELECT id FROM vip)
表子查询N 行 N 列FROM (SELECT ...) AS t

按「能否独立执行」分两类:

  • 非相关子查询:内层不依赖外层,执行一次即可。
  • 相关子查询(Correlated):内层引用外层表的列,外层每一行都触发一次内层,很容易写成 N+1 性能陷阱。

2.1 EXISTS vs IN vs JOIN

经典面试题:查出下过单的用户。3 种写法:

sql
-- 写法 1:IN
SELECT * FROM ch5_users
WHERE id IN (SELECT user_id FROM ch5_orders);

-- 写法 2:EXISTS(相关子查询)
SELECT * FROM ch5_users u
WHERE EXISTS (SELECT 1 FROM ch5_orders o WHERE o.user_id = u.id);

-- 写法 3:JOIN + DISTINCT
SELECT DISTINCT u.*
FROM ch5_users u JOIN ch5_orders o ON u.id = o.user_id;

怎么选?

  • EXISTS短路求值,找到第一条就返回,小表驱动大表时快。
  • IN:优化器通常把它重写成 Semi-Join,EXISTS 在 PG 里性能几乎一样(PG 11+ 后差别微乎其微)。
  • JOIN + DISTINCT:要去重,反而多一步。能用 EXISTS / IN 就不要 JOIN 去重。

📌 与 MySQL 的区别 MySQL 5.6 以前 IN (subquery) 会退化成 DependentSubquery,每行触发一次。PG 从来没这个坑,放心用。

2.2 NOT IN 的 NULL 陷阱(面试必考)

sql
-- 找出没下过单的用户
SELECT * FROM ch5_users WHERE id NOT IN (SELECT user_id FROM ch5_orders);

:如果 orders.user_id 里出现了 NULL,整个结果会变成空集

原因:id NOT IN (1, 2, NULL) 等价于 id <> 1 AND id <> 2 AND id <> NULL,而 id <> NULL 永远是 UNKNOWN,导致整行被过滤。

修复:用 NOT EXISTS 代替,它把 NULL 当「找不到」处理:

sql
SELECT * FROM ch5_users u
WHERE NOT EXISTS (
  SELECT 1 FROM ch5_orders o WHERE o.user_id = u.id
);

3. CTE(公共表表达式)

3.1 概念

CTE 就是给一段子查询起个名字,像变量一样在主查询里引用。语法:

sql
WITH 别名 AS (
  SELECT ...
)
SELECT ... FROM 别名;

生活类比:把 SQL 当做菜菜谱,子查询就是「先切土豆丝」这一步,CTE 就是你把切好的丝放进一个盘子并贴标签,后面的菜随时拿来用,比起在最终那口锅里现场切要清晰得多。

3.2 基础用法 + 多个 CTE

sql
WITH 
  paid_orders AS (
    SELECT * FROM ch5_orders WHERE status = 'paid'
  ),
  user_total AS (
    SELECT user_id, SUM(amount) AS total
    FROM paid_orders
    GROUP BY user_id
  )
SELECT u.name, t.total
FROM user_total t
JOIN ch5_users u ON u.id = t.user_id
WHERE t.total > 1000
ORDER BY t.total DESC;

好处

  1. 可读性极高,代码逻辑一条线读下来。
  2. 一个 CTE 可在下游多次引用
  3. 复杂报表可以按「步骤」拆解,远比层层嵌套子查询清晰。

3.3 可写 CTE(PG 特色!)

这是 PG 最酷的特性之一:在 CTE 里写 INSERT / UPDATE / DELETE 并配合 RETURNING,结果可以继续在主查询里被引用

场景:一次 SQL 同时完成「移档 + 记日志」:

sql
-- 把 2020 年之前的订单从 orders 迁到归档表 orders_archive,并返回迁移数量
WITH moved AS (
  DELETE FROM ch5_orders
  WHERE created_at < '2020-01-01'
  RETURNING *
),
ins AS (
  INSERT INTO ch5_orders_archive
  SELECT * FROM moved
  RETURNING id
)
SELECT COUNT(*) AS moved_cnt FROM ins;

等价逻辑在 MySQL 里需要 3 条独立 SQL + 手动事务,PG 一条搞定,且原子性天然保证。

注意:可写 CTE 的执行顺序由 PG 决定,不保证 moved 先于 ins 的「可见性」。所有 CTE 看到的是同一个快照(查询开始时刻)。需要严格串行就别写可写 CTE。

3.4 WITH MATERIALIZED / NOT MATERIALIZED(PG 12+)

  • PG 11 及以前:CTE 永远被「物化」成临时结果(类似虚拟表),优化器不能把外层 WHERE 下推进去。
  • PG 12+:默认「能内联就内联」,CTE 像视图一样展开到主查询里优化。
  • 需要强制物化(例如 CTE 是有副作用的 INSERT,或想故意切断优化器):
sql
WITH raw AS MATERIALIZED (SELECT ... FROM huge_table WHERE ...)
SELECT ...;
  • 阻止物化
sql
WITH raw AS NOT MATERIALIZED (SELECT ...) SELECT ...;

经验:PG 12+ 不用动,默认策略很聪明。遇到 CTE 比等价子查询慢,才加 NOT MATERIALIZED


4. 递归 CTE:树形结构查询神器

4.1 场景

  • 组织架构:员工 → 经理 → 部门总监 → CEO
  • 评论树:一级评论 → 回复 → 回复的回复
  • BOM 物料清单:成品 → 零件 → 原材料
  • 地区:国家 → 省 → 市 → 区

这类「自引用」数据,用普通 JOIN 做不了(因为你不知道树有多深)。

4.2 语法

sql
WITH RECURSIVE cte_name AS (
  -- 1) 基础(anchor)查询:起始行
  SELECT ...
  UNION ALL
  -- 2) 递归查询:引用 cte_name 自己
  SELECT ...
  FROM t JOIN cte_name ON ...
)
SELECT * FROM cte_name;

执行过程:

  1. 先跑 anchor → 得到第 0 层
  2. 拿第 0 层的每一行,喂给递归部分 → 得到第 1 层
  3. 拿第 1 层再喂 → 第 2 层……
  4. 直到某一轮返回 0 行,停止。
  ┌────────────┐  anchor: CEO
  │  Level 0   │
  └─────┬──────┘
        │ 递归
  ┌─────┴──────┐  Level 1: 总监们
  │  Level 1   │
  └─────┬──────┘
        │ 递归
  ┌─────┴──────┐  Level 2: 经理们
  │  Level 2   │
  └────────────┘

4.3 示例:员工组织架构展开

表结构(见本章 init.sql):

sql
CREATE TABLE ch5_employees (
  id        INT PRIMARY KEY,
  name      TEXT,
  manager_id INT REFERENCES ch5_employees(id)
);

查询:从 CEO(manager_id IS NULL)开始,展开整棵树,标上层级和路径

sql
WITH RECURSIVE org AS (
  SELECT id, name, manager_id, 1 AS lvl,
         name::TEXT AS path
  FROM ch5_employees
  WHERE manager_id IS NULL               -- 1) anchor:根节点
  UNION ALL
  SELECT e.id, e.name, e.manager_id, o.lvl + 1,
         o.path || ' > ' || e.name
  FROM ch5_employees e
  JOIN org o ON e.manager_id = o.id      -- 2) 递归:子节点
)
SELECT lpad('', (lvl-1)*2) || name AS tree, path
FROM org
ORDER BY path;

输出(psql 示意):

      tree       |            path
-----------------+-------------------------------
 CEO             | CEO
   VP Sales      | CEO > VP Sales
     Sales Mgr A | CEO > VP Sales > Sales Mgr A
     Sales Mgr B | CEO > VP Sales > Sales Mgr B
   VP Eng        | CEO > VP Eng
     Eng Mgr     | CEO > VP Eng > Eng Mgr
       Alice     | CEO > VP Eng > Eng Mgr > Alice
       Bob       | CEO > VP Eng > Eng Mgr > Bob

4.4 防死循环

如果数据有环(A 的上级是 B,B 的上级是 A),递归 CTE 会无限增长直到 OOM。防御方法:

  1. 限制深度:递归里加 WHERE o.lvl < 100
  2. 路径去重:在路径里检测 NOT (e.id = ANY(o.path_ids))
  3. CYCLE 子句(PG 14+):
sql
WITH RECURSIVE t AS (
  ...
) CYCLE id SET is_cycle USING path
SELECT * FROM t WHERE NOT is_cycle;

5. 窗口函数(Window Function)

5.1 什么是窗口函数?

一句话:像聚合函数(SUM/COUNT)那样计算,但不把行合并掉

对比:

sql
-- 聚合:1000 行 → 1 行
SELECT SUM(amount) FROM ch5_orders;

-- 窗口:1000 行 → 1000 行,每行都带一个累加值
SELECT id, amount, SUM(amount) OVER (ORDER BY id) AS running_total
FROM ch5_orders;

生活类比:聚合是「班里所有人成绩加起来 → 一个总分」;窗口是「每个人身上都贴一张小纸条,写着『你以及排你前面的所有人总分』」。

5.2 核心语法:OVER (PARTITION BY ... ORDER BY ... 窗口框架)

sql
函数名(参数) OVER (
  [PARTITION BY 分区列]     -- 把数据分桶,每个桶内独立计算
  [ORDER BY 排序列]         -- 桶内排序(排名、偏移、移动窗口都依赖它)
  [ROWS/RANGE BETWEEN ...]  -- 窗口框架:当前行的「邻居」范围
)

PARTITION BY 类似 GROUP BY,但不合并行。ORDER BY 在窗口里的作用是定义「顺序」,不是最终结果的排序。

5.3 排名函数

函数相同值处理下一个排名
ROW_NUMBER()不并列,强制唯一连续
RANK()并列跳号(1,1,3)
DENSE_RANK()并列不跳(1,1,2)
NTILE(n)均分 n 桶
sql
-- 每个类目里的商品按价格排名
SELECT category, name, price,
       ROW_NUMBER() OVER (PARTITION BY category ORDER BY price DESC) AS rn,
       RANK()       OVER (PARTITION BY category ORDER BY price DESC) AS rk,
       DENSE_RANK() OVER (PARTITION BY category ORDER BY price DESC) AS drk
FROM ch5_products;

经典问题:每个类目的 TOP 3

sql
SELECT * FROM (
  SELECT category, name, price,
         ROW_NUMBER() OVER (PARTITION BY category ORDER BY price DESC) AS rn
  FROM ch5_products
) t WHERE rn <= 3;

5.4 偏移函数

函数作用
LAG(col, n, default)取同分区内当前行往前第 n 行的值
LEAD(col, n, default)取同分区内当前行往后第 n 行的值
FIRST_VALUE(col)窗口首行的值
LAST_VALUE(col)窗口末行的值(默认框架坑见下)
NTH_VALUE(col, n)第 n 行的值

同比 / 环比

sql
-- 每月销售额 + 上月销售额 + 环比
SELECT month, revenue,
       LAG(revenue)  OVER (ORDER BY month) AS last_month,
       (revenue - LAG(revenue) OVER (ORDER BY month))::NUMERIC
         / NULLIF(LAG(revenue) OVER (ORDER BY month), 0) AS mom
FROM ch5_monthly_sales;

5.5 聚合当窗口用 + 窗口框架

sql
-- 7 日移动平均
SELECT day, revenue,
       AVG(revenue) OVER (
         ORDER BY day
         ROWS BETWEEN 6 PRECEDING AND CURRENT ROW
       ) AS ma7
FROM daily_sales;

窗口框架 2 种模式

  • ROWS BETWEEN N PRECEDING AND M FOLLOWING按物理行数
  • RANGE BETWEEN ... INTERVAL ... PRECEDING AND ...按值距离(PG 11+ 支持 INTERVAL

默认框架

  • ORDER BY 时:RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW
  • ORDER BY 时:RANGE BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING(整个分区)

LAST_VALUE 的坑:默认框架是「起点到当前行」,所以 LAST_VALUE 永远等于当前行的值!要正确用需要显式写:

sql
LAST_VALUE(price) OVER (
  PARTITION BY category ORDER BY day
  ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
)

5.6 窗口函数执行顺序

SQL 语句的逻辑执行顺序:

FROM → JOIN → WHERE → GROUP BY → HAVING → 窗口函数 → SELECT → DISTINCT → ORDER BY → LIMIT

窗口函数发生在 HAVING 之后、SELECT 之前,所以:

  • 不能在 WHERE 里直接写窗口函数(窗口还没算呢)
  • 需要对窗口结果再过滤 → 用子查询或 CTE 包一层
sql
-- 错:WHERE rn <= 3 —— rn 还不存在
-- 对:用子查询 / CTE 包一层
WITH ranked AS (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY cat ORDER BY price DESC) AS rn FROM ch5_products
)
SELECT * FROM ranked WHERE rn <= 3;

6. 分组扩展:GROUPING SETS / ROLLUP / CUBE

6.1 需求:一次性出「分维度小计 + 总计」

老板要看「每个类目 + 每个省 + 总计 + 类目小计 + 省小计」。普通写法:3 个 UNION ALL,每个一次扫表。

用分组扩展一次 SQL 搞定:

sql
-- GROUPING SETS:自由组合要统计的维度
SELECT category, province, SUM(amount)
FROM sales
GROUP BY GROUPING SETS ( (category, province), (category), (province), () );
-- 4 个小计:类目+省 / 只类目 / 只省 / 总计

-- ROLLUP:层级小计(从右往左逐级卷起)
SELECT category, province, SUM(amount)
FROM sales
GROUP BY ROLLUP (category, province);
-- 等价于 GROUPING SETS ( (cat,prov), (cat), () )

-- CUBE:所有维度组合
SELECT category, province, SUM(amount)
FROM sales
GROUP BY CUBE (category, province);
-- 等价于 GROUPING SETS ( (cat,prov), (cat), (prov), () )

6.2 GROUPING() 函数分辨小计类型

小计行里非分组列会是 NULL,怎么区分「真 NULL」和「小计标记」?

sql
SELECT
  CASE WHEN GROUPING(category)=1 THEN '所有类目' ELSE category END AS category,
  CASE WHEN GROUPING(province)=1 THEN '所有省份' ELSE province END AS province,
  SUM(amount)
FROM sales
GROUP BY ROLLUP (category, province);

GROUPING(col) 返回 0/1:0 = 正常分组值,1 = 小计行的「泛化」值。


7. RETURNING 子句(PG 特色)

DML 语句返回受影响行,不用二次 SELECT

sql
-- INSERT 返回新生成主键
INSERT INTO ch5_users (name, email) VALUES ('Alice', 'a@x.com')
RETURNING id, created_at;

-- UPDATE 返回改前 / 改后值(配合子查询可以返回改前)
UPDATE ch5_orders SET status='paid' WHERE id=100
RETURNING id, status, updated_at;

-- DELETE 返回删掉的整行
DELETE FROM tmp_session WHERE expire_at < NOW()
RETURNING *;

-- 与 CTE 联动,一条 SQL 做两件事
WITH deleted AS (
  DELETE FROM cart WHERE user_id=1 RETURNING product_id, qty
)
INSERT INTO ch5_orders(user_id, product_id, qty)
SELECT 1, product_id, qty FROM deleted;

📌 与 MySQL 的区别

  • MySQL 没有通用 RETURNING(MariaDB 10.5+ 才有)。
  • MySQL 8.0 的 INSERT ... RETURNING 需要 MariaDB 才有;原生 MySQL 只能靠 LAST_INSERT_ID() + 再查。
  • 想知道「改了哪些行」,MySQL 只能先 SELECTUPDATE,有并发安全隐患(需要 FOR UPDATE)。

8. ON CONFLICT:PG 的 UPSERT(插入或更新)

8.1 语法

sql
INSERT INTO ch5_users (id, name, email) VALUES (1, 'Alice', 'a@x.com')
ON CONFLICT (id) DO UPDATE
SET name = EXCLUDED.name,
    email = EXCLUDED.email,
    updated_at = NOW();

-- 冲突时什么也不做(忽略)
INSERT INTO t VALUES (...) ON CONFLICT DO NOTHING;

-- 基于部分索引冲突
INSERT INTO t (..) ON CONFLICT (email) WHERE deleted_at IS NULL DO UPDATE ...;

EXCLUDED 伪表:指代「本次企图插入但冲突的那一行」。

8.2 和 MySQL INSERT ... ON DUPLICATE KEY UPDATE 对比

维度PG ON CONFLICTMySQL ON DUPLICATE KEY UPDATE
指定冲突列必须指定 (col)ON CONSTRAINT name无法指定,所有唯一键冲突都触发
冲突值引用EXCLUDED.colVALUES(col)(8.0.20 后标记废弃,改用别名)
冲突处理DO UPDATE / DO NOTHING只能 UPDATE
部分索引冲突✅ 支持 WHERE 条件
返回行RETURNING

9. LATERAL JOIN:横向连接

9.1 动机

普通子查询 / JOIN 里,右表不能引用左表的列

sql
-- 错误!子查询里的 u.id 不可见
SELECT u.name,
  (SELECT o.id FROM ch5_orders o WHERE o.user_id = u.id ORDER BY o.created_at DESC LIMIT 3)
FROM ch5_users u;

LATERAL 打破这个限制:右侧的 FROM / 子查询可以引用左侧的列,相当于「对左表每一行都执行一次右边的查询」。

9.2 经典场景:每个用户最近 N 笔订单

sql
SELECT u.id, u.name, o.id AS order_id, o.amount, o.created_at
FROM ch5_users u
LEFT JOIN LATERAL (
  SELECT id, amount, created_at
  FROM ch5_orders
  WHERE user_id = u.id
  ORDER BY created_at DESC
  LIMIT 3                      -- 这里的 LIMIT 才是「每用户 3 笔」
) o ON true;

如果没有 LATERAL,你需要用窗口函数 ROW_NUMBER() <= 3 变通,且难免全表排序。

9.3 另一个经典:UNNEST 展开数组

sql
-- 给每个订单展开标签数组
SELECT o.id, tag
FROM ch5_orders o,
     LATERAL unnest(o.tags) AS tag;

没有 LATERAL 关键字时,PG 也会隐式做 LATERAL(因为 unnest(o.tags) 明显引用了左表),但显式写上可读性更好。


10. DISTINCT ON:按字段去重并保留极值行

10.1 问题

每个用户的最新订单,只要一行。普通 DISTINCT 做不到(它对整行去重)。

10.2 PG 标准写法

sql
SELECT DISTINCT ON (user_id) user_id, id, amount, created_at
FROM ch5_orders
ORDER BY user_id, created_at DESC;   -- DISTINCT ON 的字段必须是 ORDER BY 开头

含义:对相同 user_id,只保留 ORDER BY 排序的第一行

MySQL 要写:

sql
-- MySQL 写法:用窗口函数 + 子查询
SELECT * FROM (
  SELECT *, ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY created_at DESC) AS rn
  FROM ch5_orders
) t WHERE rn = 1;

一眼看懂:PG 3 行,MySQL 5 行。


11. 集合运算:UNION / INTERSECT / EXCEPT

sql
-- 并集:两边合起来,去重
SELECT name FROM customers
UNION
SELECT name FROM suppliers;

-- 并集:不去重(快得多)
SELECT ... UNION ALL SELECT ...;

-- 交集
SELECT ... INTERSECT SELECT ...;

-- 差集:A 有但 B 没有
SELECT ... EXCEPT SELECT ...;

要求:两边列数、类型必须兼容。排序只对最终结果生效:

sql
(SELECT ... LIMIT 10) UNION ALL (SELECT ...) ORDER BY id LIMIT 20;

📌 与 MySQL 的区别:MySQL 8 才有 INTERSECT / EXCEPT(前者之前叫 EXCEPT),且性能不如 PG。


12. TABLESAMPLE:随机抽样(PG 9.5+)

给亿级表做问题排查、报表预览,不想扫全表:

sql
-- BERNOULLI:逐行投硬币(准确但慢),抽样 1%
SELECT * FROM ch5_orders TABLESAMPLE BERNOULLI (1);

-- SYSTEM:按页随机(快但聚集倾斜),抽样 1%
SELECT * FROM ch5_orders TABLESAMPLE SYSTEM (1);

-- 可重复:REPEATABLE (seed)
SELECT * FROM ch5_orders TABLESAMPLE SYSTEM (1) REPEATABLE (42);

13. 底层原理:高级查询如何被执行?

来一段 JOIN + 子查询的执行计划解读:

sql
EXPLAIN (ANALYZE, BUFFERS)
SELECT u.name, o.amount
FROM ch5_users u
JOIN ch5_orders o ON o.user_id = u.id
WHERE o.amount > 500 AND u.city = 'Beijing'
ORDER BY o.amount DESC
LIMIT 10;

典型输出(批注版):

Limit  (cost=1024.55..1024.58 rows=10 width=40) 
       (actual time=3.241..3.243 rows=10 loops=1)
  Buffers: shared hit=42
  ->  Sort  (cost=1024.55..1026.05 rows=600 width=40) 
            (actual time=3.240..3.241 rows=10 loops=1)
        Sort Key: o.amount DESC
        Sort Method: top-N heapsort  Memory: 26kB
        ->  Hash Join  (cost=24.25..1011.62 rows=600 width=40)
              Hash Cond: (o.user_id = u.id)
              Buffers: shared hit=42
              ->  Seq Scan on ch5_orders o  (cost=0..950.00 rows=3000 width=16)
                    Filter: (amount > 500)
                    Rows Removed by Filter: 7000
              ->  Hash  (cost=22.00..22.00 rows=180 width=32)
                    ->  Index Scan using idx_ch5_users_city on ch5_users u
                          Index Cond: (city = 'Beijing'::text)
Planning Time: 0.412 ms
Execution Time: 3.280 ms

关键字段逐行解读

字段含义
cost=1024.55..1024.58优化器估算的启动/总代价(代价单位是抽象的,数值大约 = 读 1 页顺序块的代价 = 1.0)
rows=10估计返回行数(和实际对比看统计信息准不准)
width=40平均每行字节数
actual time=3.241..3.243实际启动 / 总耗时(ms)
loops=1这个节点被执行了几次(Nested Loop 的右侧会 > 1)
Buffers: shared hit=42共享缓冲区命中 42 页,read 表示磁盘读(hit 远多于 read 才健康)
Sort Method: top-N heapsort因为 LIMIT 10,PG 用堆排只保留 top N,比普通 Sort 省内存
Hash Join / Hash Cond走哈希连接,先构建 users 的哈希表,再探测 orders

排错口诀

  1. rows 估算和 actual 差多少 → 差得多:ANALYZE 或建扩展统计信息。
  2. 看有没有 Seq Scan 在大表上 → 确认是否该加索引。
  3. Buffers: read 是否远大于 hit → 缓存没装下这部分数据。
  4. Sort 有没有 external merge Diskwork_mem 小了。

14. 与 MySQL 对比速查

特性PostgreSQLMySQL 8.0+
FULL OUTER JOIN❌(需 UNION 模拟)
CTE✅ 可写 CTE、MATERIALIZED 提示✅ 但不支持可写 CTE
递归 CTE✅ + CYCLE 检测(14+)
窗口函数✅ 全支持 + 自定义✅ 8.0 开始支持
LATERAL JOIN❌(类似能力靠 JSON_TABLE)
DISTINCT ON
RETURNING✅ 全 DML 支持❌(MariaDB 10.5 有)
ON CONFLICT✅ 灵活,可指定冲突列INSERT ... ON DUPLICATE KEY UPDATE,不可指定
GROUPING SETS/CUBE✅ 8.0 开始(仅 ROLLUP,无 CUBE)
INTERSECT / EXCEPT✅ 8.0+
TABLESAMPLE

15. 小结:一张思维导图带走

高级查询
├─ 多表拼接
│  ├─ JOIN 全家桶(INNER/LEFT/RIGHT/FULL/CROSS)
│  ├─ 自连接
│  └─ LATERAL JOIN(右可引用左)
├─ 嵌套查询
│  ├─ 子查询 4 种形状
│  └─ EXISTS vs IN vs JOIN
├─ 分步查询
│  ├─ CTE(可读性)
│  ├─ 可写 CTE(PG 特色)
│  └─ 递归 CTE(树、图)
├─ 结果集加工
│  ├─ 窗口函数(排名、偏移、移动)
│  ├─ GROUPING SETS / ROLLUP / CUBE
│  ├─ DISTINCT ON(PG 特色去重)
│  └─ 集合运算(UNION/INTERSECT/EXCEPT)
└─ 写入 + 查询组合
   ├─ RETURNING(PG 特色)
   ├─ ON CONFLICT upsert(PG 特色)
   └─ TABLESAMPLE 抽样

🎮 配套演示

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

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


16. 面试高频题(7 题)

Q1:EXISTSIN 到底哪个快?

考察点:子查询优化、半连接(Semi Join)、NULL 语义。

答案

  1. 在现代 PostgreSQL(11+)里两者几乎等价:优化器都会把 IN (subquery) 改写成 Semi Join,执行计划几乎一致。差异只在写法习惯。
  2. 在相关子查询里 EXISTS 更贴近语义:强调「存在就行」,优化器看到后会短路求值。
  3. NOT IN 有 NULL 陷阱:当子查询结果含 NULL 时,NOT IN 会把整条查询变空集,因为 x <> NULL 永远是 UNKNOWNNOT EXISTS 没有这个问题。
  4. 返回大量字段时 INJOIN 更干净JOIN + DISTINCT 为了去重会做 Sort/Hash,而 IN / EXISTS 走 Semi Join 不会有去重开销。
  5. 经验法则:查存在性用 EXISTS/IN,查「另一张表的字段」用 JOIN,不要为了风格统一硬改。

加分项:提到 MySQL 5.6 以前 IN (subquery) 会退化成 DependentSubquery(对外层每行执行一次),5.7 之后才优化,PG 从来没这个历史包袱。


Q2:窗口函数和 GROUP BY 有什么区别?什么时候该用窗口?

考察点:窗口函数语义、SQL 执行顺序。

答案

  1. 行数不同GROUP BY 把 N 行合并成一行;窗口函数保留 N 行,只是每行带一个聚合值。
  2. 执行位置不同GROUP BYHAVING 之前,窗口函数在 HAVING 之后、SELECT 之前。所以不能在 WHERE/HAVING 里用窗口函数,必须用子查询或 CTE 包一层。
  3. 能力不同:窗口可以「看一组里相邻的几行」(LAG/LEAD、移动平均),这是 GROUP BY 做不到的。
  4. 使用场景
    • 要计算排名、百分位、累计、同比环比 → 窗口。
    • 要最终只输出汇总值 → GROUP BY
    • 要在保留明细的同时带上汇总(比如每行显示「和所在部门平均的差」) → 窗口。
  5. 性能:窗口需要排序,若 PARTITION BY 列已有索引能复用;否则会触发 Sort 节点。GROUP BY 可以走 HashAggregate,一般更快。

Q3:PostgreSQL 的 ON CONFLICT 跟 MySQL 的 ON DUPLICATE KEY UPDATE 对比,有什么优势?

考察点:UPSERT 语义、并发控制、唯一索引设计。

答案

  1. 能指定冲突目标:PG 必须指明是 (col) 还是 ON CONSTRAINT name;MySQL 任何唯一键冲突都会触发同一段 UPDATE,表有多个唯一键时极易误伤。
  2. 支持 DO NOTHING:直接忽略冲突行,MySQL 没有原生等价(只能 INSERT IGNORE,但它连外键错误、非空错误都忽略,危险)。
  3. 支持部分索引冲突ON CONFLICT (email) WHERE deleted_at IS NULL 可以只对软删除未生效的唯一约束做 UPSERT,MySQL 没有部分索引。
  4. 冲突值引用更清晰:PG 用 EXCLUDED.col 指「试图插入的那一行」;MySQL 旧写法 VALUES(col) 被废弃,新写法用行别名 AS new
  5. 原子性 + 可组合性:PG 可 ON CONFLICT DO UPDATE ... RETURNING *,一条 SQL 完成插入 / 更新 / 返回新值;MySQL 没有 RETURNING
  6. 并发:两者都用行锁实现 UPSERT 无死锁,但 PG 的实现更透明(pg_try_advisory_xact_lock 等机制可补强)。

加分项:提到 PG 的 UPSERT 是 PG 9.5 引入,此前要写「先 INSERT ... WHERE NOT EXISTS,冲突再 UPDATE」的模式,或用 MERGE(PG 15+)。


Q4:递归 CTE 如何防止死循环?生产上用的时候要注意什么?

考察点:递归 CTE 执行模型、环检测。

答案

  1. 执行模型:递归 CTE 用「working table + result table」双表迭代。每轮把本轮产生的行放进 working table,下一轮拿它再 JOIN,直到某轮产出 0 行停止。
  2. 死循环来源:数据里出现环(A 的 parent 是 B,B 的 parent 是 A),或者 JOIN 条件写错导致每轮都产出新行。
  3. 防御
    • 深度上限WHERE lvl < 50,超过就停。
    • 路径去重:把已访问的 id 收集到一个数组,NOT (e.id = ANY(path_ids))
    • CYCLE 子句(PG 14+):语法上直接声明 CYCLE id SET is_cycle USING path,遇到环时打标记。
  4. 性能要点
    • 递归 CTE 不会用索引直接跳层,每轮都要 JOIN,复杂度和树的节点数成正比。
    • 大量数据的树展开,比业务代码在应用层做 N+1 查询要快得多,但不如物化的「闭包表」(closure table)。
    • 能用 ltree 扩展的场景(路径预存),查询效率更高。
  5. 生产经验:把递归 CTE 放在只读从库 / 报表库里跑,避免长事务影响主库 MVCC 回收。

Q5:什么场景下该用 LATERAL JOIN?举一个用窗口函数做不到的例子。

考察点:LATERAL 语义、与窗口函数的边界。

答案

  1. LATERAL 语义:让右侧的子查询 / 函数能引用左侧的列,等价于「对左表每一行执行一次右边」。
  2. 和窗口函数的区别:窗口函数能「在整张表排序完后取 TOP K」,但它要把全表都排;LATERAL 对每个分区独立执行子查询,可以命中 (user_id, created_at) 组合索引的 Top-N,在有索引时快得多。
  3. 窗口函数做不到的例子
    • 每个用户最近 3 笔订单 + 每笔订单的前两个商品(两层嵌套 TOP N),窗口函数难以嵌套;LATERAL 可以层层套。
    • 对每一行调用一个返回多行的函数SELECT u.*, h.score FROM ch5_users u, LATERAL my_func(u.id) AS h
    • 需要在子查询里动态构造参数LATERAL (SELECT ... FROM t WHERE t.created_at > u.signup_at + INTERVAL '30 days')
  4. 性能注意:LATERAL 默认走 Nested Loop,左表大 + 右表无索引时,要警惕 N×M 爆炸。建好 (user_id, created_at DESC) 这种索引是关键。
  5. 兼容性:LATERAL 是 SQL:1999 标准,PG、Oracle、SQL Server 都支持;MySQL 8 有 LATERAL DERIVED(别名写法不同)。

Q6:DISTINCTDISTINCT ON 有什么区别?

考察点:PG 特色语法、与其他数据库的可移植性。

答案

  1. DISTINCT:对 SELECT 列表的整行做去重,两行完全一致才算重复。
  2. DISTINCT ON (cols):按指定列分组,每组只保留 ORDER BY 排序后的第一行。
  3. 经典用法
    sql
    SELECT DISTINCT ON (user_id) user_id, id, created_at
    FROM ch5_orders
    ORDER BY user_id, created_at DESC;
    含义:取每个用户最新的一笔订单。
  4. 约束DISTINCT ON 的列必须是 ORDER BY前缀,否则哪一行被保留是未定义行为。
  5. 和窗口函数对比
    • 等价写法:ROW_NUMBER() OVER (PARTITION BY user_id ORDER BY created_at DESC) = 1,但多一层子查询。
    • DISTINCT ON 在小规模数据下可读性高,性能有时候更好(可以走 Skip Scan)。
  6. 可移植性DISTINCT ON 是 PG 独有语法,迁移到 MySQL/Oracle/SQL Server 都要改成窗口函数写法。
  7. 常见误区:以为 DISTINCT ON (a) b, c 是「只对 a 去重,bc 任取」——实际上「任取」由 ORDER BY 决定,没 ORDER BY 就是乱取。

Q7:窗口函数 ROWSRANGE 的区别?为什么 LAST_VALUE 经常返回「当前行值」?

考察点:窗口框架、默认值陷阱。

答案

  1. ROWS vs RANGE
    • ROWS BETWEEN N PRECEDING AND M FOLLOWING按物理行数,前 N 行 + 当前 + 后 M 行。
    • RANGE BETWEEN ...按排序列的值距离,相同值的行会被视为同一个「逻辑位置」。
  2. 等值行的处理:如果 ORDER BY day 里有多个同一天的行,ROWS 把它们算不同行;RANGE 会把同一天都纳入同一个边界。
  3. 默认框架的坑
    • ORDER BY 时,默认是 RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW,即「窗口起点到当前行(含同值行)」。
    • 所以 LAST_VALUE(x) OVER (ORDER BY y) 永远等于「当前行的 x」,因为「窗口里最后一行」就是当前行。
  4. 正确写法:要取分区里真正的最后一个值,必须显式写框架:
    sql
    LAST_VALUE(x) OVER (
      PARTITION BY g ORDER BY y
      ROWS BETWEEN UNBOUNDED PRECEDING AND UNBOUNDED FOLLOWING
    )
  5. 性能ROWS 框架优化器容易内存连续滑动;RANGE INTERVAL(PG 11+)虽然强大但可能退化成全分区扫描。建议数据量大时用 ROWS 表达。
  6. 实战口诀:写窗口函数时,如果用到 LAST_VALUE / NTH_VALUE / FIRST_VALUE,一定显式写窗口框架,否则掉坑是必然的。

本章示例对应的 init.sql、HTML 演示、5 个 Python 脚本均在同名目录 05_advanced_query/ 下。


🔗 延伸阅读

🎬 可视化演示

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

💻 示例代码

python
"""公共工具:psycopg v3 连接参数 + 友好打印。

依赖:pip install "psycopg[binary]"
所有脚本统一使用这里的 `conninfo()` / `connect()` 以便集中修改连接参数。
"""
from __future__ import annotations

import os
from contextlib import contextmanager
from typing import Iterable, Sequence

import psycopg


def conninfo() -> str:
    """返回连接串。支持环境变量覆盖(便于 CI / 容器环境)。"""
    host = os.getenv("PGHOST", "127.0.0.1")
    port = os.getenv("PGPORT", "5432")
    db   = os.getenv("PGDATABASE", "learn_pg")
    user = os.getenv("PGUSER", "postgres")
    pwd  = os.getenv("PGPASSWORD", "")
    parts = [f"host={host}", f"port={port}", f"dbname={db}", f"user={user}"]
    if pwd:
        parts.append(f"password={pwd}")
    return " ".join(parts)


@contextmanager
def connect():
    """with connect() as conn: ..."""
    with psycopg.connect(conninfo(), autocommit=False) as conn:
        yield conn


def print_table(rows: Sequence[Sequence], headers: Iterable[str]) -> None:
    """极简表格打印,避免额外依赖 tabulate。"""
    headers = list(headers)
    str_rows = [[("" if v is None else str(v)) for v in r] for r in rows]
    widths = [len(h) for h in headers]
    for r in str_rows:
        for i, v in enumerate(r):
            if i < len(widths):
                widths[i] = max(widths[i], len(v))
    line = "+" + "+".join("-" * (w + 2) for w in widths) + "+"
    print(line)
    print("| " + " | ".join(h.ljust(widths[i]) for i, h in enumerate(headers)) + " |")
    print(line)
    for r in str_rows:
        print("| " + " | ".join(r[i].ljust(widths[i]) for i in range(len(widths))) + " |")
    print(line)


def section(title: str) -> None:
    bar = "=" * 72
    print("\n" + bar)
    print(f"  {title}")
    print(bar)
python
"""01_join_subquery.py —— JOIN 与子查询性能对比

演示:
  1) 4 种 JOIN 结果差异
  2) EXISTS vs IN vs JOIN+DISTINCT 的执行计划对比
  3) NOT IN 的 NULL 陷阱

运行前请先执行 init.sql。
"""
from __future__ import annotations

from _common import connect, print_table, section


def demo_join_types(cur) -> None:
    section("1. 4 种 JOIN 行数对比(ch5_orders / ch5_users)")
    for jt in ("INNER", "LEFT", "RIGHT", "FULL"):
        cur.execute(
            f"""
            SELECT COUNT(*) FROM ch5_users u
            {jt} JOIN ch5_orders o ON u.id = o.user_id
            """
        )
        (n,) = cur.fetchone()
        print(f"  {jt:5s} JOIN -> {n:>6d} 行")

    cur.execute("SELECT COUNT(*) FROM ch5_users u CROSS JOIN ch5_orders o")
    (n,) = cur.fetchone()
    print(f"  CROSS JOIN -> {n:>6d} 行 (笛卡儿积)")


def demo_exists_vs_in(cur) -> None:
    section("2. EXISTS / IN / JOIN 执行计划对比")
    sqls = {
        "IN":     "SELECT COUNT(*) FROM ch5_users WHERE id IN (SELECT user_id FROM ch5_orders)",
        "EXISTS": """
            SELECT COUNT(*) FROM ch5_users u
            WHERE EXISTS (SELECT 1 FROM ch5_orders o WHERE o.user_id = u.id)
        """,
        "JOIN":   """
            SELECT COUNT(*) FROM (
              SELECT DISTINCT u.id FROM ch5_users u JOIN ch5_orders o ON o.user_id = u.id
            ) t
        """,
    }
    for name, sql in sqls.items():
        cur.execute("EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) " + sql)
        plan = cur.fetchone()[0][0]
        top = plan["Plan"]
        print(f"\n[{name}]  Total Cost={top['Total Cost']:.2f}  "
              f"Actual={top['Actual Total Time']:.3f}ms  "
              f"Rows={top['Actual Rows']}")


def demo_not_in_null_trap(cur) -> None:
    section("3. NOT IN 的 NULL 陷阱")
    cur.execute("CREATE TEMP TABLE tmp_u(id INT)")
    cur.execute("INSERT INTO tmp_u VALUES (1),(2),(3)")
    cur.execute("CREATE TEMP TABLE tmp_o(uid INT)")
    cur.execute("INSERT INTO tmp_o VALUES (1), (NULL)")

    cur.execute("SELECT id FROM tmp_u WHERE id NOT IN (SELECT uid FROM tmp_o)")
    rows_notin = cur.fetchall()

    cur.execute("""
        SELECT id FROM tmp_u u
        WHERE NOT EXISTS (SELECT 1 FROM tmp_o o WHERE o.uid = u.id)
    """)
    rows_notexists = cur.fetchall()

    print_table(
        [
            ["NOT IN   (遇到 NULL 全空)",  str(rows_notin)],
            ["NOT EXISTS (正确)",          str(rows_notexists)],
        ],
        ["写法", "结果"],
    )


def main() -> None:
    with connect() as conn:
        with conn.cursor() as cur:
            demo_join_types(cur)
            demo_exists_vs_in(cur)
            demo_not_in_null_trap(cur)


if __name__ == "__main__":
    main()
python
"""02_recursive_cte.py —— 员工组织架构递归展开

演示:
  1) 从 CEO 递归向下展开,生成 (name, lvl, path)
  2) 从任意员工向上找所有上级
  3) 用 CYCLE(PG 14+)防死循环(若版本不支持则跳过)
"""
from __future__ import annotations

from _common import connect, print_table, section


def expand_down(cur) -> None:
    section("1. 从 CEO 向下展开组织架构")
    cur.execute(
        """
        WITH RECURSIVE org AS (
          SELECT id, name, manager_id, 1 AS lvl,
                 name::TEXT AS path
          FROM ch5_employees
          WHERE manager_id IS NULL
          UNION ALL
          SELECT e.id, e.name, e.manager_id, o.lvl + 1,
                 o.path || ' > ' || e.name
          FROM ch5_employees e
          JOIN org o ON e.manager_id = o.id
        )
        SELECT lpad('', (lvl-1)*2, ' ') || name AS tree, lvl, path
        FROM org
        ORDER BY path
        """
    )
    print_table(cur.fetchall(), ["tree", "lvl", "path"])


def find_up(cur, start_name: str = "Alice") -> None:
    section(f"2. 从 {start_name} 向上查找所有上级")
    cur.execute(
        """
        WITH RECURSIVE mgr AS (
          SELECT id, name, manager_id, 1 AS lvl
          FROM ch5_employees WHERE name = %s
          UNION ALL
          SELECT e.id, e.name, e.manager_id, m.lvl + 1
          FROM ch5_employees e
          JOIN mgr m ON m.manager_id = e.id
        )
        SELECT lvl, name FROM mgr ORDER BY lvl
        """,
        (start_name,),
    )
    print_table(cur.fetchall(), ["距离", "姓名"])


def demo_cycle(cur) -> None:
    """PG 14+ 支持 CYCLE 子句,用于防止环导致死循环。"""
    section("3. CYCLE 子句防死循环(PG 14+)")
    cur.execute("SHOW server_version_num")
    (ver,) = cur.fetchone()
    if int(ver) < 140000:
        print("  当前 PG 版本 < 14,跳过 CYCLE 子句演示。")
        return

    cur.execute(
        """
        WITH RECURSIVE org AS (
          SELECT id, name, manager_id, 1 AS lvl
          FROM ch5_employees WHERE manager_id IS NULL
          UNION ALL
          SELECT e.id, e.name, e.manager_id, o.lvl + 1
          FROM ch5_employees e JOIN org o ON e.manager_id = o.id
        )
        CYCLE id SET is_cycle USING path
        SELECT name, lvl, is_cycle, path FROM org ORDER BY path LIMIT 5
        """
    )
    print_table(cur.fetchall(), ["name", "lvl", "is_cycle", "path"])


def main() -> None:
    with connect() as conn:
        with conn.cursor() as cur:
            expand_down(cur)
            find_up(cur, "Alice")
            demo_cycle(cur)


if __name__ == "__main__":
    main()
python
"""03_window_functions.py —— 窗口函数实战

演示:
  1) 每个类目的商品价格排名(ROW_NUMBER / RANK / DENSE_RANK / NTILE)
  2) 销售额的 3 个月移动平均 + 累计 + 同比
  3) 每个用户最近一笔订单(DISTINCT ON)
"""
from __future__ import annotations

from _common import connect, print_table, section


def rank_demo(cur) -> None:
    section("1. 每个类目的商品价格 TOP 3(4 种排名对比)")
    cur.execute(
        """
        WITH t AS (
          SELECT category, name, price,
                 ROW_NUMBER() OVER w AS rn,
                 RANK()       OVER w AS rk,
                 DENSE_RANK() OVER w AS drk,
                 NTILE(4)     OVER w AS bucket
          FROM ch5_products
          WINDOW w AS (PARTITION BY category ORDER BY price DESC)
        )
        SELECT category, name, price, rn, rk, drk, bucket
        FROM t WHERE rn <= 3
        ORDER BY category, rn
        """
    )
    print_table(
        cur.fetchall(),
        ["category", "name", "price", "ROW_NUMBER", "RANK", "DENSE_RANK", "NTILE(4)"],
    )


def moving_avg(cur) -> None:
    section("2. 月销售 3 个月移动平均 + 累计 + 同比")
    cur.execute(
        """
        SELECT
          month,
          revenue,
          ROUND(AVG(revenue) OVER (
            ORDER BY month
            ROWS BETWEEN 2 PRECEDING AND CURRENT ROW
          ), 2) AS ma3,
          SUM(revenue) OVER (ORDER BY month) AS running_total,
          LAG(revenue, 12) OVER (ORDER BY month) AS same_month_last_year,
          ROUND(
            ((revenue - LAG(revenue, 12) OVER (ORDER BY month))
              / NULLIF(LAG(revenue, 12) OVER (ORDER BY month), 0) * 100)::NUMERIC,
            2
          ) AS yoy_pct
        FROM ch5_monthly_sales
        ORDER BY month
        """
    )
    print_table(
        cur.fetchall(),
        ["month", "revenue", "MA3", "cum", "YoY基数", "YoY%"],
    )


def latest_order_per_user(cur) -> None:
    section("3. 每个用户最近一笔订单(DISTINCT ON)")
    cur.execute(
        """
        SELECT DISTINCT ON (user_id)
               user_id, id AS order_id, amount, status, created_at
        FROM ch5_orders
        ORDER BY user_id, created_at DESC
        LIMIT 10
        """
    )
    print_table(
        cur.fetchall(),
        ["user_id", "order_id", "amount", "status", "created_at"],
    )


def main() -> None:
    with connect() as conn:
        with conn.cursor() as cur:
            rank_demo(cur)
            moving_avg(cur)
            latest_order_per_user(cur)


if __name__ == "__main__":
    main()
python
"""04_upsert.py —— ON CONFLICT UPSERT 演示

演示:
  1) DO NOTHING:忽略冲突
  2) DO UPDATE + EXCLUDED:更新冲突行
  3) RETURNING 结合 UPSERT:一次拿到最新行
  4) 基于部分索引的冲突条件
"""
from __future__ import annotations

from _common import connect, print_table, section


DDL = """
DROP TABLE IF EXISTS ch5_up_users;
CREATE TABLE ch5_up_users (
  id         INT PRIMARY KEY,
  name       TEXT NOT NULL,
  email      TEXT NOT NULL,
  updated_at TIMESTAMPTZ NOT NULL DEFAULT NOW(),
  deleted_at TIMESTAMPTZ
);

-- 部分唯一索引:软删除用户的 email 允许重复
CREATE UNIQUE INDEX ch5_up_users_email_active
  ON ch5_up_users (email) WHERE deleted_at IS NULL;
"""


def demo(cur) -> None:
    section("准备表 ch5_up_users 与部分唯一索引")
    cur.execute(DDL)

    section("1. DO NOTHING:主键冲突则忽略")
    cur.execute(
        "INSERT INTO ch5_up_users(id, name, email) VALUES (1, 'Alice', 'a@x.com')"
        " ON CONFLICT (id) DO NOTHING RETURNING id"
    )
    print(f"  首次插入返回: {cur.fetchall()}")

    cur.execute(
        "INSERT INTO ch5_up_users(id, name, email) VALUES (1, 'Alice-dup', 'x')"
        " ON CONFLICT (id) DO NOTHING RETURNING id"
    )
    print(f"  再次插入冲突返回: {cur.fetchall()}  (空代表被忽略)")

    section("2. DO UPDATE + EXCLUDED:冲突则更新")
    cur.execute(
        """
        INSERT INTO ch5_up_users(id, name, email) VALUES (1, 'Alice-v2', 'a2@x.com')
        ON CONFLICT (id) DO UPDATE
        SET name = EXCLUDED.name,
            email = EXCLUDED.email,
            updated_at = NOW()
        RETURNING id, name, email, updated_at
        """
    )
    print_table(cur.fetchall(), ["id", "name", "email", "updated_at"])

    section("3. 基于部分唯一索引的冲突:软删用户 email 可重复")
    cur.execute(
        "INSERT INTO ch5_up_users(id, name, email, deleted_at) VALUES "
        "(2, 'Bob', 'b@x.com', NOW())"
    )
    cur.execute(
        """
        INSERT INTO ch5_up_users(id, name, email) VALUES (3, 'Bob-new', 'b@x.com')
        ON CONFLICT (email) WHERE deleted_at IS NULL DO UPDATE
        SET name = EXCLUDED.name
        RETURNING id, name, email, deleted_at
        """
    )
    print_table(cur.fetchall(), ["id", "name", "email", "deleted_at"])
    print("  说明:软删除的 Bob 不参与唯一约束,所以新 Bob 被成功 INSERT。")

    section("4. 批量 UPSERT:一次 INSERT N 行")
    rows = [(10 + i, f"user_{i}", f"u{i}@x.com") for i in range(5)]
    cur.executemany(
        "INSERT INTO ch5_up_users(id, name, email) VALUES (%s,%s,%s)"
        " ON CONFLICT (id) DO UPDATE SET name = EXCLUDED.name",
        rows,
    )
    cur.execute("SELECT id, name, email FROM ch5_up_users ORDER BY id")
    print_table(cur.fetchall(), ["id", "name", "email"])


def main() -> None:
    with connect() as conn:
        with conn.cursor() as cur:
            demo(cur)
            conn.commit()


if __name__ == "__main__":
    main()
python
"""05_lateral_join.py —— LATERAL JOIN 实战

场景:
  A) 每个用户最近 3 笔订单
  B) 用窗口函数做同样的事,对比执行计划
  C) LATERAL + 聚合:每个用户近 30 天消费总额
"""
from __future__ import annotations

from _common import connect, print_table, section


def ensure_index(cur) -> None:
    # (user_id, created_at DESC) 是 LATERAL Top-N 的最佳索引
    cur.execute(
        "CREATE INDEX IF NOT EXISTS idx_ch5_orders_user_time "
        "ON ch5_orders (user_id, created_at DESC)"
    )


def lateral_topN(cur) -> None:
    section("A. LATERAL:每个用户最近 3 笔订单")
    cur.execute(
        """
        SELECT u.id, u.name, o.id AS order_id, o.amount, o.created_at
        FROM ch5_users u
        LEFT JOIN LATERAL (
          SELECT id, amount, created_at
          FROM ch5_orders
          WHERE user_id = u.id
          ORDER BY created_at DESC
          LIMIT 3
        ) o ON true
        WHERE u.id <= 5
        ORDER BY u.id, o.created_at DESC NULLS LAST
        """
    )
    print_table(
        cur.fetchall(),
        ["uid", "name", "order_id", "amount", "created_at"],
    )


def window_topN(cur) -> None:
    section("B. 窗口函数版(等价但全表排序)")
    cur.execute(
        """
        WITH ranked AS (
          SELECT u.id AS uid, u.name, o.id AS order_id, o.amount, o.created_at,
                 ROW_NUMBER() OVER (PARTITION BY u.id ORDER BY o.created_at DESC) AS rn
          FROM ch5_users u
          LEFT JOIN ch5_orders o ON o.user_id = u.id
        )
        SELECT uid, name, order_id, amount, created_at
        FROM ranked
        WHERE rn <= 3 AND uid <= 5
        ORDER BY uid, created_at DESC NULLS LAST
        """
    )
    print_table(
        cur.fetchall(),
        ["uid", "name", "order_id", "amount", "created_at"],
    )


def compare_plans(cur) -> None:
    section("C. 两种写法的执行计划对比")
    sqls = {
        "LATERAL": """
            SELECT u.id, o.id FROM ch5_users u
            LEFT JOIN LATERAL (
              SELECT id FROM ch5_orders
              WHERE user_id = u.id
              ORDER BY created_at DESC LIMIT 3
            ) o ON true
        """,
        "Window":  """
            WITH r AS (
              SELECT u.id AS uid, o.id,
                     ROW_NUMBER() OVER (PARTITION BY u.id ORDER BY o.created_at DESC) rn
              FROM ch5_users u LEFT JOIN ch5_orders o ON o.user_id = u.id
            ) SELECT uid, id FROM r WHERE rn <= 3
        """,
    }
    for name, sql in sqls.items():
        cur.execute("EXPLAIN (ANALYZE, BUFFERS, FORMAT JSON) " + sql)
        p = cur.fetchone()[0][0]["Plan"]
        print(f"  {name:8s}  cost={p['Total Cost']:>9.2f}  "
              f"actual={p['Actual Total Time']:>8.3f}ms  "
              f"rows={p['Actual Rows']}")


def lateral_agg(cur) -> None:
    section("D. LATERAL + 聚合:每个用户近 30 天订单数 & 金额")
    cur.execute(
        """
        SELECT u.id, u.name, s.cnt, s.total
        FROM ch5_users u
        LEFT JOIN LATERAL (
          SELECT COUNT(*) AS cnt, COALESCE(SUM(amount), 0) AS total
          FROM ch5_orders
          WHERE user_id = u.id
            AND created_at >= NOW() - INTERVAL '30 days'
        ) s ON true
        WHERE u.id <= 10
        ORDER BY s.total DESC NULLS LAST
        """
    )
    print_table(cur.fetchall(), ["uid", "name", "近30天单数", "近30天金额"])


def main() -> None:
    with connect() as conn:
        with conn.cursor() as cur:
            ensure_index(cur)
            conn.commit()
            lateral_topN(cur)
            window_topN(cur)
            compare_plans(cur)
            lateral_agg(cur)


if __name__ == "__main__":
    main()
markdown
# 第 5 章 配套代码

> 演示 PG 高级查询:JOIN 全家桶、子查询、CTE / 递归 CTE、窗口函数、UPSERT、LATERAL JOIN。配套表统一以 `ch5_` 前缀,避免与其他章节冲突。

## 准备工作

1. 跑初始化脚本:
   ```bash
   psql -h 127.0.0.1 -U postgres -d learn_pg -f ../init.sql

会建出 7 张表(ch5_users / ch5_products / ch5_orders / ch5_order_items / ch5_employees / ch5_orders_archive / ch5_monthly_sales)并灌好 1000 行订单 / 24 个月销售 / 16 人组织树。 2. 安装依赖:

bash
pip install "psycopg[binary]>=3.1"
  1. (可选)通过环境变量覆盖默认连接:
    bash
    export PG_DSN="host=127.0.0.1 port=5432 dbname=learn_pg user=postgres"

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

脚本一句话说明关键 PG 特性
01_join_subquery.py4 种 JOIN + EXISTS / IN / JOIN 计划对比 + NOT IN NULL 陷阱LEFT/RIGHT/FULL/CROSS JOIN / 半连接 / NOT EXISTS
02_recursive_cte.py员工树向下展开 + 向上找上级 + CYCLE 防环(PG 14+)WITH RECURSIVE / CYCLE … SET … USING …
03_window_functions.py类目排名 + 移动平均 + 同比 + DISTINCT ON 取每用户最新一笔ROW_NUMBER/RANK/DENSE_RANK/NTILE / LAG / ROWS BETWEEN …
04_upsert.pyON CONFLICT DO NOTHING/UPDATE + 部分唯一索引 + 批量 UPSERTEXCLUDED / ON CONFLICT (col) WHERE … / 部分唯一索引
05_lateral_join.py每用户最近 N 笔订单(LATERAL vs 窗口函数)+ EXPLAIN 对比LEFT JOIN LATERAL / (user_id, created_at DESC) 索引

预期输出

05_lateral_join.py 跑完会看到 LATERAL 与窗口函数的执行计划差异:

A. LATERAL:每个用户最近 3 笔订单
  uid | name      | order_id | amount  | created_at
  ----+-----------+----------+---------+------------
   1  | user_1    |   902    | 1820.34 | 2026-04-15 …
   1  | user_1    |   771    | 1244.10 | 2026-04-12 …


C. 两种写法的执行计划对比
  LATERAL   cost=    35.42  actual=   2.184ms  rows=15
  Window    cost=   168.91  actual=  18.347ms  rows=15

常见报错

  • connection refused → PG 没起 / 端口不对
  • relation "ch5_orders" does not exist → 没跑 ../init.sql
  • column "rn" does not exist → 把窗口函数 ROW_NUMBER() OVER … 直接用在 WHERE 里了,须用子查询 / CTE 包一层
  • duplicate key value violates unique constraint "ch5_users_email_key" → 04_upsert.py 跑过半被中断后再跑,先 TRUNCATE ch5_up_users 或重跑 init.sql

_common.py ↗ · 01_join_subquery.py ↗ · 02_recursive_cte.py ↗ · 03_window_functions.py ↗ · 04_upsert.py ↗ · 05_lateral_join.py ↗ · README.md ↗