跳至内容
17 合纵连横:分析加速与分布式选型

第 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 B120,000 / 120,000 行
本地日汇总2,880 行
最终月报32 行
租户 3 的 4 月7,500 笔,69,375.00
朴素 FDW 返回协调端240,000 行
两阶段聚合返回协调端960 行
业务校验和42fb8ab5444469eba1f104a8e1e529dd
月报校验和644d45544ebbc2a80c42270c38ac6885

这里的“返回行数”描述数据流形状,不是网络字节,也不是耗时。三个数据库都在 一台机器的一个 PostgreSQL 18.4 实例中,没有独立 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.4 文档要求 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 数据平台与替代边界 · 查看全书目录 · 查看索引中心

最后更新于