Citus 分布式分片与水平扩展

Citus 分布式 PostgreSQL 实战:coordinator / worker 架构、分布式表与参考表、分片键选择与共置(colocation)、查询路由与下推、分片再平衡与扩容、高可用与故障恢复、以及 Citus 的适用边界与常见限制。

单机 PostgreSQL 的垂直扩展有明确上限:一台 64 核 512GB 的机器能撑起大多数业务,但当数据量到 TB 级、写入 QPS 到十万级时,继续加配置的性价比会急剧下降。水平扩展(Scale Out) 的路径是把数据分片(Sharding)到多台机器,让每台只处理一部分数据与查询。Citus 是这条路径上最成熟的开源方案——它不是一个独立的数据库,而是 PostgreSQL 的一个扩展,把标准的 PostgreSQL 变成分布式数据库,应用仍然用 SQL 与 PostgreSQL 驱动,只是背后多了协调与并行下推。

Citus 的价值在于「对应用近乎透明」:大部分单节点 SQL 在分布式表上继续可用,代价是分片键的选择会深刻影响性能,且一部分操作(跨分片 JOIN、跨分片唯一约束)有明确限制。本文从架构、分片模型、共置、查询路由讲到扩容与高可用,最后给出适用边界。

一、Citus 架构

Citus 集群由两类节点组成:

角色职责
Coordinator接收客户端连接、解析 SQL、拆分并下发给 worker、汇聚结果
Worker存储实际分片,执行下推的子查询,返回部分结果

Coordinator 本身不存储业务数据(除元数据与参考表的副本),所有分布式表的数据都在 worker 上。一条查询到达 coordinator 后:

  1. 解析 SQL,判断涉及哪些分片。
  2. 生成「分片查询(shard query)」并下发给相关 worker。
  3. 并行收集各 worker 的结果。
  4. 在 coordinator 做最终聚合、排序、LIMIT。

这个过程对应用不可见,但可以通过 EXPLAIN 观察下推情况。

安装 Citus 分两种形态:一种是单节点 Citus(coordinator 与 worker 合一,用于开发与中小规模),另一种是多节点集群。单节点形态适合先验证 SQL 兼容性,再平滑迁移到多节点。

-- 在 coordinator 上添加 worker
SELECT citus_add_node('worker-1', 5432);
SELECT citus_add_node('worker-2', 5432);

-- 查看集群拓扑
SELECT * FROM citus_get_active_worker_nodes();

生产上通常把 coordinator 与 worker 部署在 Kubernetes 上,用 StatefulSet 保证稳定的网络标识与持久卷。节点编排与滚动升级的细节可参考 Kubernetes 集群运维 。

1.1 单节点快速验证

在决定是否采用 Citus 之前,先用单节点镜像验证 SQL 兼容性,成本最低:

docker run -d --name citus \
  -p 5432:5432 \
  -e POSTGRES_PASSWORD=secret \
  citusdata/citus:12.1

psql -h localhost -U postgres -c "SELECT citus_version();"

单节点形态下,coordinator 与唯一的 worker 在同一进程里,create_distributed_table 依然可用,分片只是被放在同一实例的不同 schema 里。它足以验证:分片键是否让查询路由到单分片、JOIN 是否能共置、哪些 SQL 无法下推。确认无误后再迁移到多节点,迁移的只是数据,SQL 基本不用改。

1.2 节点角色与元数据

Citus 的集群拓扑记录在 coordinator 的元数据表里:

SELECT nodeid, nodename, nodeport, noderole, isactive
FROM pg_dist_node;

noderole 取 primary / secondary;isactive = false 的节点不会被分配新分片。所有分片的位置记录在 pg_dist_shard_placement:

SELECT shardid, shardstate, nodename, nodeport
FROM pg_dist_shard_placement
WHERE shardid = 102008;

shardstate = 1 表示正常,3 表示正在迁移。理解这两张元数据表是排查「分片找不到」「查询报 node unreachable」类问题的前提。

二、分片模型:三种表

Citus 把表分成三类,选择哪一类决定了数据如何分布。

2.1 分布式表(Distributed Table)

按分片键(distribution column)把行哈希到 N 个分片(shard),每个分片是一个普通 PostgreSQL 表,落在某个 worker 上。

CREATE TABLE orders (
    id         bigint,
    tenant_id  int NOT NULL,
    amount     numeric(12,2),
    created_at timestamptz NOT NULL,
    PRIMARY KEY (tenant_id, id)
);

SELECT create_distributed_table('orders', 'tenant_id');

