跳转到主要内容

17 合纵连横:分析加速与分布式选型

“数据越来越多,所以要上分布式”不是一个架构结论,只是一句尚未完成的 问题描述。

同样一条慢月报,可能分别来自:

统计信息失真
  -> 规划器选错路径
缺少合适索引
  -> 选择性查询扫描过多数据
work_mem 不足
  -> 排序或哈希落到临时文件
每次重算历史事实
  -> 缺少可接受新鲜度的汇总层
OLTP 与 OLAP 争用资源
  -> 缺少负载隔离
单节点资源确已越界
  -> 才可能需要横向拆分

若不先辨认瓶颈,把数据分到更多节点只会把一个可观测的本地问题变成网络、 路由、远端事务、再平衡和部分失败共同参与的问题。

本章坚持一条次序:

先定义服务目标,再证明单机边界;先减少无效工作,再隔离负载;只有明确 哪一种资源无法在单节点满足目标后,才比较分布式候选。

这并不是反对分布式。恰恰相反,只有把进入条件、分布键、数据局部性、 失败语义和撤退路线写清楚,分布式才是一项可评审的工程决策,而不是对增长 焦虑的技术性反射。

本章完成后

你应当能够:

  • 把“分析慢”改写为数据量、并发、P50/P95/P99、吞吐、新鲜度、正确性、 RPO/RTO 和成本目标;
  • 区分 CPU、存储 I/O、缓存、临时文件、锁等待、计划误差和远端传输瓶颈;
  • EXPLAIN (ANALYZE, BUFFERS)、系统统计和冻结工作负载建立单机证据;
  • 从计划中识别 Gather、parallel scan、partial/final aggregate 与实际 worker 数;
  • 解释“计划允许并行”与“执行时拿到 worker”为什么是两件事;
  • 用 covering B-tree、BRIN、物化汇总和批处理分别解决不同访问形状;
  • 解释 Index Only Scan 的 visibility map 前置条件,不把一次偶然计划当 稳定合同;
  • work_mem 的外排/内排反例说明为什么不能按单查询峰值做全局调参;
  • 区分 PostgreSQL 原生物化视图的完整刷新与应用维护的增量汇总;
  • 识别 OLTP/OLAP 共存时对 CPU、buffer、temp、WAL、vacuum 和副本延迟的 竞争;
  • 写出进入分布式评审的硬门槛,而不是只写“未来数据会增长”;
  • 选择候选分布键,计算数据倾斜,并审计跨分片查询、JOIN、事务与唯一性;
  • 解释 PostgreSQL HASH 分区 remainder 为什么不等于整数 % modulus
  • 比较 PostgreSQL 扩展、兼容数据库与专用 OLAP 时区分 SQL、类型、事务、 扩展、运维和故障兼容;
  • 用相同冻结输入、相同查询和相同失败条件比较候选;
  • 通过 postgres_fdw 计划分清过滤下推、聚合下推、协调端聚合与行传输;
  • 解释“数据同分片”为什么仍不能自动证明某条 JOIN 已被下推;
  • 定义单分片不可达时,单租户读、全局读、写入与重试分别应如何表现;
  • 在 Pigsty 中把分析读隔离到 offline replica,或声明一个待验收的 Citus 拓扑,同时不把配置片段当生产验收;
  • 输出一份包含证据、限制、生产代价、复审触发器和退出路线的 ADR。

贯穿本章的销售分析

实验生成一份完全确定的合成数据:

8 tenants
50 accounts per tenant
120 days
5 sales per account per day

400 accounts
240,000 sales
1,200,000 units
2,256,000.00 amount

同一份业务月报由四条路径计算:

local raw facts
local daily materialized summary
partitioned postgres_fdw parent
two-stage remote daily + coordinator monthly aggregation

四条路径都必须逐字节等于 frozen-monthly.csv 的 32 行。正确性不一致 时,不允许继续比较计划或性能。

冻结事实:

项目 数值
本地事实 240,000 行
shard A / shard B 120,000 / 120,000 行
本地日汇总 2,880 行
最终月报 32 行
租户 3 的 4 月 7,500 笔,69,375.00
朴素 FDW 返回协调端 240,000 行
两阶段聚合返回协调端 960 行
业务校验和 42fb8ab5444469eba1f104a8e1e529dd
月报校验和 644d45544ebbc2a80c42270c38ac6885

这里的“返回行数”描述数据流形状,不是网络字节,也不是耗时。三个数据库都在 一台机器的一个 PostgreSQL 18.6 实例中,没有独立 CPU、磁盘、网络或故障域。

先看单机还能做什么

冻结计划证明四件不同的事。

月聚合可并行:

Finalize HashAggregate
  -> Gather
       Workers Planned: 2
       Workers Launched: 2
       -> Partial HashAggregate
            -> Parallel Seq Scan on sales_fact

租户 3 的选择性查询可走 covering index:

Index Only Scan using sales_fact_tenant_day_idx
  actual rows=7500
  Heap Fetches: 0

Heap Fetches: 0 并不是 INCLUDE 自动保证的。实验重建后显式执行 VACUUM (ANALYZE),让 visibility map 建立 all-visible 信息,再验证零回表。 如果跳过这一步,新表可能合理地使用 Bitmap Heap Scan。

同一排序在两个会话级配置下呈现不同资源路径:

work_mem=64kB -> external merge, Disk ~= 5.9MB
work_mem=32MB -> quicksort, Memory ~= 13.6MB

这只说明 spill 可被计划证据观察。一个查询需要 32MB,不等于应该把全局 work_mem 设置成 32MB;一个并发查询可以包含多个 sort/hash 节点,还有 并行 worker 和并发会话共同放大内存。

物化日汇总把月报输入从 240,000 行降为 2,880 行,但随之引入:

freshness target
refresh schedule
refresh failure recovery
late-arriving correction
locking and WAL cost
definition version

PostgreSQL 的物化视图持久保存查询结果,读取时像表;数据不会自动保持最新, 需要 REFRESH MATERIALIZED VIEW。官方 Materialized Views 把“读得更快”与“可能不新鲜”明确放在同一项权衡里。

再看分布式改变了什么

实验的协调端由 LIST 分区父表接管两个外表:

sales_fact_distributed PARTITION BY LIST (tenant_id)
├── sales_fact_dist_0: tenants 2,4,6,8 -> pg36_shard_a
└── sales_fact_dist_1: tenants 1,3,5,7 -> pg36_shard_b

租户 3 的查询只访问 shard B,计划中的远端 SQL 带上租户和日期:

Foreign Scan on sales_fact_dist_1
Remote SQL:
  SELECT amount
  FROM shop_ch17_shard.sales_fact
  WHERE occurred_on >= '2026-04-01'
    AND tenant_id = 3

全局月报若直接从分区父表聚合,两个 Foreign Scan 各返回 120,000 行, 协调端接收 240,000 条事实后聚合。改成每个远端先按租户、日期聚合:

shard A: 120,000 facts -> 480 daily aggregates
shard B: 120,000 facts -> 480 daily aggregates
coordinator: 960 daily aggregates -> 32 monthly rows

结果相同,传输形状完全不同。这是分布式查询最重要的思维之一:

尽量让过滤、连接和聚合靠近数据发生;但必须用实际计划证明下推,不可从 SQL 外观或拓扑图推断。

反例也被固定下来:租户 3 的账户与销售位于同一分片,查询通过两个分区外表 父表连接时,实测仍在协调端执行 Hash Join,接收 7,500 条销售与 50 条 账户。postgres_fdw 的远端优化、代价、fetch_size、连接与事务管理以 PostgreSQL 官方 postgres_fdw 文档为准。

HASH 分区不是整数取模

本章第一版失败原型使用:

remote fixture routing = tenant_id % 2
coordinator routing = PARTITION BY HASH (tenant_id)

它们不是同一算法。PostgreSQL HASH 分区先使用数据类型的哈希支持函数,再按 MODULUS/REMAINDER 判断分区;REMAINDER 0 不表示“偶数值”。当查询带 tenant_id 时,协调端会按自己的哈希算法裁剪到一个分区,而目标租户可能被 生成器放在另一个数据库,于是出现“全表看似有数据,按租户裁剪却静默漏数” 的危险结果。

冻结实验改用显式 LIST 路由,使物理分片和协调端边界完全一致。生产系统不应 手写八个租户清单,而应让同一个经过版本化的路由算法或分片元数据成为写入、 读取、再平衡和恢复的共同事实来源。

PostgreSQL 官方 Table Partitioning 说明 HASH 分区以 modulus/remainder 描述分区边界;不能把这些名词误读为对 原始整数直接做 %

Pigsty 中的两条候选路径

本章不把 Pigsty 等同于某一种分布式数据库。它首先提供一种声明和交付运行 环境的方法。

路径 A 是保留一个 PostgreSQL 数据体系,把 OLAP/ETL/交互慢查询隔离到 offline replica 或带 pg_offline_query 标签的副本:

all:
  children:
    pg-analytics:
      hosts:
        10.10.10.11: { pg_seq: 1, pg_role: primary }
        10.10.10.12:
          pg_seq: 2
          pg_role: replica
          pg_offline_query: true
      vars:
        pg_cluster: pg-analytics
        pg_conf: olap.yml

路径 B 是在硬门槛满足后评估 Citus。Pigsty 4.5 文档要求 Citus 拓扑声明 pg_mode: cituspg_shard、各分片的 pg_grouppg_primary_db, 并配置数据节点间访问规则;完整生产设计还必须补齐 coordinator/worker HA、 服务路由、备份恢复、再平衡、监控和升级。

Pigsty 的 配置入口集群/实例类型 给出了 offline 与 Citus 的当前声明方式。本章资产 pigsty-declaration.example.yml 只是两个互斥候选的草图,未执行 L1,不能直接合并进生产 inventory。

实验资产

规范与决策:

生成与建立:

结果与计划:

审计与退出:

快速运行

本地开发数据库先完成第 4 章的角色与物理模型,然后提供管理员 service:

export PGSERVICEFILE=/path/to/pg_service.conf
export PGSERVICE=pg36-admin

PG36_EVIDENCE_DIR="$PWD/evidence/ch17" \
  ./static/labs/ch17/task.sh all

all 会执行两个完整周期:

bootstrap retained database shells
  -> rebuild shard A and B
  -> rebuild coordinator
  -> export and compare four monthly paths
  -> collect local and distributed plans
  -> prove application write denial
  -> make shard B temporarily unreachable
  -> prove shard A scoped read still works
  -> prove global read fails with 08001
  -> review all evidence
  -> prove reset token/target/active-worker guards
  -> pre-verify all three databases
  -> exact per-database reset
  -> rebuild everything
  -> repeat evidence and review

正式实测输出:

status=ok
fixture=frozen-byte-identical-four-paths
single_node=parallel+index+summary+spill
distributed=tenant-pruning+fdw+two-stage
counterexamples=hash-is-not-modulo+join-not-pushed
failure=healthy-shard-read+global-08001
guards=P3660+P3661+P3663
postgres_fdw=1.2
pigsty_l1=not-run
release_candidate_checksum=3dcb7308cf6983122ee860ad3dc2a4b44651549e3d5631770839bb9a0be450c6

破坏边界

task.sh all 会在精确身份、marker、对象清单、权限和数据校验和匹配后, 删除并重建 shop_ch17shop_ch17_ext、两个 foreign server、六个 user mapping,以及两个分片数据库中的 shop_ch17_shard。它保留空的 pg36_shard_apg36_shard_b 数据库壳。只可用于本书受控开发 fixture, 不得在生产执行。

本章目录

17.1 先证明单机边界

17.2 单机分析能力

17.3 何时需要分布式

17.4 比较分布式候选

17.5 部署最小分布式 PoC

17.6 实战:从单机证据到选型 ADR


上一章:经天纬地:时序、空间与时空查询 · 返回上卷导读 · 下一章:万法归宗:PostgreSQL 数据平台与替代边界 · 查看全书目录 · 查看索引中心

17.1 先证明单机边界

“单机扛不住”必须是证据结论,而不是架构会议里的气氛。

最常见的误判有两类:

局部问题被说成容量问题
  一条坏 SQL / 一个缺失索引 / 一次统计失真
  -> “PostgreSQL 不适合分析”

容量问题被说成局部问题
  工作集、写入、维护窗口或故障域已越过单节点
  -> “再调一个参数就好”

本节不预设答案。先把目标、负载和瓶颈拆成可测量的对象,再决定应该优化、 隔离、扩容,还是分布。

17.1.1 定义数据量、并发、延迟与新鲜度目标

“数据量”至少有六种尺寸

只报“十亿行”几乎没有决策价值。十亿个窄整数与十亿个宽 JSONB 的存储、 缓存和扫描成本不同;十亿行均匀访问与 99% 查询只看最近一天也不同。

基线至少记录:

维度 示例问题 可验证证据
逻辑规模 行数、租户数、时间跨度? count(*)、业务目录
物理规模 heap、TOAST、索引各多大? pg_relation_sizepg_total_relation_size
工作集 查询真正反复访问哪部分? 计划 buffers、时间谓词、缓存命中
增长 每日新增、更新、删除多少? 时序采样与容量预测
倾斜 最大租户/日期/键占多少? percentile、top-N、直方图
生命周期 热、温、冷数据如何变化? 保留、归档与访问统计

本章 fixture 的身份不是一句“24 万行”:

8 tenants
400 accounts
2026-01-01 .. 2026-04-30
240,000 sales
2,000 heap pages in the verified run
local facts + daily summary + two remote shard copies

它还有明确限制:数据均匀、确定、合成,无法代表真实倾斜、缓存冷启动、 网络、WAL、vacuum 或生产并发。

并发不是 QPS 的同义词

分析系统常见四种并发:

arrival concurrency
  同时到达多少请求

active database concurrency
  同时在 PostgreSQL 内执行多少语句

in-query parallelism
  一条语句使用多少 parallel workers

background concurrency
  autovacuum、checkpoint、备份、复制、ETL、刷新同时做什么

一个 dashboard 打开时发出 30 条 SQL,不等于数据库应该同时运行 30 个重 聚合。连接池可以排队,应用可以合并请求,汇总层可以复用结果。反过来, “线上只有 20 个连接”也不表示压力小:每条查询可能启动多个 worker、多个 sort/hash 节点并产生大临时文件。

因此基线应同时记录:

request rate
queue time
active sessions
parallel workers planned/launched
statements per request
rows scanned / returned
temporary bytes
CPU and I/O saturation

延迟要有分位数和查询类别

“平均 800ms”会掩盖两类事实:

  • 99% 查询 10ms,1% 查询 80s;
  • 所有查询稳定在 800ms。

两者的容量和用户体验完全不同。至少按 workload class 报告:

类别 典型目标
单租户交互明细 P50/P95/P99 与超时率
dashboard 聚合 首屏、完整加载与刷新周期
批量报表 完成窗口与失败重跑时间
数据导出 吞吐、并发上限与资源封顶
ETL/刷新 截止时间、WAL/lag 与恢复点

不要把一次 EXPLAIN ANALYZE 的执行时间直接当 SLO。它只是一条 SQL 在某个 缓存、数据、参数和系统负载下的一次观察。SLO 需要在代表性并发、冷暖缓存和 运行周期下统计。

新鲜度独立于查询速度

分析请求常把两个目标混成一个:

query latency: 用户发出查询后多久返回
data freshness: 返回的数据距离真实业务现在有多旧

