
1. 为什么要用Citus给PostgreSQL做分布式先聊点实际的。单机PostgreSQL跑得好好的为什么要折腾分布式我见过太多团队业务涨到一定量级后数据库开始出现三个典型症状单表数据量破亿后即使建了索引查询延迟也开始变得不稳定写入吞吐上不去主库CPU常年打满从库只能扛读写路径始终是单点最棘手的是随着数据增长VACUUM、索引重建、备份恢复这些运维操作的时间窗口越来越长动不动就要锁表。这时候一般有两条路一是做应用层分库分表自己维护路由规则和跨库查询工作量大而且一旦路由键设计失误后面改起来要命二是引入NewSQL或者分布式数据库但迁移成本高团队还得重新学习一套新体系。Citus走的是一条中间路线它不是一个全新的数据库而是PostgreSQL的一个扩展直接把单机PostgreSQL升级成分布式集群。装上Citus之后你的数据库对外表现仍然是一个PostgreSQL连接方式一样SQL语法几乎不用改psql照常用ORM框架照常连PG的生态工具备份、监控、迁移基本都能沿用。区别在于数据会被自动分片到多个worker节点上查询会被协调器节点自动改写、下推、汇总。这篇博文就按照我从零搭建一套Citus集群的完整过程来写包括环境规划、安装配置、分片设计、数据迁移、性能调优以及实际运维中踩过的坑和排查思路。如果你正在考虑把PostgreSQL从单机引向分布式或者已经在用Citus但觉得性能没达到预期这篇文章可以给你一个完整的参考路径。我用的版本是PostgreSQL 15 Citus 12.1这是一个相对稳定且匹配度较好的组合。操作系统是Ubuntu 22.04 LTS。后面所有操作步骤我都按这个环境来展开版本差异带来的坑我会单独标注。2. 动手之前先搞清楚Citus的工作原理2.1 协调器与工作节点各干什么活Citus集群里有两类角色协调器Coordinator和工作节点Worker。协调器就是你在psql里连接的那个节点对外看起来像一个普通的PostgreSQL实例。当用户发起一个查询时协调器会做三件事解析SQL根据元数据判断涉及哪些分片生成查询计划把能下推的部分下推到对应worker节点执行汇总结果把各worker返回的部分结果做聚合、排序、关联然后返回给客户端。工作节点负责实际存储数据分片并执行查询。数据不是均匀乱撒的而是根据分布列distribution column的哈希值按分片数量取模映射到不同的分片shard再按照分片亲和性shard affinity分配到各个worker节点上。每个分片在物理上就是worker节点上的一个普通PostgreSQL表表名形如orders_102008。这个设计意味着Citus并没有发明一套新的存储引擎它复用的是PostgreSQL本身的存储、索引、MVCC机制。所以PostgreSQL的行存储、TOAST、WAL日志、流复制这些底层能力在Citus里依然有效只是多了一层分布式的调度和路由逻辑。2.2 数据怎么分片查询怎么下推Citus支持两种数据分布方式哈希分布和参考表。哈希分布是最核心的方式。创建分布式表时指定分布列Citus会根据分布列的值计算哈希决定数据落到哪个分片。关键在于查询下推如果查询带了分布列的等值过滤条件比如订单表按user_id分布查询WHERE user_id 123协调器可以直接定位到那个具体的分片去执行这个叫单分片查询性能接近单机如果不带分布列条件比如查最近一周的订单协调器就得把查询广播到所有分片各查各的再汇总这个叫分布式查询性能取决于分片数和worker数量。参考表则是把数据量小、经常需要关联的表复制到所有worker节点上每个节点都保存一份完整的副本。典型例子是省份表、类别表、配置表。关联时Citus可以在本地完成JOIN避免跨节点数据传输。这里有一个新手比较容易误解的地方Citus并不是让所有查询都变快它擅长的是“分布列点查聚合下推”这种模式。如果业务查询几乎不带分布列过滤又需要频繁跨分片JOIN那Citus的收益会大打折扣。所以后面章节我会详细讲怎么选分布列这个决策直接决定集群的性能上限。2.3 为什么是Citus而不是其他方案对比应用层分库分表Citus最大的优势是透明性不需要改应用代码不需要维护路由中间件分布式细节被封装在扩展内部。对比Greenplum这类分析型MPP数据库Citus更偏向OLTP轻量OLAP混合场景它支持完整的PostgreSQL事务语义包括ACID、外键有约束条件、索引、触发器而Greenplum在事务处理上有明显短板。选型时我建议自己对照这张表做评估对比维度Citus应用层分库分表Greenplum对应用侵入性低PG生态无缝衔接高需改造DAO层低事务支持完整支持跨节点分布式事务2PC通常牺牲跨库事务较弱实时写入能力强适合高并发写入中弱批量导入为主查询类型适配点查聚合下推视路由设计而定复杂分析为主运维复杂度低一个扩展搞定高需维护中间件较高3. 集群规划与版本选型3.1 节点规划三节点起步还是两节点够用Citus没有强制要求必须多少个节点最小拓扑是一个协调器加一个worker就能跑起来这在测试环境没问题。生产环境我建议至少两到三个worker节点这样分片可以分布到多个节点上发挥并行能力同时避免协调器单点故障——注意协调器本身也可以做流复制备库但Citus的协调器自动故障转移需要额外的运维手段MySQL的MHA那套思路在Citus里并不直接适用。我这次搭建用的是三个节点192.168.10.10 — coordinator192.168.10.11 — worker1192.168.10.12 — worker2三台机器都是4核8G的虚拟机磁盘SSD。生产环境建议至少8核32G起步因为Citus节点本质还是PostgreSQL内存越充裕shared_buffers和work_mem可以给得越大排序和哈希操作的性能越好。3.2 版本组合PostgreSQL 15 Citus 12.1版本选择是很多人忽略但至关重要的一步。Citus是紧跟PostgreSQL主版本迭代的不同版本的Citus对应的PostgreSQL版本范围不同。我选的是PostgreSQL 15 Citus 12.1理由有三点PG 15的public模式权限变更更合理逻辑复制增强pg_stat_statements等插件集成度更高。Citus 12.x对PG 15支持是官方明确列出的稳定组合。同时Citus 12开始分片修剪和分布式查询优化器成熟度很高很多之前需要手动hint的查询现在能自动生成不错的执行计划。提示不要盲目追求最新版本。PG 16甚至17配合最新Citus虽然已经可用但其行为和插件兼容性可能需要社区验证更多时间。我身边有同事在生产上用PG 16 Citus踩过扩展编译问题。稳定优先。3.3 网络与防火墙的准备工作Citus的coordinator和worker之间需要互相通信走的是PostgreSQL协议默认端口5432。集群内部需要在防火墙放行5432端口并且建议只允许集群内网IP访问。另外coordinator需要能SSH连接到各个worker部分管理命令工具会用到也一并放行。一个比较容易踩的坑如果启用了SELinux或者AppArmor可能需要为PostgreSQL进程单独配置网络访问策略。Ubuntu默认的AppArmor对PostgreSQL没有额外限制一般不需要管。4. 从零开始搭建Citus集群4.1 在所有节点上安装PostgreSQL 15三个节点都需要安装PostgreSQL服务和Citus扩展操作一致。Ubuntu 22.04默认软件源里带的PostgreSQL版本是14要装15需要先添加官方APT源。sudo apt update sudo apt install -y postgresql-common sudo /usr/share/postgresql-common/pgdg/apt.postgresql.org.sh -y sudo apt update sudo apt install -y postgresql-15 postgresql-client-15安装完成后PostgreSQL服务会自动启动数据库数据目录在/var/lib/postgresql/15/main。先把服务停掉后面改配置再启动。sudo systemctl stop postgresql4.2 安装Citus扩展Citus官方APT源通常已经包含在postgresql-common脚本配置的源里直接安装sudo apt install -y postgresql-15-citus-12.1装完后检查一下扩展文件是否到位ls /usr/lib/postgresql/15/lib/citus.so ls /usr/share/postgresql/15/extension/citus.control看到这两个文件就说明Citus已经安装成功接下来需要在每个节点的postgresql.conf里加载Citus库并设置相关参数。4.3 配置节点的关键参数编辑/etc/postgresql/15/main/postgresql.conf在所有节点上做三块配置。第一块是加载Citusshared_preload_libraries citus第二块是通用的PostgreSQL性能参数注意这里是初调后面的性能调优章节还会细调max_connections 200 shared_buffers 2GB effective_cache_size 6GB work_mem 16MB maintenance_work_mem 256MB checkpoint_completion_target 0.9 wal_buffers 16MB第三块是Citus特有参数coordinator和worker侧略有差异。我在所有节点上统一先配置citus.node_conninfo sslmodeprefer citus.max_worker_nodes_per_ha 32第三块里的citus.node_conninfo表示节点间通信的连接参数。如果PG启用了SSL这里需要配置对应的SSL选项没启用就用默认即可。4.4 启动集群并建立节点关系所有节点上的PostgreSQL服务启动sudo systemctl start postgresql sudo systemctl enable postgresql然后切换到postgres用户进入psql先在coordinator上创建扩展sudo -u postgres psqlCREATE EXTENSION citus;接下来在coordinator上录入worker节点。这个命令Citus会自动把元数据写入协调器的citus相关系统表并且初始化与worker节点的连接SELECT citus_add_node(192.168.10.11, 5432); SELECT citus_add_node(192.168.10.12, 5432);查看集群节点状态SELECT * FROM citus_get_active_worker_nodes();如果返回两行记录主节点列表里出现这两个worker的IP和端口集群就已经建立成功了。这里要说一句后期如果要把自动故障转移和节点管理交给专门的工具可以考虑在协调器上启用citus.enable_metadata_sync、citus.replicate_reference_tables_on_activate这些参数默认值在大多数场景下是合理的。4.5 验证集群基本功能先建一个简单的分布式表验证CREATE TABLE test_dist ( id bigint, val text ); SELECT create_distributed_table(test_dist, id);然后插入数据并查询INSERT INTO test_dist SELECT g, value- || g FROM generate_series(1, 10000) g; SELECT count(*) FROM test_dist; SELECT * FROM test_dist WHERE id 12345;通过EXPLAIN查看执行计划如果出现类似Distributed Query、Custom Scan (Citus)这样的节点说明查询确实走了分布式执行。5. 数据建模与分布列设计5.1 分布列选错后面全是坑Citus性能好坏的命门就在分布列的选择上。分布列决定数据如何跨节点分布也决定哪些查询能够下推成单分片操作。原则一分布列应该是业务中最常见的高频等值查询条件。以订单表为例如果业务大部分查询是“查某个用户的订单”那么user_id是理想的分布列。如果按order_id分布user_id过滤就无法裁剪分片每次都要跑全节点。原则二分布列应该保证数据足够分散避免数据倾斜。比如按gender分布就非常糟糕男/女两个桶数据只落在少数分片上。如果业务只有几千个大客户但需要按客户维度做FK关联这种场景需要评估是否需要重新建模。原则三需要JOIN的大表应该尽量用同一个分布列。Citus处理JOIN时如果两边表都按相同的分布列分布那么可以在本地节点完成匹配不需要shuffle否则协调器会执行重分布repartition join数据传输开销很大。我见过一个真实的案例有人把店铺的订单表按store_id分片又把用户表按user_id分片结果两边做JOIN的时候Citus只能把所有表数据重分布一遍8个worker跑一个JOIN查询耗时比单机还慢两三倍。5.2 参考表的使用场景与注意事项参考表适合数据量小几万行以内变化不频繁需要和其他分布式表频繁JOIN。比如商品类目、区域配置、枚举字典。创建参考表很简单CREATE TABLE regions ( id int primary key, name text ); SELECT create_reference_table(regions);参考表会在每个worker节点上保留完整副本。新增节点时Citus会自动同步参考表数据。需要注意参考表如果有外键指向其他参考表没问题但如果引用分布式表需要验证业务约束是否允许。参考表的数据更新会触发所有节点同步频繁写入参考表会成为性能瓶颈所以它只适合低频更新的维表不适合存放订单、日志这类持续写入的数据。5.3 外键、唯一约束在分布式环境下的限制Citus对分布式表的约束支持是有取舍的主键或唯一约束必须包含分布列。也就是说如果表按user_id分布唯一索引可以是(user_id, order_id)不能只建order_id唯一索引。这是为了保证唯一性检查可以在单个分片内完成不需要跨节点通信。分布式表之间的外键要求外键列必须等于被引用表的分布列。也就是说如果orders按user_id分布payments也要按user_id分布然后外键orders(user_id) REFERENCES users(user_id)才被允许。分布式表与参考表之间可以建外键参考表作为被引用方没有限制。这个限制让很多从单机迁移过来的应用需要调整建表语句。建议迁移前先梳理现有表结构把不符合条件的约束找出来评估是调整业务逻辑还是放弃某些约束。不要把外键满天飞的单机模型直接搬过来Citus不是这么用的。6. 数据迁移与分布式化改造6.1 从单机PostgreSQL数据导出常见的迁移方式是用pg_dump导出数据。这里要特别注意先迁移表结构创建分布式表再导入数据顺序不能反。如果在普通表里先导入大量数据再调用create_distributed_table()Citus需要额外的数据重分布过程耗时长且占用大量临时空间。正确流程在源库导出schema不导出数据pg_dump -h 源库IP -U postgres -d 业务库 --schema-only -f schema.sql在Citus集群的coordinator上执行schema.sql建出普通表。分析哪些表适合做分布式表哪些做参考表执行相应的转换命令SELECT create_distributed_table(orders, user_id); SELECT create_reference_table(regions);数据导入。如果源库数据量在GB级别直接用pg_dump --data-only导出然后再psql导入coordinator即可。数据量更大时建议用\copy按表分批导入或者使用pg_dump --formatcustom加上-j参数并行导出。6.2 利用copy命令高效灌入数据对于大表COPY是最高效的导入方式。建议在源库先把数据按分片键排序导出排序后的数据在Citus灌入时能减少分片切换的随机IO提升不少速度。psql -h 源库 -c \copy (SELECT * FROM orders ORDER BY user_id) TO /tmp/orders.csv WITH CSV HEADER然后在coordinator上执行\copy orders FROM /tmp/orders.csv WITH CSV如果数据量上百GB单线程COPY可能比较慢可以考虑按时间或者按分布列范围拆分文件并行执行多个\copy会话。但要注意并发COPY时coordinator和worker的CPU/IO都会升高建议观察负载逐步加并发。6.3 应用层的改造点数据迁过去了应用侧也不是完全零改动。主要有三个地方要注意连接串改为连接coordinator的地址和端口而不是原来的单机PG。如果做了读写分离需要额外配置负载均衡或故障转移方案。高频写入的业务插入语句最好显式指定分布列不要让数据库自己去推断。虽然Citus支持从主键推断分布列但显式指定查询计划更可控。事务涉及跨节点操作时要重新评估。Citus支持分布式事务通过两阶段提交协议实现但跨节点的写事务性能远低于单节点事务。如果应用里大量使用跨分片的大型事务性能会明显下降。我在项目中遇到过一个批量导入逻辑在单机PG跑没问题上Citus后发现每个事务要写多个分片延迟从几十毫秒涨到几秒。后来改成按分布列分批提交才恢复正常。7. 性能调优从参数到SQL的完整链路7.1 分片数量怎么定从规划到实测验证分片数由citus.shard_count控制默认是32。分片数不是越大越好也不是越小越好。先按规则估算一个基准值每个worker节点的分片数建议在2到4之间。也就是说3个worker节点的集群初始分片数在6到12之间。分片数太少会导致单分片过大数据不均衡太多会增加协调器维护分片元数据的内存开销某些查询计划生成也会变慢。我这次3 worker的集群设为citus.shard_count 12每个worker平均4个分片。创建分布式表时显式指定SET citus.shard_count 12; SELECT create_distributed_table(orders, user_id);验证分片分布情况SELECT nodename, nodeporter, count(*) FROM citus_shards GROUP BY 1, 2;如果发现某个worker的分片数明显多说明分片分配不够均匀。可以重建表或使用citus_move_shard()调整。注意citus_move_shard()是重活要在低峰期操作。7.2 协调器与worker的参数差异化配置协调器负责查询路由和结果汇总它本身也存储元数据表。worker节点承担真正的数据存储和查询执行。两者的参数调优侧重点不同。worker节点上重点配置shared_buffers给到物理内存的25%比如32G内存给8G。effective_cache_size给到物理内存的60%到70%帮助优化器判断是否走索引扫描。work_mem在内存有余量的前提下可以适当加大。排序、哈希关联都在这个内存里做默认4MB在分析型查询下远远不够。我这边给到64MB但要注意work_mem是每个排序操作单独分配连接数高时会翻倍消耗内存不要无脑调大。max_parallel_workers_per_gather设为2到4让并行查询能跑起来。后台worker总预算max_parallel_workers也要相应调整。协调器节点上的共享内存参数可以适当低于worker因为大表数据不在本地但shared_buffers也不能过小因为元数据查询和结果汇总也消耗资源。7.3 查询下推优化避免昂贵的重分布JOIN再看一遍执行计划凡是出现Repartition Join的SQL都值得审查。重分布JOIN是最费资源的操作协调器把两张表数据按JOIN键全部重哈希分发数据在节点间网络传输集群再大也会被拖垮。优化手段主要集中在三方面第一检查两个JOIN表的分布列是否一致。不一致的分析业务是否可以把两表调整到同一个分布列。比如orders和order_items本来就按order_id关联那就都按order_id分布。第二小表指定为参考表。如果一个表只有几千行与其做大表重分布JOIN不如直接建参考表让每个节点都有本地副本JOIN就在本地完成。第三把高频查询条件下推。SQL条件里尽量带上分布列等值过滤。比如查询订单明细业务上如果能带上user_id条件即使order_id是主键也建议把user_id一起放进WHERE让Citus能精确裁剪分片。7.4 大宽表与聚合查询的调优建议Citus在OLAP场景下的优势在于能够把聚合计算下推。比如SELECT count(*), date_trunc(day, created_at) FROM orders GROUP BY 1Citus会把每个分片都算一遍局部聚合再汇总到协调器。这类查询的调优重点排序和聚合最好在worker节点完成协调器只做轻量汇总。如果涉及高基数的GROUP BYwork_mem要足够大否则worker排序阶段会溢写到磁盘拖慢查询。避免在SELECT列表里写SELECT *然后丢给协调器过滤最理想的是在worker节点完成过滤和投影减少网络传输量。对经常做聚合分析的大表可以创建物化视图并由cron定时刷新。Citus对物化视图的支持比较好但要注意物化视图如果建在分布表上刷新时每条记录都要经过协调器调度耗时较长。7.5 实测一次从250ms到18ms的调优过程说一个具体的优化案例。我负责的一个订单查询接口单机PG只需要30ms左右上了Citus集群后竟然要250ms。用EXPLAIN ANALYZE定位后发现两个严重问题SQL里带了两个JOIN其中一张小表没有建参考表导致全节点重分布JOIN查询条件里只带了order_id没带user_idCitus无法裁剪分片。优化动作如下先看执行计划的关键部分EXPLAIN (ANALYZE, BUFFERS) SELECT o.id, o.total_amount, u.name FROM orders o JOIN users u ON u.id o.user_id WHERE o.id 1008611;执行计划中显示Repartition Join说明两个表没有使用同一个分布键。随后我做了两处调整第一把users表改为参考表它是典型的维表量小且经常联合查询。第二把查询条件改为SELECT o.id, o.total_amount, u.name FROM orders o JOIN users u ON u.id o.user_id WHERE o.user_id 123456 AND o.id 1008611;应用层在查询时想办法从会话上下文里带上user_id。改完后执行计划变成了Router Query只访问一个分片延迟降到18ms。这个案例说明Citus不是装上就完事分布列选择、查询写法、执行计划分析这三件事是持续投入才能拿回报的。8. 常见问题与排查技巧实录8.1 问题排查的常规路径Citus集群出问题我的排查顺序是先确认是不是数据分布问题再看执行计划是否下推成功然后检查节点间的网络开销最后考虑参数和资源瓶颈。几个常用SQL-- 查看分片在节点上的分布是否均衡 SELECT nodename, count(*) FROM citus_shards GROUP BY 1; -- 查看当前正在运行的分布式查询 SELECT * FROM citus_stat_activity WHERE query ILIKE %SELECT%; -- 查看分片修剪是否生效执行计划里有没有Shard pruning字样 EXPLAIN SELECT * FROM orders WHERE user_id 123;8.2 数据倾斜为什分片ID特别大有时候你会看到某个分片体积异常大其他分片都很小。多半是分布列的选择导致哈希碰撞集中或者业务数据本身分布极不均匀。排查方法SELECT nodeport, shardid, shard_size FROM citus_shards ORDER BY shard_size DESC;如果倾斜严重两个思路一是更换分布列换更高基数的字段二是改用citus.rebalance_strategy中的策略配合citus_move_shard()手动迁移独立分片。这里要注意Citus 12的自动再平衡可能需要pro版本才开放全部能力社区版建议手动干预。8.3 连接数被打满coordinator很快没响应Citus集群里coordinator需要同时维护与客户端和worker的PostgreSQL连接。后端连接数很容易被占满表现为FATAL: sorry, too many clients already。解决思路有几层提高max_connections只是表面解法协调器内存会被堆满。应用侧使用连接池比如PgBouncer把客户端连接收敛。注意PgBouncer要配置transaction模式避免会话状态跨事务残留。用citus.max_client_connections控制协调节点能够使用的后端连接数上限避免worker连接不够用。我实际踩过没有连接池时100个应用实例每个建10条连接直接打爆了coordinator。上了PgBouncer事务模式后后端连接稳定在40左右集群立刻稳了。8.4 修改分片数要谨慎垃圾数据清理要注意有些人在建表时没考虑好分片数后面想改。Citus社区版不支持直接修改已存在分布式表的分片数正确做法是新建一张表设置好citus.shard_count把数据INSERT INTO ... SELECT或者\copy过去再删掉旧表改指向新表的视图或者应用连接。还有一点容易忽略Citus的分布式表删除后worker节点上的物理分片可能不会立即删除干净。用DROP TABLE之后检查一下worker上是否还有*_xxxxx的残留表如果有需要手动清理。我在一次释放存储空间时发现由于之前删表和重建过于频繁多个worker上残留了几百个无主分片白白占用了几十GB磁盘。8.5 VACUUM与膨胀问题在分布式环境下的表现Citus每个分片都是一个独立的表VACUUM的工作量成倍增加。高频更新的表建议把autovacuum调积极一点比如把autovacuum_vacuum_scale_factor从默认的0.2改小到0.05。Citus 12对VACUUM做了并行化可以同时清理多个分片但coordinator上执行VACUUM时仍然可能收到部分分片锁冲突的错误日志。遇到这种问题时建议在低峰期对分布表执行VACUUM (ANALYZE, PARALLEL 4) orders;因为Citus会把VACUUM操作下推到各个worker让它并行执行可以有效控制表膨胀。8.6 节点宕机与恢复worker节点宕机coordinator上的查询会出现在线分片不可用错误。应对思路数据安全性上建议每个worker开启流复制备库Citus层面不做副本靠PG原生的备库来保数据。一旦worker宕机恢复重启后确认citus扩展加载正常再用SELECT citus_activate_node(worker_ip, 5432);重新激活节点。如果某个分片确实损坏先把该分片的元数据置为不可用SELECT citus_mark_not_available(worker_ip, 5432);优先保住集群整体可用再恢复数据。8.7 快速参考常见问题速查表症状可能原因排查命令/动作查询延迟高JOIN未下推分片裁剪失效EXPLAIN查看是否出现Repartition JOIN或广播查询单分片数据过大分片数不足查询citus_shards看分布评估重建表数据分布不均分布列哈希碰撞更换分布列或手动move shard连接耗尽并发连接数过高引入连接池控制max_connectionsINSERT超时跨节点事务过多检查事务是否只涉及单客户端序列改写分批提交节点宕机后查询失败分片不可用citus_mark_not_available后恢复节点9. 对这套集群的一些复盘和心得做完整套从零搭建到调优的过程我的感受是Citus把分布式数据库的门槛降低了很多但它的天花板其实是你的模型设计水平。分布列选对了查询按着分布列走集群可以很轻松扩展选错了或者应用层随意写SQL跑出来的性能甚至不如单机。最后分享一个我后期运维时总结的习惯每次新接入一个大表都先跑一遍EXPLAIN确认所有高频查询都走到了单分片路由Router Query或至少是下推执行Distributed Query而不是重分布JOIN。然后过一周再看一次分片大小分布。数据分布会随着业务变化而变化定期体检比等出问题再救火踏实得多。