默认分片数是 32(citus.shard_count)。分片数一旦确定,后期调整需要重新分片,因此要在建表时估算好未来规模。经验值:分片数 = worker 数 × 每 worker 期望的分片数(通常 4~8),让每个 worker 承载多个分片以便再平衡时能搬动。

2.2 参考表(Reference Table)

整张表在每个 worker 上都保留一份完整副本,用于「小维表 + 大事实表」的 JOIN。

CREATE TABLE countries (
    code text PRIMARY KEY,
    name text NOT NULL
);

SELECT create_reference_table('countries');

参考表的写操作会广播到所有 worker(通过两阶段提交保证一致性),因此只适合读多写少的维表。它让「事实表 JOIN 维表」可以在 worker 本地完成,避免跨节点拉数据。

2.3 本地表(Local Table)

不调用 create_distributed_table 的表就是本地表,只存在于 coordinator。它们不参与分布式查询的下推,任何 JOIN 本地表的分布式查询都要把数据拉到 coordinator 处理,代价高。因此业务表应尽量建成分布式表或参考表,本地表只用于运维元数据。

三类表的对照:

类型数据分布写代价适用场景
分布式表按分片键哈希到 N 个分片单分片写(路由)或广播大事实表、业务主表
参考表每 worker 一份完整副本广播 + 两阶段提交读多写少的维表
本地表仅在 coordinator本地写运维元数据、临时表

一个常见的性能陷阱是「把大表建成参考表」:参考表的每次写都要广播到所有 worker 并走两阶段提交,写入吞吐会随 worker 数下降。参考表只该放「小且几乎不变」的维表,如国家码、币种、状态字典。

三、分片键选择与共置

分片键(distribution column)是 Citus 最重要的设计决策。选错了,性能可能比单机还差。

3.1 选择原则

  • 高基数:分片键的取值要足够多,避免数据倾斜。tenant_id 通常是好选择,status(只有几个值)是灾难。
  • 出现在大多数查询的 WHERE 里:这样查询能路由到单个分片(router query),避免广播到所有分片。
  • 用于表间 JOIN:两张表用同一个分片键可以「共置」,JOIN 在 worker 本地完成。

3.2 共置(Colocation)

共置是指分片键相同、分片数相同的表,其对应分片落在同一个 worker 上。这样两张表的分片级 JOIN 就在本地进行:

SELECT create_distributed_table('orders',    'tenant_id');
SELECT create_distributed_table('order_items', 'tenant_id');

orders 与 order_items 都用 tenant_id 分片,且分片数相同(默认 32),Citus 会把它们放进同一个共置组(colocation group)。验证共置:

SELECT logicalrelid, colocationid, shardcount
FROM pg_dist_partition p
JOIN pg_dist_colocation c USING (colocationid);

共置是 Citus 性能的核心。如果 orders 用 tenant_id 分片而 order_items 用 order_id 分片,两张表无法共置,JOIN 退化为「重分区(repartition)」——运行时会按 JOIN 键重新洗牌数据,开销巨大。

3.3 重新分片与共置组变更

如果分片键选错了,Citus 提供在线重分片:

SELECT citus_rebalance_start();

但重分片是重量级操作,会占用大量 IO 与网络。更好的做法是设计阶段就用真实查询模式验证分片键:把生产查询的 WHERE 与 JOIN 条件列出来,确认分片键在其中高频出现。

四、查询路由与下推

Citus 把查询分成三种执行模式,性能差异巨大。

4.1 Router 查询(单分片)

当 WHERE 中带有分片键的等值条件时,coordinator 知道数据在哪个分片,直接路由过去:

SELECT * FROM orders WHERE tenant_id = 42 AND id = 1001;

这条查询只访问一个分片,延迟与单机 PostgreSQL 几乎一致。这是 Citus 上最理想、也最该被鼓励的查询形态。用 EXPLAIN 验证:

Custom Scan (Citus Adaptive)
  Task Count: 1
  Tasks Shown: All
  ->  Task
        Node: host=worker-1 port=5432
        ->  Index Scan using orders_pkey_102008 on orders_102008

Task Count: 1 就是单分片路由。

4.2 多分片并行查询

当查询没有分片键谓词,但可并行时,coordinator 把查询下发给所有分片并行执行,再汇聚:

SELECT count(*) FROM orders WHERE created_at >= '2026-01-01';
Custom Scan (Citus Adaptive)
  Task Count: 32
  ->  Task (32 tasks)

Task Count: 32 意味着 32 个分片全部参与。这类查询的收益来自并行,但每个分片都要扫一遍,总工作量随分片数放大。它适合大范围聚合,不适合高频点查。

4.3 跨分片 JOIN