一个物化汇总可以在 20ms 返回昨天的数据;一条扫描原表的查询可以在 2s 返回刚提交的数据。谁更好取决于合同,而不是毫秒数。

新鲜度目标应写成可验证形式:

event-time freshness <= 5 minutes at P99
daily financial close complete by 02:00 UTC
late events within 24 hours must be included in next rebuild
dashboard may lag primary commit by 60 seconds

若使用副本,还要区分:

source event lag
ingestion lag
replication replay lag
summary refresh lag
cache lag

只看其中一个指标会把旧数据误报成“查询很快”。

正确性是第一项 SLO

所有候选必须在相同输入下得到相同业务结果。本章冻结 32 行月报,并同时固定:

business checksum = 42fb8ab5444469eba1f104a8e1e529dd
monthly checksum  = 644d45544ebbc2a80c42270c38ac6885
CSV SHA-256       = 64b045809e10364fd84a587121d919e8562a15335c4c6c015e91a0ead3a44323

四条计算路径逐字节比较:

cmp frozen-monthly.csv monthly-local.csv
cmp frozen-monthly.csv monthly-summary.csv
cmp frozen-monthly.csv monthly-distributed.csv
cmp frozen-monthly.csv monthly-two-stage.csv

如果某个候选“快很多”但少一个租户,它不是优化,而是错误。

用目标表替代形容词

一个可评审的初始目标可以长这样:

指标 当前 目标 测量条件
单租户明细 P95 1.8s < 500ms 50 并发、30 日窗口
全局月报完成时间 24min < 10min 冷缓存、完整月
dashboard 新鲜度 P99 12min < 5min 按事件时间
temp write/小时 800GB < 100GB 正常峰值
primary CPU P95 92% < 70% OLTP+分析同时
replica replay lag P99 9min < 60s 报表窗口

当前值未知时写 unknown,随后安排测量。不要用“应该没问题”填表。

工作负载清单

选型前收集每类查询:

SQL fingerprint
business owner
read/write
frequency and concurrency
parameters and selectivity
rows scanned / returned
latency distribution
temporary I/O
lock behavior
freshness requirement
retry/idempotency behavior
failure consequence

同时固定 schema、统计信息、参数、数据生成方式和版本。否则两次跑分比较的 可能不是同一个系统。

本章的 fixture-manifest.json 保存生成器、行数、分片、校验和与限制;生产基线还应保存脱敏 workload manifest 和运行环境 manifest。

17.1.2 区分 CPU、I/O、内存、锁与计划瓶颈

先问“时间花在哪里”

慢查询的第一层分类:

waiting
  lock / I/O / client / WAL / remote / worker

running
  CPU expression / decompression / hash / sort / aggregation

planned badly
  row estimate / join order / access path / partition pruning

doing too much work
  wrong grain / no predicate / repeated calculation / data transfer

分类不是互斥的。错误估算可能选择大量随机 I/O;内存不足可能产生 temp I/O; 锁等待可能让 CPU 很空但延迟很高。

用执行计划建立因果链

推荐从:

EXPLAIN (
  ANALYZE,
  BUFFERS,
  WAL,
  SETTINGS,
  VERBOSE
)
SELECT ...;

开始,但要理解风险:

  • ANALYZE 会真正执行语句;
  • 对写语句使用时会真的修改数据,除非放在可回滚且外部副作用可控的事务中;
  • BUFFERS 展示 PostgreSQL buffer/I/O 计数,不等于操作系统层面的完整因果;
  • 一次计划不是延迟分布;
  • planner estimate 与 actual 的差距比节点名字本身更重要。

PostgreSQL 官方 Using EXPLAIN 解释 plan tree、cost、actual rows、loops、buffers 与不同节点的读法。

先检查:

actual rows × loops
estimated rows versus actual rows
rows removed by filter
heap fetches
sort method / memory / disk
hash batches
shared/local/temp buffers
workers planned / launched
partition subplans actually visited
remote SQL and returned rows

CPU 瓶颈

常见信号:

  • runnable CPU 长期接近可用核心上限;
  • 查询主要读取 cached buffers,物理 I/O 不高;
  • 大量表达式、JSON、正则、排序、哈希、聚合或 JIT 消耗;
  • 增加并发只增加排队,吞吐不再提高;
  • parallel worker 增加后单查询变快、系统总吞吐却下降。

CPU 证据必须区分:

database process CPU
kernel CPU
steal/throttling
per-query CPU
background maintenance CPU

不能从 PostgreSQL Execution Time 单独推导 CPU 时间。

可尝试的方向:

  • 减少扫描和返回行;
  • 改善连接顺序与聚合粒度;
  • 避免对每行重复做昂贵表达式;
  • 使用预计算/物化;
  • 审计并行度与并发;
  • 扩大单机 CPU;
  • 只有工作可被安全分片时再横向扩 CPU。

I/O 瓶颈

常见信号:

  • cache miss 后读取延迟高;
  • shared read 与系统块设备队列共同上升;
  • 顺序大扫描把 OLTP 热页挤出缓存;
  • temp read/write 大量增长;
  • checkpoint、backup、vacuum 与分析抢同一存储;
  • 增加 CPU 不改善吞吐。

要区分三类 I/O:

base relation/index I/O
temporary spill I/O
WAL/checkpoint/backup/replication I/O

它们的修复不同。缺索引与低选择性扫描不是同一问题;给全表聚合增加 B-tree 也未必比顺序扫描好。

内存与 spill

本章用同一排序证明:

SET work_mem = '64kB';
-- external merge, temp read/write

SET work_mem = '32MB';
-- quicksort in memory

计划来自:

冻结 24 万行上观察到:

64kB: external merge, Disk about 5920kB
32MB: quicksort, Memory about 13645kB

不要据此设置:

work_mem = 32MB

然后乘上几百连接。work_mem 是许多执行节点各自可以使用的预算,不是整个 查询或实例的硬上限;并行查询还会放大消费者。正确步骤是:

  1. 找到真实 spill 的 SQL 和节点;
  2. 判断能否通过索引、过滤、聚合顺序减少数据;
  3. 估算峰值并发 × 每查询节点 × worker;
  4. 优先用角色、数据库、会话或任务级设置;
  5. 同时设置超时、并发和 temp_file_limit 一类护栏;
  6. 用压力回放核对实例 RSS、OOM 与总吞吐。

PostgreSQL 的 Resource Consumptionwork_memhash_mem_multiplier、maintenance memory 与 huge pages 等 参数的版本基准。

锁瓶颈

分析查询通常只读,不等于不会造成并发问题:

  • 长事务延长 snapshot 生命周期,阻碍 vacuum 清理;
  • DDL 等待或被 ACCESS SHARE 阻塞;
  • REFRESH MATERIALIZED VIEW 的锁行为影响读者;
  • 报表函数可能隐含写临时/业务表;
  • 导出事务可能持有 snapshot 很久;
  • standby 上长查询可能与 WAL replay 冲突。

诊断要同时看:

SELECT
  pid,
  wait_event_type,
  wait_event,
  xact_start,
  query_start,
  state,
  application_name
FROM pg_catalog.pg_stat_activity
WHERE datname = current_database();

以及 blocking graph,而不是只数连接。第 10、12 章的事务、锁与慢查询诊断 方法在这里继续适用。

计划瓶颈

错误计划常见来源:

stale or insufficient statistics
correlated columns not represented
parameter-sensitive selectivity
implicit casts/collations
function-wrapped predicates
partition key not exposed
wrong join cardinality
generic plan versus custom plan
foreign table statistics drift

本章对外表执行 ANALYZE。官方 postgres_fdw 文档指出:本地统计可以减少 远端估算开销,但远端频繁变化时会很快过期;use_remote_estimate 则会增加 远端 planning 往返。两者都不是无条件更好。

“做太多工作”比节点选择更根本

原始月报与日汇总都得到 32 行:

raw plan:
  Parallel Seq Scan on sales_fact
  240,000 facts contribute

summary plan:
  Seq Scan on daily_tenant_summary
  2,880 summaries contribute

即使原始扫描计划完全正确,它仍在重复计算已经稳定的历史粒度。若业务允许 分钟或日级新鲜度,汇总可能比继续微调原表扫描更有效。

同理,分布式计划若把 240,000 行传到协调端再聚合,远端每个 Seq Scan 都 可能是“正确计划”,整体数据流却仍不合理。

一张瓶颈—证据—动作表

怀疑 至少需要的证据 优先动作
CPU CPU 饱和、每查询 CPU、计划工作量 少做工作、审计并行与表达式
base I/O buffer/read、设备延迟、访问形状 索引、裁剪、缓存/存储
temp I/O sort/hash 方法、temp bytes 减少输入、局部内存与并发
blocker、wait event、事务年龄 缩短事务、调度/锁语义
计划 estimate/actual、统计、参数 统计、SQL、索引、版本基线
重复计算 相同历史范围反复聚合 汇总、缓存、批处理
远端传输 Remote SQL、返回行、网络 下推、局部聚合、分布键

17.1.3 单机未被正确使用前不急于分布式

“单机优先”是一条证据顺序

合理的升级阶梯:

1. 业务口径与 SQL 正确
2. 统计、索引、分区裁剪正确
3. 内存与并行在并发预算内
4. 重复分析有汇总/批处理
5. OLTP 与 OLAP 有资源隔离
6. 单节点纵向容量仍不足
7. 分布键与主要查询天然对齐
8. 团队能承担分布式运维
9. 才进入横向分布

这不是要求永远把单机压到 100%。生产需要安全余量、维护窗口和故障容忍。 “正确使用”是达到经过评审的安全上限,而不是让事故替你找到极限。

先拒绝伪瓶颈

一个值得写进 ADR 的反例:

症状:
  租户 3 的 4 月明细聚合慢

错误推断:
  表有 24 万行,因此需要分片

证据:
  合适 covering index 后只读 7,500 个索引项
  Index Only Scan
  Heap Fetches: 0

结论:
  当前问题是访问路径,不是节点容量

索引定义:

CREATE INDEX sales_fact_tenant_day_idx
ON shop_ch17.sales_fact (
  tenant_id,
  occurred_on,
  account_id
)
INCLUDE (amount, units, channel);

fixture 重建结束后显式:

VACUUM (ANALYZE) shop_ch17.sales_fact;

这是计划合同的一部分。刚装载的 heap 尚未有足够 all-visible 位时,PostgreSQL 可能选择 Bitmap Heap Scan;不能把之前一次 autovacuum 留下的状态当可重复 前置条件。

BRIN 是相关性工具,不是“更小的 B-tree”

本章还创建:

CREATE INDEX sales_fact_day_brin_idx
ON shop_ch17.sales_fact
USING brin (occurred_on)
WITH (pages_per_range = 16);

BRIN 对“列值与物理位置天然相关”的大表按 block range 保存摘要,索引很小, 但返回候选 page range 后仍需 recheck,是 lossy 路径。它适合追加顺序与时间 大体一致的巨大事实表,不适合替代每种选择性 B-tree。

冻结小表只验证目录中存在 date_minmax_ops,并观察 BRIN 比 covering B-tree 小;不宣称这个查询上 BRIN 更快。官方 BRIN Indexes 说明 block range、物理相关性、lossy recheck、pages_per_range 与 summarization 行为。

单机边界应是曲线,不是一个点

容量实验应逐级增加:

data scale
concurrency
query mix
ingest rate
background maintenance
cache state

记录:

throughput
P50/P95/P99
queueing
CPU
read/write IOPS and latency
temp bytes
WAL
checkpoint
vacuum debt
replica lag
error/timeout rate

理想结果是一组曲线:

低并发:延迟稳定,吞吐线性增长
接近饱和:排队上升,吞吐增幅变小
过载:延迟和错误率急升,吞吐可能下降

生产容量线应位于拐点之前,并包含节点故障、维护和增长余量。

什么时候单机证据足以支持“继续单机”

可以暂缓分布式,当:

  • 调优后 SLO 在峰值与故障演练下满足;
  • 未来容量预测仍位于安全余量内;
  • 物化/批处理的新鲜度合同可接受;
  • offline replica 能隔离读负载;
  • 主要风险是可通过纵向扩容或存储升级解决;
  • 业务需要大量跨实体事务与灵活 JOIN,分片会显著破坏局部性;
  • 团队尚未具备分片备份、恢复、再平衡与值班能力。

什么时候不能再用“继续调优”拖延

应正式进入分布式评审,当代表性证据显示:

  • 单节点 CPU、内存、存储容量或 I/O 已越过安全上限;
  • 维护、vacuum、备份或恢复无法在窗口内完成;
  • 即使隔离到副本,分析吞吐仍受单节点资源限制;
  • 业务故障域或地域要求不能由一个集群满足;
  • 主要访问天然按租户/实体局部化,跨分片比例可控;
  • 硬件纵向升级的边际成本和上限不再可接受;
  • 团队已经定义跨分片事务、部分失败、重平衡和退出流程。

本节的停止条件

在以下问题没有答案前,不进入“选哪个分布式产品”:

目标是什么?
当前瓶颈是哪一种资源?
哪条 SQL、哪个粒度、哪类并发造成?
单机优化后曲线在哪里拐弯?
未来多久越过安全容量?
哪些查询可以按一个分布键局部化?
哪些事务一定跨边界?
如果一个节点不可用,业务允许什么结果?

下一节先把 PostgreSQL 单节点内部可用的并行、索引、BRIN、分区、物化和 负载隔离工具讲透,再讨论真正的分布式门槛。


返回本章目录 · 下一节:单机分析能力 · 查看全书目录 · 查看索引中心

17.2 单机分析能力

PostgreSQL 的“单机”不是“单进程、单线程、每次从原表重算”。

在引入分布式之前,至少有五个正交杠杆:

减少访问的数据       -> 选择性索引、分区裁剪
并行处理必要的数据   -> parallel scan/join/aggregate
缩小每次处理的粒度   -> 物化汇总、批处理
利用物理相关性       -> BRIN、聚簇/装载顺序
隔离不同负载         -> 会话护栏、连接池、offline replica

每个杠杆解决不同问题。把它们都叫“性能优化”会丢失决策边界。

17.2.1 并行扫描、连接、聚合与限制

并行计划的基本结构

PostgreSQL 在计划树中使用 GatherGather Merge 汇集 worker 的结果:

leader
  Gather / Gather Merge
    worker 0 -> parallel-aware subtree
    worker 1 -> parallel-aware subtree
    ...

Gather 不保留 worker 输出顺序;Gather Merge 合并已经排序的并行流。 Gather 下面并非每个节点都自动并行。只有 parallel-aware 的 scan、join、 aggregate 等节点能让 workers 分担输入;普通节点可能在每个 worker 内分别 执行,也可能只在 leader 上执行。

PostgreSQL 官方 Parallel Query 把并行扫描、连接、聚合、append 与 parallel safety 分开说明。读计划时应 沿 plan tree 判断“谁分担数据、谁合并结果”,而不是只搜索一个 Gather

本章的并行聚合

local-parallel-plan.sql 为冻结查询设置 一个可重复的实验上下文:

SET max_parallel_workers_per_gather = 2;
SET min_parallel_table_scan_size = 0;
SET parallel_setup_cost = 0;
SET parallel_tuple_cost = 0;

EXPLAIN (
  ANALYZE,
  BUFFERS,
  COSTS OFF,
  SUMMARY OFF,
  TIMING OFF
)
SELECT
  tenant_id,
  date_trunc('month', occurred_on::timestamp)::date
    AS month_start,
  count(*) AS sale_count,
  sum(units)::bigint AS unit_count,
  sum(amount)::numeric(18,2) AS amount_total
