跳至内容

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:

查询/事务频率候选键单分片?跨分片代价
租户 dashboardtenant_id
租户账户与销售 JOINtenant_id
全局月报tenant_idpartial/final
跨租户对账tenant_idfan-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

监测:

[ \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 兼容”标签把不同系统 压成同一类。


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

最后更新于