若 JOIN 的两侧表未共置,或 JOIN 条件不含分片键,Citus 需要「重分区」或在 coordinator 做归并:

-- orders 按 tenant_id 分片,users 按 id 分片,JOIN 键 user_id 不含任一分片键
SELECT o.id, u.name
FROM orders o JOIN users u ON u.id = o.user_id;

这类 JOIN 会触发 repartition join:把两张表按 JOIN 键重新哈希分片,代价是两次全量数据洗牌。避免方法:

  • 让 JOIN 表共置(同一分片键)。
  • 把小表改成参考表。
  • 把 JOIN 键也纳入分片键设计。

判断 JOIN 是否被下推,看 EXPLAIN 里是否出现 Task Count: N 且每个 Task 内部是完整的 JOIN,而不是 coordinator 顶层的 Hash Join。

4.4 下推的限制

并非所有 SQL 都能下推。以下操作只能在 coordinator 执行,代价高:

  • ORDER BY + LIMIT 的组合:各 worker 先各取一批,coordinator 再全局排序取前 N(除非 LIMIT 很小且有分片键)。
  • 窗口函数跨分片:row_number() OVER (...) 若分区范围跨分片,需要 coordinator 全局处理。
  • 跨分片的 DISTINCT、GROUPING SETS。
  • 递归 CTE、FOR UPDATE 跨分片。

这些限制的根源是「分片间没有全局的排序/事务视图」。设计查询时应主动规避,把它们拆成可下推的形式。

4.5 判断下推是否发生

EXPLAIN 是唯一可靠的判据。看 Task Count 与 Task 内部的节点类型:

现象含义处置
Task Count: 1单分片路由理想,无需处理
Task Count: N(N = 分片总数),Task 内是完整聚合分片级并行可接受,注意总工作量
顶层 Hash Join / Sort 在 coordinator未下推检查 JOIN 键与分片键是否一致
Task Count: N 且出现 repartition重分区 JOIN改为共置或参考表
Subplan 里有 Fetch跨节点拉数据数据分布不匹配,需重构

一个实用手法是给关键查询加 EXPLAIN (ANALYZE, VERBOSE),在输出里搜索 host=worker 的行数——行数越少,说明越多的计算在 worker 本地完成,coordinator 越轻。

五、扩容与再平衡

5.1 加节点

Citus 扩容分两步:加节点、搬分片。

SELECT citus_add_node('worker-3', 5432);
SELECT citus_rebalance_start();

citus_rebalance_start() 会计算分片在节点间的理想分布,然后逐个移动分片(默认在后台异步进行),过程中保持集群可读写。可以通过下面的视图监控进度:

SELECT * FROM citus_rebalance_status();

关键参数:

参数作用
rebalance_strategy再平衡策略(by_disk_size / by_shard_count)
shard_transfer_modeauto / block_writes / force_logical
parallel_transfer_count并发迁移的分片数

shard_transfer_mode = block_writes 会短暂阻塞写,但迁移更快、更安全;force_logical 用逻辑复制迁移,几乎不阻塞写,但需要主键与 REPLICA IDENTITY。大批量迁移的可靠性可以借鉴 逻辑复制实战 里的经验。

5.2 缩容与隔离

把某个 worker 下线前,必须先把它的分片迁走:

-- 阻止在该节点上放置新分片
SELECT citus_set_node_property('worker-3', 5432, 'shouldhaveshards', false);
-- 再平衡,把分片移出
SELECT citus_rebalance_start();
-- 确认无分片后再摘除节点
SELECT citus_remove_node('worker-3', 5432);

直接 citus_remove_node 而不先迁分片,会丢失数据。

5.3 分片数量与再平衡成本

再平衡的最小单位是分片,不是行。分片数越多,再平衡越灵活(可以细粒度搬动),但每个分片的元数据开销也越大。Citus 官方建议每 worker 保持 4~8 个分片。分片过少(如 4 个 worker 却只有 4 个分片),加节点后无法均匀分配;分片过多(数千个),citus_rebalance 会变得非常慢。

六、高可用与运维

6.1 每个节点的 HA

Citus 本身不提供分片级复制——它把高可用交给底层 PostgreSQL 的主从复制。每个 worker 是一个标准 PostgreSQL 实例,用流复制 + Patroni / repmgr 做故障切换。coordinator 同理,但 coordinator 挂掉的代价更大(整个集群不可用),因此必须配 HA。

Citus 提供 citus_set_coordinator_host 与 coordinator 元数据备份机制,用于协调器故障恢复。完整的 HA 方案设计与切换演练见 PostgreSQL 高可用架构 。

6.2 连接与连接池