FROM shop_ch17.sales_fact
GROUP BY tenant_id, month_start
ORDER BY tenant_id, month_start;

冻结计划:

Sort (actual rows=32 loops=1)
  -> Finalize HashAggregate (actual rows=32 loops=1)
       -> Gather (actual rows=96 loops=1)
            Workers Planned: 2
            Workers Launched: 2
            -> Partial HashAggregate (actual rows=32 loops=3)
                 -> Parallel Seq Scan on sales_fact
                      actual rows=80000 loops=3

读法:

  1. leader 与两个 worker 合计三个参与者;
  2. 每个参与者扫描约 80,000 行;
  3. 每个参与者产出 32 个 partial groups;
  4. Gather 收到约 96 行;
  5. finalize aggregate 合并成 32 行;
  6. 最后按租户和月份排序。

这比“三个人一起扫 24 万行”更精确:并行收益来自把大量输入压成少量 partial state,再让 leader 合并。若每个 worker 都输出海量行,leader 可能成为瓶颈。

partial/final aggregate 的条件

聚合要能并行拆分,必须有可合并的中间状态。概念上:

worker partial state
  count = 100
  sum   = 935.50

another worker partial state
  count = 120
  sum   = 1101.20

combine/final
  count = 220
  sum   = 2036.70

某些聚合、表达式、函数或语义无法安全拆分,就不会出现 partial/final aggregate。用户自定义函数默认不是 parallel safe;必须由作者基于真实行为 正确标记,不能为了得到并行计划而随意改 catalog。

计划能并行,不表示执行一定并行

计划显示:

Workers Planned: 2

执行证据还要看:

Workers Launched: 2

可用 worker 受多个上限和当前占用影响,例如:

max_worker_processes
max_parallel_workers
max_parallel_workers_per_gather
other sessions already using workers

如果执行时拿不到 worker,leader 可能独自执行 Gather 以下部分。因此容量 测试必须在代表性并发下观察 launched,而不是从单会话计划推断。

PostgreSQL 官方 When Can Parallel Query Be Used? 还列出写入、行锁、cursor、parallel-unsafe function、嵌套并行和 worker 资源不足等限制。

并行扫描

常见 parallel-aware 扫描包括:

Parallel Seq Scan
Parallel Index Scan
Parallel Index Only Scan
Parallel Bitmap Heap Scan

它们适合的访问形状不同:

  • 大范围低选择性读取常适合 parallel seq scan;
  • 有序 B-tree 与查询方向匹配时可并行 index scan;
  • visibility map 允许时 index-only 可减少 heap 访问;
  • bitmap 路径适合聚合多个索引命中后批量访问 heap page。

不能把 Parallel Seq Scan 当成“没用索引所以坏”。24 万行几乎全参与月聚合, 顺序读并行处理可能正是正确路径。判断标准是选择性、缓存、物理布局、并发与 总体资源,而不是节点名字的好恶。

并行连接

并行连接可能让:

outer side produced in parallel
inner side shared or rebuilt per worker
join work divided among workers

不同 join 算法的资源行为不同:

  • nested loop 的 inner scan 可能在每个 worker 重复;
  • merge join 的 inner side 可能被多次执行;
  • parallel hash 可以共享 hash table;
  • skew、错误基数和 worker 数会改变收益。

因此“两个大表 JOIN 能否并行”不能只看顶层 Gather。要看每个 input 的 actual rows/loops、hash memory/batches、排序与 buffer。

并行不是免费 CPU

一条查询从 8 秒降到 3 秒,可能消耗更多总 CPU。对单用户很有利,对 100 个 并发报表可能降低系统总吞吐。

容量要同时看:

single-query latency
system throughput
queue time
CPU saturation
workers requested/launched
OLTP tail latency

一个常见策略是:

interactive OLTP role:
  low statement timeout
  limited parallelism

batch analytics role:
  controlled concurrency
  selected higher parallelism
  explicit work_mem/temp limits

不要只提高全局 max_parallel_workers_per_gather

实验设置不是生产建议

本章把:

SET min_parallel_table_scan_size = 0;
SET parallel_setup_cost = 0;
SET parallel_tuple_cost = 0;

用于稳定地产生教学计划。它们刻意降低并行门槛,不是生产基线。生产应让 cost model 在真实数据、硬件和并发下选择,并通过回归计划验证。

17.2.2 分区、物化视图、增量汇总与批处理

四种手段解决四个问题

手段 主要减少什么 不自动解决什么
分区 无关分区扫描与维护范围 单节点总容量、所有查询
物化视图 重复计算 自动实时增量、定义演进
增量汇总表 每次重扫历史 迟到修正、幂等与对账
批处理 峰值并发与重复启动 单批本身的坏计划

它们可以组合,但不能互相替代。

分区首先是数据管理边界

原生分区适合:

按时间快速 detach/drop 历史
按边界独立装载或维护
让明确谓词裁剪无关分区
缩小部分索引和 vacuum 的工作单元

它不是把数据自动放到多台机器。PostgreSQL 原生 declarative partitioning 仍可完全位于一个实例、一个 tablespace 和一个故障域。

设计分区前回答:

主要删除/归档边界是什么?
查询是否稳定携带分区键?
分区数量与规划成本是否可控?
唯一约束能否包含分区键?
跨分区更新和 default partition 如何处理?
备份、vacuum、索引和 schema change 如何编排?

本章协调端把外表挂到 LIST 分区父表,是为了展示租户裁剪和路由,不是把 原生分区冒充成分布式引擎。

物化视图保存一个可重建结果

本章:

CREATE MATERIALIZED VIEW
  shop_ch17.daily_tenant_summary AS
SELECT
  tenant_id,
  occurred_on,
  channel,
  count(*) AS sale_count,
  sum(units)::bigint AS unit_count,
  sum(amount)::numeric(18,2) AS amount_total
FROM shop_ch17.sales_fact
GROUP BY tenant_id, occurred_on, channel
WITH DATA;

CREATE UNIQUE INDEX daily_tenant_summary_pkey
ON shop_ch17.daily_tenant_summary (
  tenant_id,
  occurred_on,
  channel
);

冻结数据得到:

240,000 raw facts
  -> 2,880 tenant/day/channel summaries
  -> 32 tenant/month rows

原表月报计划:

Parallel Seq Scan on sales_fact
actual rows=80000 loops=3

汇总月报计划:

Seq Scan on daily_tenant_summary
actual rows=2880 loops=1

两个输出逐字节相同。这个对比证明的是“缩小输入粒度”,不是物化视图对所有 查询都快。

PostgreSQL 原生 refresh 不是自动增量维护

普通物化视图需要:

REFRESH MATERIALIZED VIEW shop_ch17.daily_tenant_summary;

或在满足条件时:

REFRESH MATERIALIZED VIEW CONCURRENTLY
  shop_ch17.daily_tenant_summary;

核心 PostgreSQL 不会因为 base table 新增一行,就自动把对应增量加进这个 物化视图。CONCURRENTLY 解决读可用性的一部分,并不把刷新变成免费,也不 替你定义迟到事实、删除、修正和失败恢复。

发布合同应固定:

refresh owner
schedule and trigger
maximum freshness lag
unique index prerequisite
expected duration and WAL
lock behavior
failure alert
retry/idempotency
late-arrival window
full rebuild path
definition version
checksum/reconciliation

官方 Materialized Views 说明结果持久化、不可直接更新和 refresh 行为。

增量汇总表是一项应用协议

若完整 refresh 太贵,可以自己维护 summary table:

raw immutable events
  -> watermark / changed key set
  -> recompute affected tenant/day buckets
  -> upsert summary
  -> record batch identity and source watermark
  -> reconcile checksum

推荐按“重算受影响桶”而非“对旧值直接 +delta”开始,因为:

  • 迟到事件可能修改历史日期;
  • 事件可能撤销或更正;
  • 重试必须幂等;
  • 聚合逻辑会升级;
  • min/max/distinct 一类聚合不容易用简单减加回滚;
  • 需要从 raw truth 完整重建。

一张稳健的汇总控制表可以记录:

CREATE TABLE summary_batch (
  batch_id          uuid PRIMARY KEY,
  definition_version text NOT NULL,
  source_from       timestamptz NOT NULL,
  source_to         timestamptz NOT NULL,
  started_at        timestamptz NOT NULL,
  finished_at       timestamptz,
  status            text NOT NULL,
  source_checksum   text,
  result_checksum   text
);

这比“每五分钟跑一条 UPSERT”多了一层治理,但也使失败可恢复、结果可解释。

批处理是调度与资源控制

把 100 个 dashboard 请求合并为一个定时汇总,减少的是:

duplicate scans
query startup
concurrency spikes
cache churn
client retries

批处理仍需要:

  • 明确 batch 边界和 watermark;
  • 限制最大运行时间与并发;
  • 避免与 checkpoint、backup、vacuum 高峰重叠;
  • 在失败后从确定位置重跑;
  • 不用一个长事务覆盖整个历史;
  • 控制 WAL、temp 和副本 lag;
  • 给消费者暴露最后成功批次与数据新鲜度。

分区与汇总的组合

一个常见设计:

raw facts partitioned by event month
daily summaries keyed by tenant/day
monthly closed partitions become immutable
current/late window can be recomputed
old raw partitions retained or archived by policy

好处是:

  • 新鲜窗口小;
  • 历史汇总稳定;
  • 迟到修正有明确范围;
  • 全量重建可按分区推进;
  • 对账可以逐分区做。

风险是出现两套粒度与状态机。必须写清:

哪张表是最终事实?
汇总多久可旧?
历史是否允许更正?
定义升级如何双跑?
消费者如何选择版本?
raw 删除后是否仍能重建?

17.2.3 列式能力候选必须写入版本基线

“列式”不是一个单一功能

候选可能提供:

columnar storage
vectorized execution
compression
late materialization
parallel scan
external file scan
cache/format conversion
specialized aggregate

一项产品或扩展拥有其中一个,不表示拥有全部。也不能从“压缩率更高”推导 “点查、更新、复制和恢复都更好”。

先写 workload fit

列式路径通常更适合:

  • 只读或追加为主;
  • 扫描少数列、很多行;
  • 聚合和过滤占主导;
  • 批量装载;
  • 更新/删除少;
  • 可以接受特定事务和索引限制。

行存 PostgreSQL 通常在以下方面仍有优势:

  • 高选择性点查;
  • 频繁小事务更新;
  • 丰富 B-tree/GIN/GiST/SP-GiST 索引;
  • 完整约束、触发器与扩展组合;
  • 成熟复制、PITR 和工具链;
  • 单一数据副本与事务语义。

真实系统常混合两类负载,所以问题通常不是“行存还是列存”,而是:

哪些数据、哪些查询、在哪个新鲜度和事务边界下使用哪条路径?

版本是功能的一部分

一个可执行基线至少固定:

PostgreSQL major/minor
extension/product exact version
operating system and package source
storage format version
required shared_preload_libraries
GUC baseline
CPU architecture and instruction set
license
supported backup/restore path
supported upgrade path
replica behavior
known incompatibilities

不能写:

uses columnar extension

而应写:

candidate X exact version Y
on PostgreSQL 18.x
package repository Z
validated on every Pigsty L1 node
backup/restore drill identifier ...

本章正式实验没有安装列式扩展,因此 baseline-v1.5-proposal.json 明确只验证行存、BRIN、物化和 loopback FDW。没有运行的候选不会出现在 “已验证”清单里。

查询兼容之外的基线

列式候选还要验证:

类别 问题
DML insert/update/delete/upsert/truncate 支持到哪?
DDL alter type、default、constraint、partition 如何?
索引 哪些 access method、unique、FK 可用?
MVCC snapshot、vacuum、HOT、freeze 如何变化?
WAL/复制 physical/logical、PITR、standby 是否支持?
扩展 PostGIS、vector、FDW、UDF 能否组合?
备份 工具是否理解存储格式?
升级 大版本与扩展版本如何排序?
观测 size、I/O、bloat、query metrics 是否可见?
许可 部署、节点、商业使用与再分发条件?

“SQL 跑通”只覆盖第一行的一小部分。

基准必须包含负面工作负载

不要只跑候选擅长的宽表聚合。还要包含:

single-row lookup
selective range query
high-concurrency small reads
batch insert
small update/delete
schema evolution
vacuum/compaction
backup while serving
restore and checksum
replica catch-up
node or process restart

选型不是找一个最高分,而是确认它在必要场景上没有不可接受的零分。

17.2.4 OLTP 与分析负载在同机共存的代价

共存争用表

资源 OLTP 典型需求 OLAP 典型行为 冲突
CPU 短请求低尾延迟 长扫描/聚合吞吐 worker 抢核心
shared buffers 热索引与热点页 大范围扫描 缓存污染
OS page cache 热数据 顺序历史读 热页被挤出
memory 小且稳定 sort/hash 波动 OOM/回收
storage 小随机 I/O、WAL 大顺序/临时 I/O 队列延迟
locks/snapshot 短事务 长快照/refresh vacuum/DDL
connections 短会话/池 少量长查询 slot 与队列
replicas 低 lag replay 与只读查询 recovery conflict

同一 SQL 在夜间快、白天慢,不一定是计划变化;可能是共存资源不同。

缓存命中率不能单独判断

分析大扫描可能有很高 shared hit,因为数据已经在缓存;它仍会消耗 CPU 并 驱逐其他热页。也可能有较低命中但利用高吞吐顺序读,对自己的完成时间尚可, 却让 OLTP 随机读尾延迟变差。

需要把:

database buffers
OS I/O
query latency
system throughput
OLTP tail latency

放在同一时间轴。

会话级护栏

对分析角色可以评审:

ALTER ROLE analyst SET statement_timeout = '10min';
ALTER ROLE analyst SET lock_timeout = '2s';
ALTER ROLE analyst SET idle_in_transaction_session_timeout = '1min';
ALTER ROLE analyst SET temp_file_limit = '20GB';
ALTER ROLE analyst SET work_mem = '64MB';
ALTER ROLE analyst SET max_parallel_workers_per_gather = 2;

数值只是示意,必须按容量计算。角色设置也不是资源管理器:它不能严格保证 CPU 百分比或 IOPS,仍需要连接池并发、作业调度、操作系统资源或实例隔离。

连接池与任务队列

分析任务应有独立入口和并发上限:

application request
  -> analytics queue
  -> bounded worker pool
  -> analyst database role
  -> statement/temp/parallel limits

这样过载首先表现为可观测排队,而不是所有查询同时进入数据库后互相拖垮。 队列本身要有:

deadline
priority
cancellation
deduplication
retry policy
idempotency
queue age alert

副本隔离不是免费复制

把报表放到只读副本可以隔离部分 CPU 和读 I/O,但仍共享:

  • primary 产生 WAL 的成本;
  • 网络带宽;
  • replay lag;
  • 长查询与 recovery conflict;
  • schema/extension 版本;
  • failover 时的角色变化;
  • 备份和维护体系。

还必须接受“副本可能比 primary 旧”。如果查询要求 read-your-writes 或刚提交 即见,不能无条件路由到异步副本。

Pigsty 4.5 把 offline 实例用于慢查询、ETL、OLAP 和交互查询隔离,也允许 在现有 replica 上设置 pg_offline_query。其当前行为与服务归属见 Cluster / Instance。 这是比直接分片更低一层的候选。

单独分析集群

若副本上的物理复制语义仍不合适,可以建立:

OLTP source
  -> logical replication / CDC / batch load
  -> independent analytical PostgreSQL cluster

它进一步隔离参数、存储、索引和维护,却引入:

data pipeline
schema propagation
freshness lag
replay/idempotency
DDL compatibility
backfill
dual-system reconciliation

是否比 Citus 或专用 OLAP 更合适,要由工作负载和运行模型决定。

何时单机能力已经被合理用尽

至少满足:

  • 大查询的扫描、连接、聚合路径合理;
  • worker planned/launched 与并发预算相符;
  • 选择性查询有正确索引;
  • 分区裁剪能消除无关数据;
  • spill 被量化并有会话级边界;
  • 重复历史计算已评估物化/汇总;
  • OLTP 与分析已有入口和资源隔离;
  • backup、vacuum、checkpoint、replica lag 一同压测;
  • 硬件纵向扩容与未来增长已建模;
  • 正确性和新鲜度仍满足。

只有到这一步,“单节点哪一种资源仍越界”才有明确答案。下一节据此定义何时 需要分布式,以及分片键会把哪些数据库语义变成应用必须承担的合同。


上一节:先证明单机边界 · 返回本章目录 · 下一节:何时需要分布式 · 查看全书目录 · 查看索引中心

17.3 何时需要分布式

分布式系统的价值,是把某种不可再容纳于单一故障域的资源或责任拆开。

它的代价,是把原本由一个 PostgreSQL 实例隐式保证的事实,变成显式协议:

row lives where?
query runs where?
transaction spans where?
failure is partial or total?
backup represents which global point?
schema change reaches which nodes?
how is data rebalanced?
how do we leave?

因此,“需要分布式”的完整句子必须包含:

因为哪一种经过测量的边界,采用哪一种拆分单位,并接受哪些一致性、延迟、 可用性和运维代价。

17.3.1 容量、吞吐、地域与组织边界

容量边界

容量可以指:

heap + index + TOAST storage
working set memory
WAL generation and retention
backup repository and window
restore duration
vacuum/freeze maintenance window
index build / schema change window
replica catch-up

“磁盘还放得下”不是容量充足。如果一个 40TB 节点需要 30 小时才能从备份 恢复,而 RTO 是 2 小时,恢复窗口已经越界;如果冻结维护无法追上事务年龄, 也已经越界。

容量评审要有未来曲线:

[ \text{projected bytes}(t) = \text{current bytes}

  • \text{daily net growth} \times t
  • \text{index/WAL/maintenance headroom} ]

还要加入:

  • 增长误差区间;
  • 最大租户与平均租户差异;
  • 扩容交付时间;
  • 节点故障时的冗余;
  • 大版本升级期间的双份空间;
  • 重分片期间的额外副本。

只用平均增长会低估热点和迁移峰值。

吞吐边界

吞吐越界不是“单查询太慢”,而是调优后:

arrival rate > sustainable completion rate
queue age continuously grows
CPU/I/O at safe ceiling
tail latency and timeout rise
more concurrency no longer increases throughput

横向拆分只有在工作可并行且协调开销小于收益时有用。若所有请求都需要访问 全部分片,增加节点可能同时增加 fan-out、连接和合并成本。

应把工作负载分类:

查询形状 横向拆分潜力
单租户、单实体 高,若所有相关数据同分片
跨租户可分解聚合 中高,可做 partial/final
全局 top-N 需要各分片候选 + 全局合并
跨分片大 JOIN 低,可能需要 shuffle
全局唯一写入 需要协调或新语义
强事务跨多个实体 代价高

地域边界

地理分布可能为了:

latency
data residency
disaster recovery
network sovereignty
organizational autonomy

这些目标不能混为“多活”。例如:

  • 就近读缓存不等于可在多地写同一行;
  • 数据驻留可能禁止跨境复制;
  • DR standby 不是日常承载写入的 active-active;
  • 跨地域同步提交会把网络 RTT 加进事务延迟;
  • 异步复制会引入 RPO 和陈旧读。

先定义:

who may write where
which data may replicate where
conflict owner
failover authority
RPO/RTO per region
rejoin and divergence handling

第 26、27 章会深入复制、高可用与故障切换;本章只把地域作为进入选型的边界, 不提前把一个 loopback FDW PoC 说成多地域方案。

组织边界

有时拆分的主要原因不是硬件,而是责任:

  • 不同团队有独立发布节奏;
  • 法规要求独立访问控制和审计;
  • 某业务需要独立 SLO 与故障域;
  • 租户需要物理隔离;
  • 成本需要可归属;
  • 数据生命周期不同。

但数据库拆分不会自动解决组织问题。若跨域查询、事务和 schema 仍高度耦合, 拆库只会把内部调用变成网络调用。

评审组织拆分时画出:

data owner
schema owner
writer
reader
cross-domain transaction
cross-domain report
incident owner
backup/restore owner

没有唯一 owner 的共享表,是分布式后最容易成为争议中心的对象。

四种边界的硬证据

边界 不能只说 至少要提供
容量 数据很大 增长、工作集、维护/恢复窗口
吞吐 QPS 很高 饱和曲线、查询 mix、排队
地域 用户遍布全球 RTT、驻留、写入/冲突/RPO
组织 微服务化 ownership 与跨域依赖图

不是分布式门槛的信号

以下现象本身不足以证明:

  • 单条查询偶发慢;
  • 某次 CPU 100%;
  • 行数达到一个整齐数量级;
  • 云厂商有一项“分布式”产品;
  • 团队担心未来增长;
  • 某竞品使用很多节点;
  • 一个 demo 在笔记本上跑通;
  • “PostgreSQL 是单机数据库”。

它们可以触发测量,不能直接触发迁移。

17.3.2 分片键、数据局部性与跨分片事务

分片键决定数据库能否继续像数据库

好的分片键同时满足:

high enough cardinality
balanced data and load
stable over entity lifetime
present in major queries
present in joins and transactions
allows related rows to co-locate
supports operational moves

这些条件常冲突。tenant_id 局部性好,但一个超级租户可能形成热点; event_id 均匀,却把同一租户的查询撒到全部节点;时间范围方便归档,却可能 让“最近一天”的所有写入集中在一个分片。

从查询反推分片键

建立 query-to-key matrix:

查询/事务 频率 候选键 单分片? 跨分片代价
租户 dashboard tenant_id
租户账户与销售 JOIN tenant_id
全局月报 tenant_id partial/final
跨租户对账 tenant_id fan-out
账户迁移租户 极低 tenant_id 数据移动

不是从最大表里挑一个列名,而是看主要业务单元能否局部闭合。

本章的租户路由

冻结生成器把:

tenant_id % 2 = 0 -> pg36_shard_a
tenant_id % 2 = 1 -> pg36_shard_b

协调端使用显式 LIST 边界:

CREATE TABLE shop_ch17.sales_fact_distributed (...)
PARTITION BY LIST (tenant_id);

CREATE FOREIGN TABLE shop_ch17.sales_fact_dist_0
PARTITION OF shop_ch17.sales_fact_distributed
FOR VALUES IN (2, 4, 6, 8)
SERVER pg36_ch17_shard_a;

CREATE FOREIGN TABLE shop_ch17.sales_fact_dist_1
PARTITION OF shop_ch17.sales_fact_distributed
FOR VALUES IN (1, 3, 5, 7)
SERVER pg36_ch17_shard_b;

这让租户 3 谓词可被裁剪到 sales_fact_dist_1

关键反例:HASH remainder ≠ 整数取模

错误原型:

CREATE TABLE ... PARTITION BY HASH (tenant_id);

CREATE FOREIGN TABLE ... PARTITION OF ...
FOR VALUES WITH (MODULUS 2, REMAINDER 0);

直觉误读:

remainder 0 -> even tenant_id
remainder 1 -> odd tenant_id

实际不是。PostgreSQL HASH partitioning 对分区键调用内部哈希支持,再依据 组合哈希值选择 remainder。原值为 2,不保证进入 remainder 0。

危险路径:

remote loader places tenant 3 by 3 % 2 -> shard B
PostgreSQL hashes tenant 3 -> perhaps chooses another remainder
partition pruning trusts PostgreSQL partition bounds
only the chosen foreign partition is scanned
tenant 3 physically absent there
query returns zero or partial rows without transport error

这是“可用性正常、SQL 成功、结果错误”的最坏一类问题。

本章没有通过关闭 partition pruning 掩盖它,而是修正路由合同。生产分片系统 必须保证:

writer router
reader router
catalog metadata
rebalance tool
restore tool
application cache

使用同一个版本化映射。若更换 hash function、seed、token range 或 shard count,需要正式数据迁移,不能只改配置。

验证物理放置

不要只查父表总数。验证每个物理分片:

SELECT
  tableoid::regclass,
  count(*),
  min(tenant_id),
  max(tenant_id)
FROM shop_ch17.sales_fact_distributed
GROUP BY tableoid;

再检查:

SELECT *
FROM shop_ch17.sales_fact_dist_0
WHERE mod(tenant_id, 2) <> 0;

SELECT *
FROM shop_ch17.sales_fact_dist_1
WHERE mod(tenant_id, 2) <> 1;

冻结结果:

dist_0 = 120000 rows, tenants 2,4,6,8
dist_1 = 120000 rows, tenants 1,3,5,7

生产还应保存每个 shard 的 count、checksum、key range、路由 epoch 与采集时间。

数据局部性不止“在同一节点”

一个 JOIN 要局部执行,通常需要:

same distribution key
same key type and semantics
compatible shard mapping/colocation group
join predicate includes the key
query can be pushed by the target engine
functions/collations/types are compatible
statistics and cost favor pushdown
same user mapping where FDW requires it

本章账户与销售在同一远端库,查询也有 tenant_id = 3。但通过两个 partitioned foreign-table parents 查询时,实测:

Hash Join on coordinator
  Foreign Scan sales_fact_dist_1 -> 7500 rows
  Foreign Scan account_dim_dist_1 -> 50 rows

这证明“物理共置”不等于“当前抽象层与规划器已把 JOIN 下推”。官方 postgres_fdw 文档说明同一个 foreign server 上的外表 JOIN 可能 整体 发送到远端,但规划器仍可能判断分别取回更合适,其他限制也会阻止下推; 实际远端 SQL 要用 EXPLAIN VERBOSE 检查。参见 postgres_fdw Remote Query Optimization

跨分片查询的四种形状

  1. 单分片路由

    WHERE tenant_id = ?

    理想情况下只访问一组 colocated shards。

  2. scatter/gather

    每个分片执行相同查询,协调端合并。节点越多,fan-out 与尾延迟越显著。

  3. partial/final aggregation

    远端先聚合,协调端合并小结果。本章从 240,000 条事实缩到 960 条日汇总。

  4. repartition/shuffle

    按另一个 join/group key 跨网络重分布。功能强,但网络、磁盘和故障复杂度 高,通常是分片设计不局部的成本中心。

候选产品对四种形状的支持不同,不能只用单租户点查比较。

跨分片事务

单机事务隐含:

one WAL stream
one transaction manager
one commit decision
one snapshot domain

跨节点后要回答:

  • 是否支持原子 commit?
  • 协调端在什么时点记录决定?
  • 某个参与者 commit 后网络断开怎么办?
  • retry 会不会重复写?
  • prepared transaction 谁清理?
  • snapshot 是否跨节点一致?
  • deadlock 是否跨节点检测?
  • 一个节点长期不可用时业务阻塞还是降级?

postgres_fdw 会为本地事务打开相应远端事务,并映射 savepoint;PostgreSQL 18 官方文档明确说明,它目前不支持把远端事务 prepare 为 two-phase commit。 所以本章只做只读分析与权限拒绝,不用这个 PoC 宣称具备通用跨分片原子写。

全局约束

分片后重新评审:

PRIMARY KEY / UNIQUE
FOREIGN KEY
sequence / identity
exclusion constraint
serializable invariant
trigger
advisory lock

如果唯一键不包含分布键,可能需要:

  • 中央目录;
  • 分布式协调;
  • 业务生成全局唯一 ID;
  • 接受仅分片内唯一;
  • 改模型。

不能假设一个节点上的 local index 会检查其他节点。

热点与大租户

按 tenant 均匀 hash 只保证 key 空间近似均匀,不保证:

bytes per tenant
queries per tenant
writes per tenant
CPU per tenant
time-of-day peak

监测:

skew ratio=largest shard loadmean shard load \text{skew ratio} = \frac{\text{largest shard load}} {\text{mean shard load}}

并分别对 bytes、QPS、CPU、I/O 计算。一个超级租户可能需要独立 shard、 二级分片或专门迁移机制;这应在选型前验证,而不是上线后临时手工搬表。

17.3.3 一致性、运维复杂度与退出成本

分布式首先改变故障集合

单节点主要状态:

up / down / recovering

两分片加协调端至少有:

coordinator down
shard A down
shard B down
coordinator can reach A but not B
client can reach coordinator but coordinator DNS/TLS/auth fails to B
schema version differs
route catalog stale
one shard lagging
partial rebalance

每一种都要定义读写语义。

本章的部分失败探针

shard-failure.sql 在一个事务里暂时把 shard B server port 设为不可达,并断开旧连接:

BEGIN;

ALTER SERVER pg36_ch17_shard_b
  OPTIONS (SET port '1');

SELECT shop_ch17_ext.postgres_fdw_disconnect(
  'pg36_ch17_shard_b'
);

随后先查 tenant 2:

SELECT count(*)
FROM shop_ch17.sales_fact_distributed
WHERE tenant_id = 2;

冻结输出:

healthy_shard_tenant_2=30000

因为 LIST pruning 只访问 shard A。

再查全局:

SELECT count(*)
FROM shop_ch17.sales_fact_distributed;

需要两个 shard,固定以 SQLSTATE 08001 失败。psql 因错误断开后,事务回滚, 任务再次导出 server catalog,并要求与故障前逐字节一致。

这个实验定义了机制,不自动定义产品语义。应用必须决定:

请求 shard B 不可用时
tenant 2 明细 可继续,标注路由 epoch
tenant 3 明细 失败,不返回“空”
全局月报 失败或明确标为 partial
写 shard A 是否允许取决于全局不变量
写 shard B 失败/排队,重试必须幂等

最危险的是默默返回 partial result 却不标识。

一致性不只是隔离级别

分布式分析至少有:

snapshot consistency
replication freshness
route consistency
schema consistency
summary definition consistency
global completeness

例如两个分片分别在 10:00:01 与 10:00:05 读取,即使各自是 Repeatable Read, 合并结果也未必代表同一个业务时点。必须明确报告需要:

  • latest available;
  • bounded staleness;
  • consistent cut;
  • closed period;
  • 或允许 approximate/partial。

备份不是“每台都备一份”

全局恢复要回答:

which coordinator metadata version?
which route epoch?
which point in time on every shard?
are distributed transactions in doubt?
are reference data and sequences aligned?
how is restored placement verified?

各节点各有成功备份,不表示组合后是业务一致的恢复点。恢复演练必须从空环境 重建拓扑,恢复数据,验证路由、行数、checksum 和应用查询。

schema change 从一次 DDL 变成编排

要考虑:

  • coordinator 与 worker 执行顺序;
  • mixed-version window;
  • old/new application compatibility;
  • 某节点 DDL 失败后的补偿;
  • index build 资源;
  • backfill 与 WAL;
  • 回滚是否还可逆;
  • 路由和 schema epoch。

“支持 ALTER TABLE”不等于在十个节点、在线负载和失败条件下安全。

监控维度乘法

除了每个 PostgreSQL 节点原有指标,还要有:

per-shard bytes/QPS/CPU/IO
skew
cross-shard query ratio
fan-out
remote rows/bytes
coordinator queue
route failures
rebalance progress
schema version drift
partial result count
distributed transaction state

平均指标尤其危险:平均 shard CPU 40% 可能掩盖一个 100% 的热点 shard。

再平衡是一项在线数据迁移

增加节点不会让旧数据自动均匀且无代价地移动。再平衡需要预算:

source read
destination write
network
WAL and replication
cache coldness
lock/metadata change
double storage
failure and resume
verification

还要决定迁移期间路由如何读写,以及最后切换点。

退出成本在进入前计算

退出路线可能是:

distributed -> larger single PostgreSQL
distributed -> independent tenant clusters
distributed -> new sharding engine
OLAP system -> PostgreSQL summaries
FDW federation -> copied local tables

提前保存:

  • canonical schema;
  • raw data export;
  • distribution map;
  • globally stable IDs;
  • ordering/watermark;
  • row counts/checksums;
  • dual-read comparison;
  • last reversible step;
  • DNS/service rollback;
  • decommission proof。

若候选使用专有类型、SQL、存储或事务语义,退出成本要显式计价。

跨数据库退出不是一个本地事务

本章 reset 需要删除协调库与两个分片库中的受管对象。每个数据库内部:

verify exact state
BEGIN
drop exact objects without CASCADE
COMMIT

但三个数据库不能被一个本地 DDL 事务原子包住。任务采用:

pre-verify coordinator + A + B
reset coordinator
reset A
reset B
rebuild all
verify all

如果中途失败,依靠可重复 reset/setup 与证据进行补偿。这一小段实验已经显示 分布式运维复杂度:即使同一实例里的三个数据库,也需要跨库编排;真实多节点 只会增加网络、权限和故障状态。

进入分布式的最终清单

只有以下问题有可审计答案,才进入产品比较:

[ ] 哪一种单节点资源已在代表性负载下越界?
[ ] 纵向扩容、汇总和副本隔离为什么不足?
[ ] 分布单位是租户、实体、时间、schema 还是别的?
[ ] 最大分片和倾斜是多少?
[ ] 多少查询/事务可以单分片?
[ ] 跨分片查询如何聚合或 shuffle?
[ ] 全局唯一、FK 和事务不变量如何变化?
[ ] 单节点、网络、协调端失败时分别返回什么?
[ ] 备份能否恢复成全局一致且路由正确的系统?
[ ] schema change、再平衡和升级是否演练?
[ ] 团队值班与工具是否能承担?
[ ] 如何迁入、如何双读验证、如何退出?

下一节在这些问题的约束下比较候选,而不是用一个“SQL 兼容”标签把不同系统 压成同一类。


上一节:单机分析能力 · 返回本章目录 · 下一节:比较分布式候选 · 查看全书目录 · 查看索引中心

17.4 比较分布式候选

候选比较的第一步,是保留“不分布”作为对照组。

如果表格只有三种分布式产品,团队会被迫在它们之间选一个;如果加入:

optimized primary
offline replica
independent analytical PostgreSQL
materialized summary

很多问题会暴露为负载隔离或数据粒度问题,而不是横向分片问题。

本节不替读者宣布一个普遍赢家,而是建立同口径的比较方法。

17.4.1 PostgreSQL 扩展、兼容数据库与专用分析系统

候选地图

可以按“离 PostgreSQL 原生语义有多远”分层:

Layer 0: same PostgreSQL instance
  indexes / partitions / parallel query / summaries

Layer 1: same physical PostgreSQL data
  read replica / offline replica

Layer 2: another PostgreSQL data copy
  logical replication / CDC / ETL to analytical PostgreSQL

Layer 3: PostgreSQL extension or federation
  postgres_fdw / Citus / workload-specific extensions

Layer 4: PostgreSQL-wire or SQL-compatible distributed database
  separate storage, transaction and operations implementation

Layer 5: specialized analytical system
  columnar/vectorized engine, separate ingestion and lifecycle

层数越高不表示越先进,只表示需要重新验证的语义越多。

对照组:优化后的 PostgreSQL

它应包含:

  • 正确 schema 与统计信息;
  • 代表性索引和 partition pruning;
  • 并行计划与并发预算;
  • 合理物化/汇总;
  • workload guardrails;
  • 当前硬件与一个可行纵向规格;
  • offline replica 或独立分析副本候选。

若分布式候选只比未经调优的原表扫描快,比较没有意义。

postgres_fdw:联邦访问与机制实验

postgres_fdw 把远端 PostgreSQL 表映射为 foreign table:

CREATE EXTENSION
  -> CREATE SERVER
  -> CREATE USER MAPPING
  -> CREATE FOREIGN TABLE / IMPORT FOREIGN SCHEMA
  -> SELECT / DML

它能做:

  • 远端过滤和列裁剪;
  • 某些 JOIN/聚合下推;
  • 外表分区;
  • 联邦读取;
  • 迁移期间的过渡;
  • 把远端数据物化到本地;
  • 显示远端 SQL 和数据流。

它不应被默认理解为一个完整的透明分布式数据库。路由、rebalance、全局 约束、协调端 HA、全局备份点和很多运维职责仍需自行设计。PostgreSQL 18 的 postgres_fdw 远端事务也不支持 prepare 为两阶段提交。

本章用它是因为机制透明:

Foreign Scan
Remote SQL
actual rows
user mapping
server failure

都能直接观察。它是很好的教学镜子,不是本章对生产选型的默认推荐。

Citus:PostgreSQL 扩展式分片候选

Citus 把表区分为 distributed、reference、local 等类型,以分布列决定行的 放置;相关表按相同分布键 colocate 后,单租户查询和某些 JOIN 可以在一组 共置分片上执行。跨租户聚合则可由 worker 产生 partial result,再由协调端 合并。

这使它适合评估:

multi-tenant workload with tenant-local transactions
real-time aggregate workload with decomposable computation
PostgreSQL ecosystem continuity

但分布键选择成为 schema 与查询合同。Citus 官方 Choosing Distribution Column 强调 tenant/entity key、co-location 与跨节点数据移动之间的关系;小型共享 维表可评估 reference table。

需要验证:

  • row-based 还是 schema-based sharding;
  • distribution column 是否出现在 PK/FK/查询;
  • colocated tables 与 reference tables;
  • 单租户与全局查询比例;
  • rebalance 与大租户隔离;
  • coordinator/worker HA;
  • distributed DDL;
  • backup/restore;
  • 版本升级和扩展组合;
  • Pigsty L1 的实际拓扑。

本章 loopback FDW 的“同分片 JOIN 未下推”不能直接外推为 Citus 行为;它只 提醒读者对目标产品的目标 SQL 看实际计划。

PostgreSQL-compatible distributed database

这类系统可能支持 PostgreSQL wire protocol、部分 SQL、驱动和工具,让迁移 起步更容易。必须拆开“兼容”:

wire protocol
parser syntax
catalog shape
data types
functions/operators
transaction semantics
isolation/locking
extensions
backup/restore
monitoring
operational commands

一个应用只用简单 SELECT/INSERT,兼容度可能足够;一个应用依赖 PostGIS、 自定义 operator class、logical decoding、advisory lock、trigger、COPY、 RLS 和精细 catalog 查询,迁移面完全不同。

不要用厂商兼容百分比替代自己的 feature inventory。

专用分析系统

专用 OLAP 系统通常优化:

columnar compression
large scans
vectorized execution
distributed aggregation
high analytical concurrency
object storage / tiering

它可能显著优于 PostgreSQL 处理某类宽表聚合,但会引入第二套:

ingestion / CDC
schema mapping
data freshness
deduplication
late-event handling
access control
backup/recovery
monitoring/on-call
cost model
query semantics

如果 OLTP truth 仍在 PostgreSQL,必须定义两个系统不一致时谁是事实来源。

候选能力矩阵

以下不是产品评分,而是评审问题:

维度 单机/副本 FDW Citus 类扩展 兼容分布式库 专用 OLAP
原生 PG 语义 最高 本地/远端 PG 高但分片有边界 必须实测 通常较低
单租户局部性 原生 手工路由 分布键核心 依实现 依模型
全局聚合 单节点 可下推/合并 分布式 partial/final 依实现 核心场景
跨分片事务 不适用 有明显限制 需按版本/形状验证 需验证 通常非 OLTP 重点
扩展生态 原生 两端一致性 需验证组合 通常有限 不适用/自有
运维体系 已有 多 PG + 编排 coordinator/workers 新体系 第二套体系
新鲜度 即时/复制 lag 远端即时视连接 即时视事务 依实现 CDC/批次 lag
退出成本 中高

矩阵中的每个“需验证”都要变成 PoC 用例。

所有权边界

候选不仅有技术 owner:

责任 必须有人承担
schema 与分布键 数据模型 owner
query migration 应用 owner
ingestion/CDC 数据平台 owner
cluster lifecycle DBA/SRE
correctness reconciliation 业务 owner + 数据 owner
incident decision on-call
cost 预算 owner
exit 项目 sponsor

缺少 owner 的候选不能因为跑分快进入生产。

17.4.2 SQL 兼容不等于事务、扩展和运维兼容

建立兼容性分层

推荐至少分八层:

L1 protocol
L2 syntax
L3 type and expression semantics
L4 transaction and concurrency
L5 schema objects and extensions
L6 planner/performance behavior
L7 operations and observability
L8 failure/recovery and lifecycle

只有 L1/L2 通过,应用仍可能在 L3–L8 失败。

协议兼容

验证:

  • TLS、SCRAM、GSS/SSO;
  • connection parameters;
  • prepared statements;
  • binary/text format;
  • COPY;
  • cancellation;
  • notices/errors 与 SQLSTATE;
  • connection pool transaction/session mode;
  • driver features;
  • failover reconnect。

能用 psql 登录只是起点。

SQL 与类型语义

测试应用真实使用的:

numeric precision and rounding
timestamp/time zone
collation and locale
NULL ordering
JSON/JSONB
arrays/ranges/multiranges
generated columns
identity/sequence
UPSERT/MERGE/RETURNING
CTE/window/lateral
recursive query

同样语法若类型、collation 或时区不同,会返回不同结果。

postgres_fdw 官方文档也建议 foreign table 的类型和 collation 与远端精确 匹配,否则本地与远端对条件的解释可能不同。它还不会自动导入除 NOT NULL 之外的约束,因为错误约束可能导致规划器做出不安全推断。

事务与并发语义

需要独立验证:

READ COMMITTED snapshot
REPEATABLE READ
SERIALIZABLE
row/table/advisory locks
deadlock detection
savepoints
DDL transactionality
cross-shard transaction
retry error classes
sequence behavior

例如“支持 serializable”不够,要验证:

  • 冲突时 SQLSTATE;
  • 是否需要 client retry;
  • 多分片是否同样保证;
  • range/predicate conflict 如何实现;
  • failover 后 in-flight transaction 的结果;
  • unknown commit 如何对账。

schema 与扩展兼容

列清单:

SELECT
  extname,
  extversion
FROM pg_catalog.pg_extension
ORDER BY extname;

对每个扩展检查:

  • 是否可安装;
  • exact version;
  • trusted/non-trusted;
  • shared preload;
  • 类型、函数、operator、index AM;
  • logical/physical replication;
  • backup/restore;
  • rolling upgrade;
  • 每个节点一致性。

不能把“兼容 PostgreSQL”理解为兼容任意 PostgreSQL 扩展。

catalog 兼容

许多工具读取:

pg_catalog
information_schema
pg_stat_*
pg_locks
pg_settings
pg_extension
pg_class/pg_attribute/pg_index

候选可能接受这些查询但字段为空、语义不同或只反映 coordinator。验证:

  • migration tool;
  • ORM introspection;
  • monitoring exporter;
  • backup tool;
  • schema diff;
  • incident runbook;
  • 自定义运维脚本。

planner 兼容

相同 SQL 的计划可以完全不同。需要观察:

where execution happens
which shards are pruned
which filters/joins/aggregates push down
how many rows cross network
how coordinator merges
what spills
what happens under skew

本章同一个 FDW 月报有两种 SQL:

parent aggregate:
  fetch 240,000 facts
  aggregate locally

explicit per-shard daily aggregate:
  fetch 480 + 480 aggregate rows
  aggregate monthly locally

语义相同,执行位置不同。兼容性测试不能只检查最终 rows。

运维兼容

对比日常动作:

动作 问题
provision 声明、包、密钥、节点身份如何?
scale 加节点是否自动,旧数据如何移动?
backup 一致点、加密、保留、校验如何?
restore 空环境恢复、route metadata 如何?
failover coordinator/worker 谁仲裁?
upgrade rolling、停机、扩展顺序?
schema change fan-out、失败补偿?
observability 全局与每 shard 指标?
security HBA、证书、user mapping、secret?
decommission 数据擦除与证明?

工具名字相同也不表示语义相同。例如在 coordinator 执行 VACUUM 是否覆盖 所有 shard,要由目标产品和版本证明。

错误兼容

应用通常围绕 SQLSTATE 决定:

retry
abort
return conflict
mark dependency unavailable

候选必须保留或重新映射这些错误语义。测试:

  • unique violation;
  • serialization failure;
  • deadlock;
  • lock timeout;
  • statement timeout;
  • connection failure;
  • read-only transaction;
  • insufficient privilege;
  • disk/full or quota;
  • shard unavailable。

只测试成功路径会让第一场故障变成兼容性测试。

精度与顺序

分析结果的隐性差异:

floating aggregate order
approximate distinct
collation sort
NULL order
time zone database version
decimal scale
non-deterministic top-N ties

冻结输出应:

  • 使用 exact numeric 或定义误差;
  • 明确 ORDER BY 与 tie-breaker;
  • 固定时区/locale;
  • 记录 approximate 算法/version/seed;
  • 对 checksum 使用稳定序列化。

17.4.3 用同一工作负载和失败条件比较

先冻结语义

比较协议应先固定:

schema
data generator / snapshot
business queries
expected results
freshness point
concurrency schedule
failure schedule
versions/config
measurement method

本章:

fixture = ch17-analytics-v1
frozen_at = 2026-07-29T00:00:00Z
monthly rows = 32
monthly checksum = 644d45544ebbc2a80c42270c38ac6885

任何候选先生成同一月报,再谈性能。

workload suite

至少包含:

  1. 高选择性单租户读

    WHERE tenant_id = 3
      AND occurred_on >= DATE '2026-04-01'

    验证 pruning、index、route、read amplification。

  2. 全局可分解聚合

    count/sum 按租户和月份,验证 partial/final 与传输。

  3. 同分片 JOIN

    account + sales,验证真实 join pushdown。

  4. 非分布键 JOIN

    故意触发 shuffle 或拒绝,量化代价。

  5. 小写入与幂等重试

    验证事务、unique、retry。

  6. 批量装载

    验证 ingest、WAL/replication、rebalance。

  7. schema change

    增列、建索引、backfill,验证 mixed-version。

  8. backup/restore

    从空环境恢复并跑 checksum。

数据规模阶梯

不要只跑一个尺寸:

S: fits in memory
M: working set near memory
L: exceeds memory
XL: near storage/maintenance target

观察曲线和拐点,而不是挑一张最好看的柱状图。

冷暖缓存

至少区分:

warm repeated query
cold/evicted data
after restart
after rebalance
after restore
after schema/index build

专用分析系统和 PostgreSQL 对缓存、编译、数据格式转换的预热不同。只比较 第十次执行或只比较第一次执行都可能偏颇。