每个客户端连接在 coordinator 上会占用一个后端进程,而 coordinator 向每个 worker 建立连接来下发分片查询。连接数会放大:100 个客户端连接 × 32 个分片并发,可能瞬间需要大量 worker 连接。因此 Citus 集群几乎必须配连接池:

  • 客户端到 coordinator:用 pgBouncer 的事务级池化,压平应用侧连接。
  • coordinator 到 worker:Citus 内部维护连接池,参数 citus.max_cached_conns_per_worker 控制每 worker 缓存的连接数。
citus.max_cached_conns_per_worker = 8
citus.max_shared_pool_size = 200

6.3 备份与恢复

分布式备份要保证所有节点的快照一致。Citus 提供 citus_create_restore_point 在 coordinator 与所有 worker 上同时创建命名还原点:

SELECT citus_create_restore_point('before-upgrade');

之后对每个节点分别做基础备份(pg_basebackup),恢复时从同一还原点回放。分片元数据(pg_dist_* 表)必须一起备份,否则恢复后无法定位分片。

6.4 监控

Citus 在标准 pg_stat_* 之外提供了集群视图:

-- 各节点的分片分布与大小
SELECT nodename, count(*) AS shards,
       pg_size_pretty(sum(shard_size)) AS total_size
FROM citus_shards
GROUP BY nodename;

-- 当前正在执行的分布式查询
SELECT * FROM citus_stat_activity;

-- 分片级统计
SELECT * FROM citus_shard_sizes;

重点关注:分片分布是否均衡(某节点分片数或体积畸高)、是否有大量跨分片查询(Task Count 长期等于总分片数)、节点间数据倾斜。

七、限制与适用边界

Citus 不是万能的。以下场景需要谨慎:

限制说明
跨分片唯一约束唯一约束必须包含分片键
跨分片外键仅支持共置表之间的外键
跨分片事务用两阶段提交,有性能开销
跨分片 JOIN需重分区,代价高
部分 SQL 特性递归 CTE、跨分片窗口函数受限
复杂 DDLALTER TABLE 会广播到所有分片

Citus 最适合的形态是多租户 SaaS(Multi-tenant SaaS):所有查询都带 tenant_id,天然路由到单分片,表间按 tenant_id 共置,JOIN 全部本地化。这类负载在 Citus 上几乎线性扩展。

它不适合:跨租户的复杂分析查询为主、需要跨分片唯一性、查询模式高度随机无法选分片键的场景。这些场景应转向列存(如 ClickHouse )或按查询模式专门设计的数仓。

7.1 Citus 与应用层分库分表对比

很多人会纠结「用 Citus 还是自己在应用层分库分表」。两者对照如下:

维度Citus应用层分库分表
对应用透明高(SQL 基本不改)低(需路由中间件/框架)
跨分片 JOIN支持(代价高)基本不支持
事务两阶段提交通常无跨分片事务
运维复杂度中(多一类节点)高(分片逻辑在应用里)
扩容citus_rebalance_start()手工迁移 + 改路由规则

结论:只要能用 SQL 表达的业务,优先 Citus;只有在「分片键与业务强耦合、且完全不接受跨分片查询」时才考虑应用层分库分表。Citus 把复杂性从应用代码搬进了数据库,长期维护成本更低。

7.2 迁移与版本升级

从单机 PostgreSQL 迁到 Citus 的路径:先用 pg_dump --schema-only 在 coordinator 建好表结构,再 create_distributed_table,最后用 COPY 把数据加载进分布式表——加载时会自动按分片键分发到各 worker。Citus 自身的升级(扩展版本)需要在 coordinator 与所有 worker 上同步升级扩展,并用 citus_create_restore_point 先建还原点。升级前务必在预发环境完整演练一遍,因为扩展升级失败可能导致元数据不一致。

小结

Citus 用扩展的方式把 PostgreSQL 变成分布式数据库,核心是「分片 + 共置 + 下推」。落地时的决策顺序是:先确定分片键(高基数、高频出现在 WHERE 与 JOIN 中),再让相关表共置,再保证查询能路由到单分片。扩容靠加节点 + 再平衡,高可用靠底层 PostgreSQL 的流复制,连接数要靠连接池压平。它的甜蜜区是多租户 SaaS 的点查与租户内 JOIN;一旦查询模式以跨分片分析为主,就要重新评估是否该用 Citus,或者配合列存做分析层。

继续阅读

探索更多技术文章

浏览归档,发现更多关于系统设计、工具链和工程实践的内容。

全部文章 返回首页

「database」更多文章

  1. PostgreSQL 锁与阻塞分析
  2. COPY 与批量数据加载优化
  3. pgvector 向量检索与混合查询