并发与到达模型

开放环和封闭环会得出不同结论:

closed-loop:
  client waits for response, then sends next
  overloaded system self-throttles

open-loop:
  requests arrive by schedule independent of completion
  queueing and overload become visible

生产若有固定到达率,基准不能只用少量客户端闭环。还要记录 client queue 与 server queue,避免 coordinated omission。

资源和成本同报

每个结果同时报告:

latency distribution
throughput
error rate
CPU seconds
memory peak
storage read/write
temp/spill
network bytes
WAL/replication
storage footprint
node count
operator time
license/cloud cost

“P95 快 2 倍但使用 8 倍节点”与“同成本快 2 倍”不是同一结论。

失败矩阵

对每个候选执行:

故障 验证
query cancel 远端工作是否停止、资源是否释放
worker/shard down 单分片与全局查询如何返回
coordinator down 新连接、已有事务、恢复
network partition timeout、unknown commit、重试
disk pressure backpressure 与告警
replica lag freshness 标识与路由
rebalance interrupted resume、重复/遗漏
schema node drift 拒绝、修复与可见性
backup during load 恢复一致性

失败结果必须是验收的一部分,不是“以后做 chaos”。

本章的最小失败条件

冻结 PoC 至少证明:

application write -> SQLSTATE 42501
shard B unreachable -> global query SQLSTATE 08001
shard B unreachable -> tenant 2 scoped read = 30000
server catalog after rollback = before failure

它没有证明:

  • 真实网络分区;
  • process kill;
  • WAL/replica behavior;
  • coordinator HA;
  • shard failover;
  • in-flight distributed write;
  • rebalance resume。

这些是生产 PoC 的追加用例。

benchmark result template

每条结论写成:

claim:
  two-stage aggregation reduces coordinator input

environment:
  PostgreSQL 18.6, postgres_fdw 1.2
  same-host/same-instance loopback

input:
  ch17-analytics-v1, 240,000 facts

evidence:
  naive Append actual rows=240000
  two-stage Append actual rows=960
  both byte-identical to frozen monthly CSV

scope:
  row-transfer shape only

not proven:
  network bytes, latency, throughput, HA, scaling

这个模板迫使作者把结论和外推边界放在一起。

评分前设置 veto

某些条件不应靠加权平均掩盖:

incorrect result
cannot meet RPO/RTO
unsupported mandatory extension
unacceptable data residency
no recoverable backup
license conflict
no exit path
unknown-commit without business reconciliation

任何 veto 失败,候选退出;不能用“查询快 30%”抵消。

决策表

通过 veto 后再评分:

维度 权重 证据 分数 不确定性
correctness veto golden/checksum pass low
workload SLO 25 representative replay
failure/RPO/RTO 20 drills
operability 15 day-2 tasks
compatibility 15 feature inventory
cost 10 same horizon/TCO
migration 10 rehearsal
exit 5 reverse rehearsal

“不确定性”单列,避免没有验证的候选因为乐观估分胜出。

本节结论

比较方法的输出不是产品排行榜,而是:

baseline
candidate contract
evidence bundle
known limitations
veto results
cost/ownership
migration and exit
decision trigger

下一节构造一个最小 FDW PoC,目的不是给候选打性能分,而是验证“租户路由、 远端聚合、权限和部分失败能否被证据化”这一个关键假设集合。


上一节:何时需要分布式 · 返回本章目录 · 下一节:部署最小分布式 PoC · 查看全书目录 · 查看索引中心

17.5 部署最小分布式 PoC

一个好的 PoC 不是“把所有组件装一遍”,而是用最小拓扑证伪一个关键假设。

本章假设是:

tenant_id 为数据局部性边界时,协调端能裁剪到单分片;全局可分解聚合 能把部分计算推到数据侧;同时,未下推 JOIN、权限和单分片故障可以被明确 观察,而不是被 demo 隐藏。

为了让 SQL、catalog 和计划完全透明,本地 PoC 使用 postgres_fdw。它不是 对 Citus 性能或生产可用性的替代测试。

17.5.1 明确 PoC 只验证一个关键假设

拓扑

one PostgreSQL 18.6 instance
one Unix-domain socket
one process/storage/failure domain

pg36_shop                 coordinator database
  shop_ch17               local facts + partitioned foreign parents
  shop_ch17_ext           postgres_fdw 1.2
  pg36_ch17_shard_a       foreign server
  pg36_ch17_shard_b       foreign server

pg36_shard_a              retained database shell
  shop_ch17_shard         tenants 2,4,6,8

pg36_shard_b              retained database shell
  shop_ch17_shard         tenants 1,3,5,7

三个数据库共享一个实例。因此它能证明:

database boundary
foreign server/user mapping
partition routing/pruning
remote SQL
row transfer shape
remote failure SQLSTATE
cross-database reset orchestration

不能证明:

network latency/bandwidth
independent CPU/storage
multi-node throughput
replication/HA
node placement
rolling upgrade
rebalance
distributed backup

PoC 的验收问题

只回答十个问题:

  1. 三份数据库数据能否由同一确定生成器重建?
  2. 本地、summary、naive FDW、two-stage FDW 是否输出同一 frozen result?
  3. 单租户谓词是否裁剪到正确物理 shard?
  4. 租户和日期过滤是否出现在 Remote SQL?
  5. 朴素全局聚合向 coordinator 返回多少行?
  6. 远端预聚合后返回多少行?
  7. 同物理分片 JOIN 是否真的被下推?
  8. application role 是否只能读且使用具名 mapping?
  9. 一个 shard 不可达时,健康 shard 的 scoped read 与全局 read 分别怎样?
  10. 受管对象能否精确退出并完整重建?

没有延迟、QPS 或扩展倍数问题,因为这个拓扑没有资格回答。

冻结数据生成

协调端和远端分别使用:

共同公式:

tenant_id       = 1..8
account_id      = 1..50
day_offset      = 0..119
sale_per_day    = 1..5
sale_id         = deterministic integer composition
channel         = deterministic cycle
units/amount    = deterministic expressions

没有随机数、当前时间或外部数据。fixture_meta 固定:

fixture_version = ch17-analytics-v1
generator_identity = fixture-generator-v1
first_day = 2026-01-01
frozen_at = 2026-07-29T00:00:00Z

两端 generator 的 SHA-256 写入 fixture-manifest.json。生成器变化意味着 fixture 身份变化,不能仍用旧 golden。

数据库壳与 schema 分开

bootstrap.sql 只创建并保留两个数据库壳:

CREATE DATABASE pg36_shard_a
  WITH OWNER pg36_owner TEMPLATE template0 ENCODING 'UTF8';

CREATE DATABASE pg36_shard_b
  WITH OWNER pg36_owner TEMPLATE template0 ENCODING 'UTF8';

并固定 database comment:

pg36 ch17 fdw shard database a; retained shell
pg36 ch17 fdw shard database b; retained shell

已有同名数据库只有在 owner、comment、template/connection 身份精确匹配时才 可复用;否则停止碰撞。reset 删除内部 shop_ch17_shard,不 drop database。

这样做的理由:

  • DROP DATABASE 破坏性更大;
  • database DDL 不能在普通事务里执行;
  • 固定壳能把重复实验聚焦于 schema/data;
  • reset 输出明确说明 database_shell=retained

远端分片

每个 shard:

CREATE SCHEMA shop_ch17_shard AUTHORIZATION pg36_owner;

CREATE TABLE shop_ch17_shard.account_dim (...);
CREATE TABLE shop_ch17_shard.sales_fact (...);

分片 A 约束:

CHECK (mod(tenant_id, 2) = 0)

分片 B:

CHECK (mod(tenant_id, 2) = 1)

每个分片固定:

4 tenants
200 accounts
120,000 sales
2026-01-01 .. 2026-04-30

远端校验和不同,因为 tenant set 不同:

shard A = 274002669404fbcd449bdecd929624e3
shard B = 0bb770361058ec76ebc81a2a7d1e2629

协调端本地基线

协调端同时生成完整本地表:

CREATE TABLE shop_ch17.sales_fact (...)
WITH (parallel_workers = 2);

CREATE INDEX sales_fact_tenant_day_idx
ON shop_ch17.sales_fact (
  tenant_id,
  occurred_on,
  account_id
)
INCLUDE (amount, units, channel);

CREATE INDEX sales_fact_day_brin_idx
ON shop_ch17.sales_fact
USING brin (occurred_on)
WITH (pages_per_range = 16);

并创建日汇总 materialized view。于是同一个 PoC 内存在单机对照组,不会拿 分布式结果和一个不存在的 baseline 比较。

外表父表

协调端:

CREATE TABLE shop_ch17.sales_fact_distributed (...)
PARTITION BY LIST (tenant_id);

CREATE FOREIGN TABLE shop_ch17.sales_fact_dist_0
PARTITION OF shop_ch17.sales_fact_distributed
FOR VALUES IN (2, 4, 6, 8)
SERVER pg36_ch17_shard_a
OPTIONS (
  schema_name 'shop_ch17_shard',
  table_name 'sales_fact'
);

CREATE FOREIGN TABLE shop_ch17.sales_fact_dist_1
PARTITION OF shop_ch17.sales_fact_distributed
FOR VALUES IN (1, 3, 5, 7)
SERVER pg36_ch17_shard_b
OPTIONS (
  schema_name 'shop_ch17_shard',
  table_name 'sales_fact'
);

账户维表使用完全相同的 LIST 边界。这让“物理共置但 JOIN 是否下推”成为可测 问题。

为什么不用 HASH 分区

早期 PoC 曾写:

PARTITION BY HASH (tenant_id)
FOR VALUES WITH (MODULUS 2, REMAINDER 0)

而远端 generator 使用 mod(tenant_id, 2)。这造成路由算法不一致。冻结版本 改用 LIST,不是因为 LIST 普遍优于 HASH,而是为了让这个八租户教学 fixture 的物理映射无歧义。

生产 PoC 应使用目标系统真实的 shard function 和 metadata,并加入:

route(key) expected shard
physical rows comply
pruned query returns golden
rebalance changes epoch atomically
old router cannot write after cutover

PoC 应主动寻找反例

本章没有把目的写成“证明 FDW 很快”,而是:

prove pushdown where it happens
prove non-pushdown where it does not

这比只展示成功计划更能检验选型假设。若一个 PoC 从不失败,通常说明验收条件 太宽或只选择了产品最擅长的路径。

17.5.2 记录组件、版本、拓扑和数据分布

版本 manifest

正式证据写入:

validation_path=direct-postgresql-loopback-fdw
server_version=18.6 ...
postgres_fdw=1.2
database=pg36_shop
shard_databases=pg36_shard_a,pg36_shard_b
distribution=explicit-list-by-tenant
pigsty_reference=4.4
pigsty_l1=not-run

此外为实验目录中每个 source file 计算 SHA-256。

为什么正式 fixture 限制 PostgreSQL 18.x:

  • postgres_fdw 行为和功能随版本变化;
  • 本章固定 extension 版本 1.2;
  • SCRAM passthrough、connection inspection 等版本功能不能模糊外推;
  • 计划文本是 18.6 的证据。

概念适用于更多版本,但复制计划前应在目标 major 重新采集。

foreign server

setup.sql 动态读取本地实例:

unix_socket_directories
port

再创建:

CREATE SERVER pg36_ch17_shard_a
FOREIGN DATA WRAPPER postgres_fdw
OPTIONS (
  host '<lab socket>',
  port '<lab port>',
  dbname 'pg36_shard_a',
  fetch_size '10000'
);

fetch_size=10000 是 fixture 基线,不是生产最优值。官方 postgres_fdw 文档说明它控制每次 fetch 取得的行数,server 级设置可被 table 级覆盖。真实 网络要在延迟、内存和结果宽度下测量。

身份映射

每个 server 有三个具名 mapping:

local postgres   -> remote postgres
local pg36_owner -> remote postgres
local pg36_app   -> remote pg36_app

两 server 共六个,没有 PUBLIC mapping。

本地隔离实验使用:

OPTIONS (
  user 'pg36_app',
  password_required 'false'
)

安全边界

password_required=false 只能由 superuser 设置,会允许映射用户利用 PostgreSQL 操作系统账户可获得的认证材料或 trust/peer 关系。官方文档明确 警告不要对 PUBLIC 设置,并要求防止映射用户借机连接成远端 superuser。 本章只在同一实例、Unix-domain socket、受控开发数据库中使用。生产不得 照抄。

PostgreSQL 18 还提供 use_scram_passthrough 选项,但它有严格条件:远端必须 请求 SCRAM,相关节点需有相同 SCRAM secret,传入会话也必须以 SCRAM 认证。 生产应在目标版本中比较:

SCRAM credentials in reviewed secret lifecycle
SCRAM pass-through
GSS delegated credentials
approved certificate/service mechanism

身份与要求见官方 postgres_fdw Connection Options

应用权限

pg36_app 只获得:

USAGE on shop_ch17
SELECT on local, summary, distributed relations/views
USAGE on two foreign servers
remote SELECT on account/sales

不获得 insert/update/delete。负例:

INSERT INTO shop_ch17.sales_fact (...)
VALUES (...);

固定失败:

SQLSTATE 42501

成功读:

SELECT
  count(*) AS sale_count,
  sum(amount)::numeric(18,2) AS amount_total
FROM shop_ch17.sales_fact_distributed
WHERE tenant_id = 3
  AND occurred_on >= DATE '2026-04-01';

输出:

7500,69375.00

catalog 证据

自动采集:

2 foreign servers
6 named user mappings
18 coordinator relations
6 local indexes
7 application privilege facts
6 size facts

关键对象清单:

4 foreign partitions
2 partitioned parents
3 local tables
1 materialized view
2 views
6 indexes

extension schema shop_ch17_ext 不含关系,只有 postgres_fdw 的五个 extension member routines;所有非 extension object 都禁止混入。

Pigsty 路径 A:offline analytics

Pigsty 是 configuration-driven 平台。当前 4.4 文档把集群定义放在:

all.children.<cluster>.hosts

并以 pg_clusterpg_rolepg_seq 等 identity 参数定义实例。一个分析 隔离草图:

all:
  children:
    pg-analytics:
      hosts:
        10.10.10.11:
          pg_seq: 1
          pg_role: primary
        10.10.10.12:
          pg_seq: 2
          pg_role: replica
          pg_offline_query: true
      vars:
        pg_cluster: pg-analytics
        pg_conf: olap.yml

这条路径的假设是:

analysis can tolerate replica lag and read-only semantics
single-node read capacity is sufficient
resource isolation solves primary interference

验收仍需:

  • replica lag 与 freshness;
  • long query/recovery conflict;
  • offline 服务路由;
  • HBA 与只读 role;
  • failover 后标签/服务行为;
  • CPU/I/O 隔离;
  • backup 与升级。

Pigsty Configuration 给出同类 pg-analytics 示例;其 Cluster / Instance 说明 offline instance 与 pg_offline_query 的职责。

Pigsty 路径 B:Citus 评估拓扑

当分片门槛满足后,Pigsty 4.5 可声明 Citus。当前文档要求:

pg_mode: citus
pg_shard: shared horizontal shard name
pg_group: shard cluster number
pg_primary_db: managed Citus database
extra HBA for local/data-node access

简化草图:

all:
  children:
    pg-citus0:
      hosts:
        10.10.20.10: { pg_seq: 1, pg_role: primary }
      vars:
        pg_cluster: pg-citus0
        pg_mode: citus
        pg_shard: pg-citus
        pg_group: 0
    pg-citus1:
      hosts:
        10.10.20.11: { pg_seq: 1, pg_role: primary }
      vars:
        pg_cluster: pg-citus1
        pg_mode: citus
        pg_shard: pg-citus
        pg_group: 1

完整 inventory 还要有 pg_primary_db、database/extension、HBA、凭据以及 全局变量。本书资产为了避免硬编码生产秘密,只保留拓扑骨架。

L1 验收至少覆盖:

package/version on every node
pg_mode/shard/group identity
coordinator/worker metadata
distributed/reference/local tables
single-tenant route and colocated joins
cross-tenant aggregate
coordinator and worker HA
backup/restore
rebalance
rolling/major upgrade
monitoring and alert
security and service routing

本地输出写 pigsty_l1=not-run,所以不能把上述 YAML 称为已部署。

声明不等于状态

Pigsty inventory 是 desired state 的重要来源,但发布证据还要从运行态回读:

inventory commit
rendered config
installed package
pg_extension
pg_settings
Patroni membership
service endpoints
HBA effective rules
shard metadata
backup status
monitoring targets

否则可能出现“YAML 正确,节点尚未收敛”。

17.5.3 不把演示集群的绝对性能外推到生产

loopback 消除了最关键的变量

本章三个数据库共享:

CPU scheduler
shared_buffers
OS page cache
storage
filesystem
socket transport
PostgreSQL installation
host failure

真实多节点新增:

network RTT/throughput/loss
TLS/authentication
independent caches
clock behavior
node skew
replication
DNS/service discovery
firewall
host maintenance

因此不能发布:

two-stage is N times faster
two shards scale linearly
FDW overhead is X ms
Citus will behave like this

本章只发布:

naive plan returns 240000 foreign rows
two-stage plan returns 960 foreign aggregate rows

EXPLAIN rows 也有边界

actual rows 表示某计划节点每 loop 的输出数量。跨网络字节还取决于:

  • row width;
  • text/binary representation;
  • protocol framing;
  • compression;
  • TLS;
  • fetch batches;
  • remote output expressions;
  • retries。

要测网络必须采集 network bytes/packets 与 server/client metrics,不能把 rows 直接乘一个猜测宽度。

人为 planner 设置

教学计划可能使用:

SET max_parallel_workers_per_gather = 2;
SET min_parallel_table_scan_size = 0;
SET parallel_setup_cost = 0;
SET parallel_tuple_cost = 0;
SET enable_seqscan = off;

它们分别用于稳定复现某种路径。强制路径证明“可执行”,不证明 planner 在 真实成本下应选择它,更不证明它在生产更快。

每份强制计划旁都应写:

why forced
what property is proven
what performance claim is not made

小数据掩盖协调成本

24 万行对现代机器很小。它可能:

  • 全在内存;
  • 规划/连接开销占比过高;
  • 看不出网络拥塞;
  • 看不出 shard skew;
  • 看不出 vacuum 和 checkpoint;
  • 看不出 rebalance;
  • 看不出 compaction/backup;
  • 无法设置生产 P99。

正式 PoC 要按 S/M/L/XL 数据阶梯,直到越过 memory 工作集和目标维护窗口。

同机故障探针的含义

把 foreign server port 改为 1,能证明:

partition pruning avoids unopened bad server
global fan-out surfaces connection failure
SQLSTATE is captured
transaction rollback restores catalog

它不能证明:

  • 半开 TCP;
  • DNS 慢失败;
  • packet loss;
  • TLS rotation;
  • remote process crash;
  • node failover;
  • long transaction during disconnect;
  • unknown distributed commit。

生产 failure matrix 要在独立节点执行这些场景。

对比 Pigsty L1 的证据层级

建议区分:

L0 design review
  schema, query, topology, safety, ADR

L1 target environment
  actual packages, nodes, config, connectivity, backup

L2 functional workload
  golden results, plans, permissions, failures

L3 representative performance
  scale, concurrency, cold/warm, resources, cost

L4 operational drills
  restore, failover, rebalance, upgrade, exit

本章 loopback 具有 L0 和部分 L2 证据;Pigsty L1 明确未运行,更没有 L3/L4。

生产 benchmark 的最小补充

[ ] independent hosts/failure domains
[ ] target network/TLS/auth
[ ] target PG/Pigsty/extension versions
[ ] representative data width and skew
[ ] cold/warm/restart runs
[ ] open-loop arrival and bounded concurrency
[ ] OLTP + OLAP mixed workload
[ ] WAL/checkpoint/vacuum/backup overlap
[ ] per-shard and coordinator metrics
[ ] network rows and bytes
[ ] worker/coordinator failure
[ ] backup/restore checksum
[ ] rebalance interruption and resume
[ ] upgrade and rollback
[ ] cost and on-call effort

PoC 的停止规则

如果发生以下任一情况,应停止扩展 demo 并回到设计:

  • golden 不一致;
  • route 与 physical placement 不一致;
  • 必需查询无法局部化;
  • 跨分片事务比例不可接受;
  • mandatory extension 不兼容;
  • backup/restore 无法证明;
  • 身份需要不安全捷径;
  • 失败返回 partial data 却无标识;
  • 生产成本/owner 不明确;
  • 没有退出路线。

本节结论

最小 PoC 的价值不是“跑起来”,而是把关键假设变成:

frozen input
exact topology
observable plan
expected success
expected failure
explicit limitation
repeatable teardown/rebuild

下一节把这些资产串成一次完整执行,并从证据直接生成 ADR,而不是先写结论再 挑选支持它的截图。


上一节:比较分布式候选 · 返回本章目录 · 下一节:实战:从单机证据到选型 ADR · 查看全书目录 · 查看索引中心

17.6 实战:从单机证据到选型 ADR

本节把前五节压成一个可重复的 1.5-proposal

freeze generator and monthly golden
  -> bootstrap two retained shard database shells
  -> build two exact remote schemas
  -> build a local single-node baseline
  -> build LIST-partitioned foreign parents
  -> compare four byte-identical result paths
  -> collect parallel/index/spill/summary plans
  -> collect pruning/pushdown/transfer plans
  -> preserve a non-pushed JOIN counterexample
  -> prove application privilege denial
  -> prove partial shard failure semantics
  -> checksum catalogs and business state
  -> reject unsafe reset attempts
  -> pre-verify three databases
  -> exact per-database reset
  -> rebuild and repeat the entire review

正式证据来自 PostgreSQL 18.6 / postgres_fdw 1.2 的受控本地开发实例。 Pigsty 4.5 的 topology 和职责已经映射,但没有执行 L1:

pigsty_l1=not-run

破坏边界

task.sh all 会删除并重建精确标记的 shop_ch17shop_ch17_ext、 两个 foreign server、六个 user mapping,以及 pg36_shard_a/pg36_shard_b 中的 shop_ch17_shard。两个数据库壳会 保留。只可在本书受控开发 fixture 中执行,禁止在生产运行。

17.6.1 证明一个边界,拒绝一个伪瓶颈

前置连接

沿用第 4 章管理员 service:

[pg36-admin]
host=/path/to/socket-or-host
port=5432
dbname=pg36_shop
user=postgres
chmod 600 /path/to/pg_service.conf
export PGSERVICEFILE=/path/to/pg_service.conf
export PGSERVICE=pg36-admin

密码应放在受控 secret/service 机制中,不出现在命令行、脚本、evidence 或 Git。

环境保护

协调端 context.sql 要求:

database = pg36_shop
writable primary
PostgreSQL major = 18
session_user = postgres superuser
can SET ROLE pg36_owner
pg36_app = constrained LOGIN non-superuser
ch04-v1 physical model exists
postgres_fdw 1.2 available or exact managed state
pg36_shard_a/b database shell identity exact
fdw host/port exactly equal current instance

远端 remote-context.sql 还核对:

expected database name
database owner/comment
shard remainder
shard marker
UTC/ISO session
timeouts

任何已有同名 schema、server、mapping、extension 或 database 身份不符都停止, 不会因名字相同就接管。

资产清单

static/labs/ch17/
├── fixture-manifest.json
├── frozen-monthly.csv
├── fixture.sql
├── fixture-remote.sql
├── bootstrap.sql
├── remote-context.sql
├── remote-setup.sql
├── context.sql
├── setup.sql
├── verify.sql
├── remote-verify.sql
├── final-state.sql
├── reset.sql
├── remote-reset.sql
├── review.py
├── task.sh
├── analytics-distributed-adr.md
├── baseline-v1.5-proposal.json
└── pigsty-declaration.example.yml

另有四份月报导出、十份计划、五份 catalog、安全边界与单分片失败探针。

冻结输入身份

fixture-manifest.json 固定:

frozen-monthly.csv
  rows=32
  sha256=64b045809e10364fd84a587121d919e8562a15335c4c6c015e91a0ead3a44323

fixture.sql
  sha256=3110d0369b0c62fffeee643200f5320d1f6bd26ad5f9950b2ab2b58991080e10

fixture-remote.sql
  sha256=da02dbca8d2cfeb0294c369b15cad3619a9433d984f98866b5043eb4508e3e91

生成器不读取当前时间,不使用随机数。冻结时间只是 fixture metadata,不参与 运行时决定。

单步建立

./static/labs/ch17/task.sh setup

setup 的顺序:

connect maintenance database postgres
  -> create/reuse exact pg36_shard_a/b shells

connect pg36_shard_a
  -> exact remote schema rebuild
  -> load even tenant IDs
  -> analyze

connect pg36_shard_b
  -> exact remote schema rebuild
  -> load odd tenant IDs
  -> analyze

connect pg36_shop
  -> exact FDW/data schema rebuild in one transaction
  -> load full local baseline
  -> create indexes/materialized summary
  -> create LIST foreign partitions/views/grants
  -> analyze
  -> commit
  -> VACUUM ANALYZE local sales fact

协调端 DDL 在单事务内,避免半成品;数据库壳创建和三个数据库的 schema 建立不能放在一个全局事务里。

为什么 setup 后显式 vacuum

覆盖索引:

CREATE INDEX sales_fact_tenant_day_idx
ON shop_ch17.sales_fact (
  tenant_id,
  occurred_on,
  account_id
)
INCLUDE (amount, units, channel);

理论上包含查询所需列,但 Index Only Scan 是否免 heap 访问还取决于 visibility map。新装载表的 page 未必 all-visible。第一轮自动验收曾真实得到:

Bitmap Heap Scan on sales_fact
  -> Bitmap Index Scan on sales_fact_tenant_day_idx

而审查器要求:

Index Only Scan
Heap Fetches: 0

正确修复不是放宽断言,而是在 COMMIT 后显式:

VACUUM (ANALYZE) shop_ch17.sales_fact;

随后两个完整周期都稳定得到:

Index Only Scan using sales_fact_tenant_day_idx
  actual rows=7500
  Heap Fetches: 0

这个经历说明:计划回归不仅依赖 DDL 和数据,也依赖 vacuum/visibility 状态。实验必须显式制造自己的前置条件。

本地事实

fixture-facts.sql

distributed_amount=2256000.00
distributed_sales=240000
distributed_units=1200000
first_day=2026-01-01
last_day=2026-04-30
local_amount=2256000.00
local_sales=240000
local_units=1200000
shard_rows=dist_0:120000,dist_1:120000
summary_rows=2880
summary_sales=240000

远端:

数据库 租户 账户 销售 units amount checksum
pg36_shard_a 2,4,6,8 200 120,000 599,988 1,188,000.00 274002…24e3
pg36_shard_b 1,3,5,7 200 120,000 600,012 1,068,000.00 0bb770…2629

总数相等不够,物理分片 checksum 也必须匹配。

并行边界

psql -X -w \
  --dbname="service=$PGSERVICE" \
  --set=fdw_host=/path/to/socket \
  --set=fdw_port=5432 \
  --file=static/labs/ch17/local-parallel-plan.sql

通常不需要手工运行;task 会动态读取 socket/port 并注入。

冻结关键路径:

Finalize HashAggregate
  -> Gather
       Workers Planned: 2
       Workers Launched: 2
       -> Partial HashAggregate
            -> Parallel Seq Scan on sales_fact
                 actual rows=80000 loops=3

证明:

parallel path exists
two workers actually launched
aggregation decomposes into partial/final
240k rows become 32 groups

不证明:

production speedup
safe system-wide worker count
behavior under concurrent reports

work_mem 边界

低内存:

Sort actual rows=240000
Sort Method: external merge
Disk: about 5920kB
temp read/write present

高内存:

Sort actual rows=240000
Sort Method: quicksort
Memory: about 13645kB
no temp read

两个文件:

这证明 spill 可被定位到一个节点;不授权全局调大 work_mem

汇总边界

raw-aggregate-plan.sql

Parallel Seq Scan on sales_fact
actual rows=80000 loops=3

summary-aggregate-plan.sql

Seq Scan on daily_tenant_summary
actual rows=2880 loops=1

两者输出同一 32 行月报。ADR 因而可以拒绝:

“全局月报每次扫描 24 万行,所以必须分片”

更准确的结论:

这条重复聚合可先用 2,880 行日粒度处理;
生产还要评审刷新、新鲜度、迟到和恢复成本。

BRIN 只记录候选

catalog 验证:

sales_fact_day_brin_idx
  access_method=brin
  operator_class=pg_catalog.date_minmax_ops
  size > 0
  size < covering B-tree in this fixture

没有强制一条 BRIN 查询并宣称更快。小表无法代表物理相关性和 block range 收益。

17.6.2 比较单机加速与一个分布式候选

四条相同结果路径

自动导出:

monthly-local.csv
monthly-summary.csv
monthly-distributed.csv
monthly-two-stage.csv

分别来自:

任务逐个:

cmp static/labs/ch17/frozen-monthly.csv \
    evidence/.../monthly-local.csv

四个文件都要 byte-identical,不做“行数相等即可”的弱比较。

单租户裁剪

tenant-pruned-plan.sql

SELECT count(*), sum(amount)
FROM shop_ch17.sales_fact_distributed
WHERE tenant_id = 3
  AND occurred_on >= DATE '2026-04-01';

关键计划:

Foreign Scan on shop_ch17.sales_fact_dist_1
  actual rows=7500
  Remote SQL:
    SELECT amount
    FROM shop_ch17_shard.sales_fact
    WHERE occurred_on >= '2026-04-01'
      AND tenant_id = 3

审查器明确要求:

sales_fact_dist_1 present
sales_fact_dist_0 absent
tenant/date in Remote SQL

这同时证明 partition pruning、column projection 和 filter pushdown。

朴素全局聚合

distributed-naive-plan.sql 直接查询 分区父表:

HashAggregate
  -> Append actual rows=240000
       -> Foreign Scan dist_1 actual rows=120000
            Remote SQL: SELECT tenant_id, occurred_on, units, amount ...
       -> Foreign Scan dist_0 actual rows=120000
            Remote SQL: SELECT tenant_id, occurred_on, units, amount ...

远端没有 GROUP BY,协调端获得 24 万条事实再聚合。

这不表示 postgres_fdw 永远不能下推 aggregate;它说明这条父表查询在本次 目标版本和 SQL 形状下没有得到期望的跨分片预聚合。

两阶段聚合

distributed-two-stage-plan.sql 显式分别查询两个 foreign partition:

SELECT tenant_id, occurred_on,
       count(*), sum(units), sum(amount)
FROM shop_ch17.sales_fact_dist_0
GROUP BY tenant_id, occurred_on

UNION ALL

SELECT tenant_id, occurred_on,
       count(*), sum(units), sum(amount)
FROM shop_ch17.sales_fact_dist_1
GROUP BY tenant_id, occurred_on;

外层再按月合并。计划:

GroupAggregate actual rows=32
  -> Sort actual rows=960
       -> Append actual rows=960
            -> Foreign Scan Aggregate dist_0 actual rows=480
                 Remote SQL ... GROUP BY 1, 2
            -> Foreign Scan Aggregate dist_1 actual rows=480
                 Remote SQL ... GROUP BY 1, 2

数据流:

960240000=0.004 \frac{960}{240000} = 0.004

即返回行数为朴素路径的 0.4%。这个比例只描述冻结 fixture 的行数,不能直接 转成“性能提升 250 倍”。

同分片 JOIN 反例

collocated-parent-plan.sql

SELECT
  account.segment,
  count(*) AS sale_count,
  sum(sale.amount) AS amount_total
FROM shop_ch17.sales_fact_distributed AS sale
JOIN shop_ch17.account_dim_distributed AS account
  ON account.tenant_id = sale.tenant_id
 AND account.account_id = sale.account_id
WHERE sale.tenant_id = 3
  AND sale.occurred_on >= DATE '2026-04-01'
GROUP BY account.segment;

物理上两表的 tenant 3 都在 shard B。实测:

Hash Join on coordinator actual rows=7500
  -> Foreign Scan sales_fact_dist_1 actual rows=7500
  -> Hash
       -> Foreign Scan account_dim_dist_1 actual rows=50

两条 Remote SQL 都没有 JOIN。这个反例写进 ADR:

对目标抽象层而言,共置只是设计前提,不是下推证据。必须检查目标产品、 版本、SQL 与 EXPLAIN VERBOSE

若生产候选是 Citus,应在真实 Citus L1 上创建目标 distributed/reference tables,重新验证 colocated join;不能用 FDW 反例替代,也不能假定一定成功。

应用身份

成功读:

sale_count=7500
amount_total=69375.00

写入负例:

exit=3
SQLSTATE 42501

catalog:

app_schema_usage=true
app_local_sales_select=true
app_local_sales_write=false
app_distributed_sales_select=true
app_distributed_sales_write=false
app_server_a_usage=true
app_server_b_usage=true

六个 mapping 都是具名;password_required=false 被 review 当成必须显式存在 的 lab-only 警告,而不是悄悄依赖环境。

单分片失败

执行器先查询 tenant 2:

healthy_shard_tenant_2=30000

再执行全局 count,固定:

exit=3
SQLSTATE 08001

任务同时捕获 stdout/stderr;如果只看 stderr,就无法证明健康分片路径曾成功。

错误断开使事务回滚后,重新导出:

server-catalog-after-failure.csv

并与故障前 server-catalog.csv 逐字节比较,确保 shard B port 没有残留为 1。

一次完整 evaluate

如果不需要 reset/rebuild 双周期:

evidence_dir="$PWD/evidence/ch17/evaluate-$(date -u +%Y%m%dT%H%M%SZ)"

PG36_EVIDENCE_DIR="$evidence_dir" \
  ./static/labs/ch17/task.sh evaluate

evaluate 会重建一次、采集所有证据并运行 review。

仅验证现有数据库:

PG36_EVIDENCE_DIR="$PWD/evidence/ch17/verify" \
  ./static/labs/ch17/task.sh verify

它分别运行 coordinator、shard A、shard B 的完整数据库内断言。

review 的职责

review.py 不连数据库,只审查 evidence:

manifest version/target/checksum
source generator and frozen CSV hashes
four byte-identical monthly exports
local/remote cardinality and checksums
server/mapping/relation/index/security/size catalogs
parallel/index/spill/summary plans
tenant pruning and Remote SQL
naive/two-stage row shape
non-pushed JOIN counterexample
application read/write
shard failure and restored server catalog
final state
coordinator/remote verify outputs
baseline ADR contract

这使 evidence 可以离线审阅,也避免“数据库后来变了,旧报告仍假装当前”。

17.6.3 输出 ADR、PoC 证据、生产代价和撤退路线

ADR 不以产品名开头

analytics-distributed-adr.md 先写背景和决策顺序:

1. fix correctness, SQL, statistics, paths
2. prove parallel/index/BRIN/spill/summary on one node
3. assess offline replica for tolerable-staleness reads
4. enter distribution only after measured resource boundary
5. evaluate Citus when PostgreSQL-compatible sharding fits
6. compare specialized OLAP only when PostgreSQL paths miss SLO

postgres_fdw 的定位是 mechanics/counterexample lab,不是生产性能结论。

ADR 的冻结证据

问题 证据 决策影响
可并行? 2 workers launched 单机还有并行路径
选择性查询? index-only,heap fetch 0 先修访问路径
内存? 64kB 外排 / 32MB 内排 按并发设局部预算
重复聚合? 240k vs 2,880 输入 先评估汇总
单租户路由? 只访问 shard B tenant_id 可局部
朴素全局? 240k foreign rows 协调端/网络风险
两阶段? 960 aggregate rows 计算靠近数据
同分片 JOIN? coordinator Hash Join 必须实测下推
shard 失败? scoped read 成功/global 08001 定义 partial semantics

决策与限制同版本

baseline-v1.5-proposal.json 固定:

target versions
default path
distribution key
routing warning
remote aggregation design
production Citus gate
fixture contracts
expected checksums
lab authentication warning
evidence inventory
rollback contract
limitations

canonical JSON SHA-256:

3dcb7308cf6983122ee860ad3dc2a4b44651549e3d5631770839bb9a0be450c6

改变 ADR contract 会改变 release candidate identity。

进入生产 PoC 前的代价

ADR 要预算:

schema/query changes for distribution key
backfill and dual-write/read
coordinator/worker nodes
HA and service routing
network/TLS/auth
backup repository and restore
monitoring/alerting
rebalance capacity
rolling/major upgrade
on-call training
license/cloud cost
exit rehearsal

不能只比较机器数量。

Pigsty 交付分支

pigsty-declaration.example.yml 保留两个候选:

Path A:
  pg-analytics primary + replica with pg_offline_query
  pg_conf: olap.yml

Path B:
  pg-citus0/1/2 groups
  pg_mode: citus
  pg_shard / pg_group

它们不是同一 inventory 的叠加方案,也不是完整生产配置。ADR 应先决定测哪条 假设,再补齐目标地址、database、users、extensions、HBA、secret、backup、 service 和 HA。

撤退路线

PoC 撤退:

stop new lab sessions
verify coordinator + A + B exact state
drop coordinator views/matview/tables
drop mappings/servers/extension
drop remote tables/schemas
retain empty database shells
verify zero remaining managed objects

生产迁移撤退则应提前设计:

canonical source of truth remains PostgreSQL
dual-read compares checksums
route change has version/epoch
old and new writers cannot both own same key
backfill has watermark and resume point
last reversible point is named
service/DNS rollback is tested
new system data can be exported back

reset 的三道 guard

协调 reset 需要:

PG36_RESET_TOKEN=RESET_CH17_ANALYTICS_FDW_LAB
PG36_RESET_TARGET=pg36_shop/shop_ch17+shop_ch17_ext+fdw

错误 action token:

SQLSTATE P3660
reset refused: invalid ch17 action token

错误 target:

SQLSTATE P3661
reset refused: invalid ch17 target token

存在 application_name LIKE 'pg36-ch17-%' 的其他 worker:

SQLSTATE P3663
reset refused: ch17 workers are active

task 会主动启动一个 sleep worker,观察到 PID 后证明 P3663,再 cancel。

精确 reset

协调端顺序:

DROP VIEW ...
DROP MATERIALIZED VIEW ...
DROP TABLE exact parents/local tables ...
DROP SCHEMA shop_ch17;

DROP USER MAPPING ...  -- six exact mappings
DROP SERVER ...        -- two exact servers
DROP EXTENSION postgres_fdw;
DROP SCHEMA shop_ch17_ext;

不使用 CASCADE。输出:

status=coordinator-reset-ok
remaining_data_schema=0
remaining_extension_schema=0
remaining_ch17_servers=0
retained_shard_databases=pg36_shard_a,pg36_shard_b

每个 remote:

full remote verify
BEGIN
drop sales_fact
drop account_dim
drop fixture_meta
drop schema
COMMIT

输出:

status=remote-reset-ok
remaining_schema=0
database_shell=retained

跨库非原子边界

任务先对三个数据库全部 preflight verify,再按:

coordinator -> shard A -> shard B

退出。每一步各自事务化,但整体不是一个事务。若 A reset 后 B 失败,系统处于 部分退出状态;恢复方式是按 exact identity 继续补偿或完整重建。

这项 limitation 同时写入 lab contract、baseline 和 ADR,不能用脚本“看起来 是一条命令”掩盖。

手工 reset

只有明确需要退出 fixture 时:

export PG36_RESET_TOKEN=RESET_CH17_ANALYTICS_FDW_LAB
export PG36_RESET_TARGET=pg36_shop/shop_ch17+shop_ch17_ext+fdw

PG36_EVIDENCE_DIR="$PWD/evidence/ch17/reset" \
  ./static/labs/ch17/task.sh reset

它会同时处理两个远端 schema。不要在生产、共享开发数据库或身份未知的目标 上执行。

17.6.4 验收采用 checklist:evidence

完整双周期

evidence_dir="$PWD/evidence/ch17/$(date -u +%Y%m%dT%H%M%SZ)"

PG36_EVIDENCE_DIR="$evidence_dir" \
  ./static/labs/ch17/task.sh all

成功输出:

status=ok
fixture=frozen-byte-identical-four-paths
single_node=parallel+index+summary+spill
distributed=tenant-pruning+fdw+two-stage
counterexamples=hash-is-not-modulo+join-not-pushed
failure=healthy-shard-read+global-08001
guards=P3660+P3661+P3663
postgres_fdw=1.2
pigsty_l1=not-run
release_candidate_checksum=3dcb7308cf6983122ee860ad3dc2a4b44651549e3d5631770839bb9a0be450c6

目录:

evidence/ch17/<run>/
├── cycle-1/
├── reset-wrong-token.*
├── reset-wrong-target.*
├── reset-active-worker.*
├── reset-exact/
└── cycle-2/

checklist:evidence

验收项 evidence 通过条件
manifest manifest.txt PG18、FDW 1.2、三库、checksums
source manifest + fixture JSON generator/golden SHA 精确
golden 四份 monthly CSV 与 frozen byte-identical
local-facts fixture-facts.csv 240k、1.2m、2.256m
remote-facts remote-*-state.csv 各 120k + shard checksum
parallel local-parallel-plan.txt planned/launched=2、partial/final
covering-index selective-index-plan.txt 7,500、Heap Fetches 0
spill low/high plans external merge vs quicksort
summary raw/summary plans 240k vs 2,880 input shape
pruning tenant-pruned-plan.txt only dist_1 + remote filters
naive naive plan 240k foreign rows
two-stage two-stage plan 480+480 remote aggregate rows
join-counterexample collocated plan coordinator Hash Join
auth mapping/security catalogs six named mappings, read-only app
app-failure app-write.* exit 3, SQLSTATE 42501
shard-failure shard-failure.* healthy=30k, global 08001
catalog-rollback two server catalogs byte-identical
database-verify three verify outputs coordinator + A + B status ok
final-state final-state.csv release/rows/checksums exact
review review.txt status=review-ok
reset-guards root failure files P3660/P3661/P3663
exact-reset reset outputs schemas/servers zero, DB shells retained
rebuild cycle-2 same complete review

最终状态

final-state.sql 固定:

business_checksum=42fb8ab5444469eba1f104a8e1e529dd
distributed_sales=240000
fixture=ch17-analytics-v1
local_sales=240000
monthly_checksum=644d45544ebbc2a80c42270c38ac6885
naive_transfer_rows=240000
postgres_fdw=1.2
release=1.5-proposal
shard_rows=dist_0:120000,dist_1:120000
summary_rows=2880
tenant3_april=7500:69375.00
two_stage_transfer_rows=960

naive_transfer_rowstwo_stage_transfer_rows 是与计划合同共同审查的固定 事实;如果 SQL 或版本改变,不能只保留硬编码值,必须重新采集 plan 并更新 proposal。

review 输出

status=review-ok
fixture=frozen-byte-identical-four-paths
single_node=parallel+index+summary+spill
distributed=tenant-pruning+fdw+two-stage
counterexamples=hash-is-not-modulo+join-not-pushed
failure=healthy-shard-read+global-08001
business_checksum=42fb8ab5444469eba1f104a8e1e529dd
monthly_checksum=644d45544ebbc2a80c42270c38ac6885
release_candidate_checksum=3dcb7308cf6983122ee860ad3dc2a4b44651549e3d5631770839bb9a0be450c6

失败排查顺序

如果 task 失败:

  1. 保留 evidence,不立即重跑覆盖;
  2. 找到最后产生的 stderr;
  3. 判断是环境 guard、业务 checksum、计划 shape、权限还是故障边界;
  4. 连接三个数据库分别运行 verify;
  5. 检查 server catalog 是否已回滚;
  6. 修复原因,不降低断言掩盖差异;
  7. 新建 evidence 目录完整重跑两个周期;
  8. 比较两个 manifest。

第一轮 index-only 失败就是这个流程的例子:证据显示 Bitmap Heap Scan, 根因是 visibility map precondition,没有把 reviewer 改成接受任意 index 路径。

生产发布仍缺什么

本地 status=ok 之后仍需:

[ ] representative production-scale data and skew
[ ] independent nodes and failure domains
[ ] target network/TLS/identity
[ ] Pigsty L1 inventory and runtime convergence
[ ] actual Citus or selected candidate, not FDW surrogate
[ ] OLTP+OLAP mixed concurrency
[ ] P50/P95/P99 and open-loop throughput
[ ] CPU/memory/I/O/temp/WAL/network cost
[ ] coordinator/worker HA
[ ] backup-to-empty restore checksum
[ ] rebalance interrupt/resume
[ ] schema change mixed-version test
[ ] rolling/major upgrade
[ ] RPO/RTO and partial-result semantics
[ ] migration dual-read and last rollback point
[ ] exit rehearsal

因此 1.5-proposal 是设计与本地机制候选,不是生产批准。

本章最终决策

在冻结工作负载上,当前可支持的结论是:

  1. 单机仍有并行、覆盖索引、spill 治理和汇总空间;
  2. 不能因 24 万行原表扫描直接宣布需要分布式;
  3. tenant_id 对单租户访问有良好局部性;
  4. 分布式全局聚合必须关注计算位置与协调端输入;
  5. 同分片不自动证明 JOIN 下推;
  6. 路由算法不一致会导致静默错误,HASH remainder 不能当整数取模;
  7. 部分失败和跨数据库退出必须成为业务/运维合同;
  8. 下一层应先比较 Pigsty offline replica;若代表性容量仍越界,再在真实 Pigsty Citus L1 上验证分布键、HA、恢复和再平衡;
  9. 专用 OLAP 只有在 PostgreSQL 路径无法满足已定义 SLO,且团队接受第二套 数据管道与值班体系时进入终选。

这就是一份合格选型 ADR 的语气:它不承诺某产品必胜,而是清楚说明当前证据 允许做什么、禁止外推什么,以及什么新证据会触发下一次决策。


上一节:部署最小分布式 PoC · 返回本章目录 · 下一章:万法归宗:PostgreSQL 数据平台与替代边界 · 查看全书目录 · 查看索引中心