跳转至

MPP 数据分布策略 —— 从 Hash 分段到一致性哈希的工程实践

作者:JiangChong | 撰写时间:2026年04月

适用场景框:当你需要理解为什么 Vertica 在 projection 级别而非表级别定义分段键、为什么一张表可以有多份不同分段的副本、以及如何为不同查询模式选择最合适的分布策略时,这篇文章适合你。

开篇声明:本文以 Vertica 为主要剖析对象,但在所有环节对比其他 MPP 系统(Greenplum、ClickHouse、Doris/StarRocks、Snowflake、Redshift)的不同实现。「MPP 共性」与「Vertica 专属」在文中明确区分——共性机制适用于所有类似架构的分布式数据库,Vertica 专属机制会标注适用范围。

关联文章

理解全文脉络

本文从「MPP 为什么需要数据分布」开始,先用跨系统对比表建立全局视角(第 1 节),然后按 MPP 共性框架的三层结构——分布粒度、分布算法、扩缩容迁移——逐层深入各系统的不同实现,其中 Vertica 作为深度剖析的载体(第 2 节),再对比 table-level 分布与 projection-level 分段的 trade-off 及跨系统设计哲学差异(第 3 节)。如果你关注的是「我的分段设计对不对」这种实操问题,第 4 节和第 5 节提供了从诊断到修复的完整闭环;第 6 节将原理浓缩为 9 条可执行的设计原则(标注了【通用】/【Vertica】适用范围)。


1. 问题背景 — MPP 中数据分布为什么是一个难题

1.1 单机数据库没有分布问题

在单机数据库中,数据分布不需要设计——所有数据都在同一台机器的磁盘上。查询只需要考虑「怎么算最快」,不需要考虑「数据在哪」。

但在 MPP 数据库中,数据分布在多台机器的本地磁盘上,网络取代磁盘 I/O 成为新的瓶颈。一个查询能否并行执行、JOIN 是否需要跨节点传输数据、GROUP BY 是否能在本地完成——这些问题的答案都取决于一个设计决策:数据按什么规则分布到各节点

1.2 数据分布决定一切

考虑一个最简场景:orders 表(10 亿行)和 customers 表(100 万行),集群 4 个节点。

如果你将 ordersorder_id 哈希分布到 4 个节点,customerscustomer_id 分布到 4 个节点,那么执行 orders JOIN customers ON orders.customer_id = customers.customer_id 时:

  • order_id = 100 的订单在 node1,但对应的客户在 node3
  • node1 无法在本地完成 JOIN——需要从 node3 获取客户数据,或者把自己的订单发给 node3

数据分布策略直接决定了查询执行策略。如果两张表按同一 JOIN 键分段,JOIN 完全本地执行——零网络开销。如果分段键不同,就需要 Resegment 或 Broadcast——网络开销取决于数据量。

不仅如此,数据分布还影响:

  • GROUP BY 效率:如果数据已按 GROUP BY 键分段,聚合可以在每个节点独立完成
  • 扩容操作:节点数变化时,数据如何重新分布?迁移量有多大?
  • 容错能力:一个节点宕机后,它的数据能否从其他节点恢复?

1.3 各 MPP 系统的应对策略一览

面对同样的「多节点并行计算需要数据分布」问题,不同 MPP 系统在分布粒度分布算法上做出了截然不同的选择:

系统 分布粒度 分布算法 分布键变更 扩缩容迁移 设计哲学
Vertica 投影级(projection-level) Ring-style hash + local segments 创建新投影(物理副本) ELASTIC_CLUSTER(segment 整块搬迁) 物理冗余换查询灵活,优化器按查询选投影
Greenplum 表级(DISTRIBUTED BY Hash modulo N ALTER TABLE SET DISTRIBUTED BY(全表重写) gpexpand(全表重分布) 简洁直观,一张表一种分布
Redshift 表级(DISTKEY / DISTSTYLE Hash modulo N 重建表 自动(RAC 集群自动重分布) 托管服务,用户选择少但运维负担低
ClickHouse 表级(Distributed 引擎 sharding_key half-md5 / murmur 重建分布式表 手动 ALTER TABLE RESHARD 本地表+分布式表的双层抽象
Doris / StarRocks 表级(DISTRIBUTED BY HASH(col) BUCKETS N murmur_hash3_32 ALTER TABLE SET DISTRIBUTION(后台异步) 自动分桶迁移 前 100 强 OLAP,开箱即用的分布策略
Snowflake 无用户可见分布键 自动 micro-partition clustering 无需变更(自动重组) 零运维(存储计算分离) 数据库完全托管,用户完全不用关心分布

本文选取 Vertica(投影级分段)、Greenplum(经典表级分布)、ClickHouse(分片键+双层抽象)、Snowflake(无分布键)四个代表性系统做深度对比——它们分别代表了「灵活性优先」「简单性优先」「双层抽象」「完全托管」四种不同的设计哲学。

1.4 传统 MPP 的做法:表级分布键

传统 MPP 数据库(如 Greenplum、Teradata)在表级别定义分布键(distribution key)。CREATE TABLE 时指定 DISTRIBUTED BY (col),之后这张表的所有数据都按这个键分布。

这是最简单直观的方式——每张表选一个分布键,数据按它分布,查询按它并行。但它的根本局限是:一张表只能有一种分布方式

事实表 orders 需要同时 JOIN customers(按 customer_id)、products(按 product_id)、dates(按 date_id)。表级分布只能优化其中的一个 JOIN——另外两个 JOIN 必然需要 Resegment 或 Broadcast。

这就是 projection-level 分段想要解决的核心矛盾:当一张表有多个维度的 JOIN 需求时,如何让每个 JOIN 都能本地执行?


2. 核心概念与机制 — MPP 数据分布的三个层次

MPP 数据分布可以按三个层次理解:分布粒度(在哪一级别定义分布——表级还是更细粒度)、分布算法(如何将行映射到节点——hash%N、ring-style、murmur 等)、扩缩容迁移(节点变化时数据如何重分布——全量重写还是增量搬迁)。这三个层次相互独立——同为表级分布,Greenplum 用 hash%N 而 ClickHouse 用 half-md5;同为 ring-style,但只有 Vertica 在投影级应用了它,且配套了 local segments 的高效迁移机制。

2.1 层次一:分布粒度 — 在哪一级别定义数据分布

MPP 共性: 所有 MPP 系统都需要回答同一个问题——数据按什么规则分布到各节点。这个规则定义在哪一级别,直接决定了你能为多少种查询模式优化数据分布。

系统 分布粒度 定义方式 一张表能有几种分布
Vertica 投影级 CREATE PROJECTION ... SEGMENTED BY HASH(col) 多种(每个投影可选不同分段键)
Greenplum 表级 CREATE TABLE ... DISTRIBUTED BY (col) 一种
Redshift 表级 CREATE TABLE ... DISTKEY(col) / DISTSTYLE KEY 一种
ClickHouse 表级(本地表) CREATE TABLE ... ENGINE = MergeTree ORDER BY col,再由 Distributed 引擎指定分片键 一种(但分布式表只是路由层,可重建)
Doris / StarRocks 表级 CREATE TABLE ... DISTRIBUTED BY HASH(col) BUCKETS N 一种
Snowflake 无用户可见粒度 自动 micro-partition clustering 零(用户完全不关心分布)

【Vertica 独有】 投影级分段(projection-level segmentation)是 Vertica 区别于所有其他 MPP 系统的核心设计。其他系统中数据分布是建表时的一次性决策,Vertica 将它变成了投影级别的多策略组合——这是「物理冗余换查询灵活」的根本来源。

以 Vertica 为例:投影级分段的运作方式。同一张 orders 表可以有三份物理副本,每份按不同的键分段:

-- Super projection:按 order_id 分段,用于主键查询和加载
CREATE PROJECTION orders_super
AS SELECT * FROM orders
ORDER BY order_date
SEGMENTED BY HASH(order_id) ALL NODES KSAFE 1;

-- 窄投影 1:按 customer_id 分段,用于客户维度的 JOIN
CREATE PROJECTION orders_by_cust
AS SELECT order_id, customer_id, order_date, amount
FROM orders
ORDER BY customer_id
SEGMENTED BY HASH(customer_id) ALL NODES KSAFE 1;

-- 窄投影 2:按 product_id 分段,用于商品维度的 JOIN
CREATE PROJECTION orders_by_prod
AS SELECT order_id, product_id, order_date, amount
FROM orders
ORDER BY product_id
SEGMENTED BY HASH(product_id) ALL NODES KSAFE 1;

优化器根据查询的 JOIN 条件自动选择最合适的投影:

  • 查询 JOIN customers ON customer_id → 选 orders_by_cust,Local Join,零网络
  • 查询 JOIN products ON product_id → 选 orders_by_prod,Local Join,零网络
  • 查询只需主键扫描 → 选 orders_super

对比传统 MPP(以 Greenplum 为例):

CREATE TABLE orders DISTRIBUTED BY HASH(order_id);
-- orders 的所有数据只按 order_id 分布
-- 只有 order_id JOIN 能本地执行
-- customer_id JOIN 和 product_id JOIN → 必然需要 Redistribute Motion

(来源:C-Store 7 Years §3.1, §3.6;verticaOptimizer §III-A)


2.2 层次二:分布算法 — 如何将行映射到节点

MPP 共性: Hash 分段是所有 MPP 数据库最主流的数据分布方式。以 4 节点集群为例,基本流程:

  1. 对每行的分段键计算哈希值,得到一个整数(通常 32 或 64 位)
  2. 将哈希空间映射到各节点——具体的映射算法因系统而异
  3. 哈希值落在哪个区间,数据就存入哪个节点

比喻:把一副扑克牌按花色分给 4 个人——红心给第一个人,黑桃给第二个人,以此类推。每张牌的去向是确定的、互不重叠的。

Hash 分段有两个关键优点:确定性(每行数据的位置由哈希函数唯一确定,不需要查询目录)和均匀性(只要分段键基数足够高且分布均匀,数据在各节点间接近平均分配)。

各系统的分布算法差异:

系统 分布算法 哈希函数 关键特征
Vertica Ring-style hash 内部哈希(64 位空间等分) 映射规则与节点数解耦,扩容时只需调整区间边界(详见 §2.3)
Greenplum 哈希映射表 hash(col) → segment mapping 简洁直观,依赖节点总数
Redshift 哈希映射表(推测) 内部哈希 AWS 托管,用户无需关心算法细节
ClickHouse half-md5 / murmur halfMD5(col) % NmurmurHash2_64(col) % N 用户可选哈希函数
Doris / StarRocks murmur_hash3_32 murmur_hash3_32(col) % buckets 固定 bucket 数,节点变化时分桶重新分配

(来源:C-Store 7 Years §3.6;ClickHouse 文档 sharding_key 章节;Doris 文档 DISTRIBUTED BY HASH 章节)

以 Vertica 为例:ring-style hash 的映射原理。

hash % N 的根本缺陷在于映射规则依赖节点总数 N。节点数一变,所有行的归属全部重算——hash % 4 = 0 的行在 hash % 6 下有 50% 概率换节点。这不是「迁移效率」的问题,而是「映射规则本身不稳定」的问题。

ring-style 的核心思路:将映射规则与节点总数解耦。 不是对 N 取模,而是把整个 64 位哈希空间切成 N 个连续区间,每个节点负责一个区间。以 $C_{MAX} = 300$ 简化,4 节点时:

Node 1:   0 ≤ hash <  75
Node 2:  75 ≤ hash < 150
Node 3: 150 ≤ hash < 225
Node 4: 225 ≤ hash < 300

每行的归属由「哈希值落在哪个区间」唯一确定——不需要知道总共有几个节点。当集群从 4 节点扩到 6 节点时,区间边界会调整(详见 §2.3),但哈希值与区间的映射逻辑本身不变——原来 hash = 50 的行仍然属于 [0, 75) 这个区间,只是该区间可能被切分后部分归属新节点。映射规则和迁移操作是正交的两件事。

对比 hash%N 和 ring-style 的映射差异:

特性 hash % N Ring-style
映射规则 hash(col) % N hash ∈ [k·C_MAX/N, (k+1)·C_MAX/N)
依赖节点总数 N ✅ 强依赖(N 变→公式失效) ❌ 不依赖(区间划分是派生结果)
扩容时映射行为 全量重算(每个 hash 值重新取模) 仅切分受影响的区间(区间内部映射不变)
查找行所属节点 O(1) 取模 O(log N) 二分查找区间
实现复杂度 极简 中等

比喻hash % N 像按学号取模分班——规则本身定义在「班级数量」上,班级数一变规则就失效。Ring-style 像按身高区间分班——每个区间是独立的,增减班级只需调整区间边界,已在区间内的人不受影响。

(来源:C-Store 7 Years §3.6)

分布算法 映射规则 映射稳定性 查找复杂度 代表系统
hash % N hash(col) % N 不稳定(依赖 N,节点数变→全量重映射) O(1) Greenplum, ClickHouse
Ring-style 哈希空间等分,区间→节点 稳定(区间独立于 N,边界可调但内部映射不变) O(log N) Vertica

2.3 层次三:扩缩容迁移 — 节点变化时数据如何重分布

MPP 共性: 集群扩缩容时,数据必须重新分布——但各系统实现方式差异巨大,体现了「迁移效率 vs 运维复杂度」的不同权衡。

系统 扩缩容机制 迁移粒度 用户操作 迁移期间可用性
Vertica ELASTIC_CLUSTER(local segment 整块搬迁) Local segment(预分段单元) SELECT REBALANCE_CLUSTER() 在线(查询正常执行)
Greenplum gpexpand(全表重分布) 整表 gpexpand + 手动重分布每张表 在线(但性能显著下降)
ClickHouse 手动 ALTER TABLE RESHARD 整表(全量重分布) 手动执行 + 可能需要停服 取决于操作方式
Doris / StarRocks 自动分桶迁移 Bucket(分桶单元) ALTER TABLE SET DISTRIBUTION(后台异步) 在线
Snowflake 零运维(存储计算分离) 无(调整 warehouse 大小即可) 始终可用

以上是各系统的宏观对比。下面以 Vertica 为例,深入其迁移机制的两个关键设计:local segments 预分段segment 边界对齐

2.3.1 连续模型:ring-style 的扩容迁移

回到 §2.2 的 ring-style 4 节点示例($C_{MAX}=300$):

扩容前(4 节点):                 扩容后(6 节点):
Node 1:   0 ≤ hash <  75        Node 1:   0 ≤ hash <  50    ← 保留前 50 单位
                                Node 2:  50 ≤ hash < 100    ← 新节点
Node 2:  75 ≤ hash < 150        Node 3: 100 ≤ hash < 150    ← 保留后 50 单位
Node 3: 150 ≤ hash < 225        Node 4: 150 ≤ hash < 200    ← 保留前 50 单位
                                Node 5: 200 ≤ hash < 250    ← 新节点
Node 4: 225 ≤ hash < 300        Node 6: 250 ≤ hash < 300    ← 保留后 50 单位

每个原节点给出 $C_{MAX}/12 = 25$ 单位(自己范围的 1/3),保留 50 单位。两个新节点(Node 2 和 Node 5)各从两个相邻原节点接收 25+25=50 单位。总迁移量 = $4 \times 25 / 300 = 100/300$ = 1/3 ≈ 33.3%——比取模方案的 66.7% 少了整整 33.3 个百分点

比喻:取模方案像按学号重新分班——规则变了,三分之二的学生重新排队。Ring-style 像把每个教室的后 1/3 排同学搬到新教室——只需要挪 1/3,其余原地不动。

但「每个原节点恰好给出 1/3」成立的前提是:数据可以按任意精度切分,切分点能精确落在 50, 100, 150, 200, 250 这些边界上。要做到这一点,需要一种在 rebalance 之前就规划好的分段机制

2.3.2 Local segments:用预分段让迁移精确对齐

§2.3.1 的连续模型假设数据可以按任意精度切分,切分点能精确落在 50, 100, 150, 200, 250 上。但实际数据以 ROS container 形式存在磁盘上,不能随意切分。

Vertica 的做法:通过 scaling factor 控制每节点的 segment 数量。 scaling factor 是 Vertica 的可配置参数,指定每节点维护多少个 local segment(官方文档示例:scaling factor=8 → 5 节点集群共 40 个 segment)。Tuple Mover 在 mergeout 时按此参数将数据等分成对应数量的 segment(论文 §4:「takes care to preserve partition and local segment boundaries when choosing merge candidates」)。到 rebalance 时,这些 segment 已经就位,直接以整段为单位搬迁(论文 §3.6:「without any rearrangement or splitting necessary」)。

scaling factor 的官方定义(来源:Vertica 26.2.x 文档「Elastic Cluster」):

  • scaling factor = 每节点 local segment 数量(即前文 S)。
  • 扩容时系统将 segment 重新分配到新老节点。如果现有 segment 无法在允许偏差范围内实现均匀分布,系统会重新划分 segment 边界(「new local segment boundaries are drawn for each node」),再行分配。
  • 如果磁盘空间不足,rebalance 会分多轮迭代完成(「the rebalance operation spans multiple iterations」)。
  • 允许的数据倾斜由参数 MAXIMUM_SKEW_PERCENT 控制。

以 4→6 扩容为例,走一遍完整流程。

第一步:设置 scaling factor。 4 节点各 S 个 segment,共 4S 个;6 节点平分,每节点得 4S/6 个。要让每节点得整数个 segment → 4S/6 必须整除 → S 是 3 的倍数。取最小值 scaling factor = 3。

第二步:本地分段。 Tuple Mover 将每个节点的数据按哈希范围等分成 3 个 segment。原 Node 1 [0, 75) 被切分为:

原 Node 1 的 3 个 segment,各 25 单位:
  [0, 25)
  [25, 50)
  [50, 75)

切分点 50 恰好是 segment 边界——这正是 scaling factor = 3 的效果。

第三步:搬迁。 2 个新节点共需 1/3 总量 = 4 个 segment。每个原节点恰好贡献 1 个——Node 1 搬 [50,75),Node 2 搬 [75,100) 段,依此类推。总迁移 = 4/12 = 33.3%,与连续模型理想值完全一致,零浪费。

如果设 scaling factor = 2 会怎样? 4×2/6 = 8/6,除不尽。segment 总数不能被 6 整除 → 不可能让 6 个节点同时分到整数个 segment,必然导致数据倾斜。

切分点 50 落入 Node 1 的第二个 segment [37.5, 75) 内部——该 segment 中只有 [50,75) 的 25 单位属于新节点,但因不可拆分,整段 37.5 单位都得搬给 Node 2。Node 2 多拿了 12.5 单位。其他节点同理,最终节点间数据量参差不齐。

官方的处理方式: 如果倾斜超过 MAXIMUM_SKEW_PERCENT 阈值,系统不会硬搬——而是重新划分 segment 边界(「new local segment boundaries are drawn for each node」),把 segment 重新均分后再搬迁。这相当于自动纠正了 scaling factor 选择不当的问题,但代价是额外的边界重算开销。

简言之:scaling factor 选整除值(如 3),segment 边界自然对齐,一步到位。选不整除值(如 2),系统也能处理——自动重划边界——但多了一道工序。

scaling factor 4×SF/6 是否整除 每原节点搬走 总迁移 效果
3 ✅ 整除(=2) 1 个 segment 4/12 = 33.3% 零浪费
6 ✅ 整除(=4) 2 个 segment 8/24 = 33.3% 零浪费(segment 更碎)
2 ❌ 不整除 不均→系统自动重划 segment 边界 需额外边界重算 能处理,但多一道工序

一句话:扩容前,根据新老节点总数算出整除的 scaling factor 值,Tuple Mover 依此做本地分段,rebalance 直接搬迁。若 scaling factor 选得不整除,系统也能自动重划 segment 边界纠正——但多一道工序。(来源:C-Store 7 Years §3.6;Vertica 26.2.x 文档「Elastic Cluster」)

2.3.3 ELASTIC_CLUSTER:整段搬迁

数据以原生格式(已排序、已编码、已压缩的 ROS container)直接迁移,目标节点接收后立即可用。这就是 rebalance 中 ELASTIC_CLUSTER 方法的底层机制——它不是「重新哈希 → 重新分发 → 重新排序 → 重新编码」,而是「识别要迁移的 local segment → 整块搬走」。

迁移方式 解决的问题 核心思路 4→6 迁移比例
hash % N 重分布 最基础的扩容 全表重算哈希,重新分发 66.7%
Ring-style 连续区间切分 大幅减少迁移量 每个原节点只给出 $M/(N+M)$ 比例的范围 $M/(N+M)$ = 2/6 = 33.3%
+ Local segments 预分段对齐 以 segment 为单位整块搬迁 设定 scaling factor 使 segment 总数被新老节点数整除,本地分段后搬迁整 segment;若不整除则自动重划边界(Vertica 独有) scaling factor=3 → 33.3%;=2 → 需重划边界

(来源:C-Store 7 Years §3.6:「When nodes are added or removed, data is quickly transferred by assigning one or more of the existing local segments to a new node and transferring the segment data wholesale in its native format, without any rearrangement or splitting necessary.」;Vertica 集群 Rebalance 完全指南 §1.2, §1.4)


2.4 分布键与数据组织的关系

分段键与排序键是解耦的【Vertica 独有,但在概念上可对比其他系统】

在 Vertica 中,分段键和排序键是独立的——分段键决定数据去哪个节点,排序键决定数据在节点内如何组织:

CREATE PROJECTION orders_proj
AS SELECT * FROM orders
ORDER BY order_date          -- 节点内按日期排序
SEGMENTED BY HASH(customer_id) ALL NODES;  -- 节点间按客户哈希分段

这种解耦让 Vertica 可以同时优化两个维度:分段键优化并行度(JOIN/GROUP BY 的本地执行),排序键优化节点内的查询效率(谓词过滤 / RLE 压缩 / Merge Join)。

比喻:分段键决定书放在哪个图书馆(城市),排序键决定书在图书馆内的哪个书架(按出版日期排列)。

对比其他系统——虽然分段与排序解耦是 Vertica 独有的显式设计(因为其他系统没有投影级分段),但其他系统也有类似的「两层组织」概念:

系统 跨节点分布 节点内排序 两者是否独立
Vertica SEGMENTED BY HASH(col) ORDER BY col ✅ 完全独立
Greenplum DISTRIBUTED BY (col) 无排序键(堆表)或 AO 表的 sort key 表级耦合
Redshift DISTKEY(col) SORTKEY(col) ✅ 独立
ClickHouse sharding_key(Distributed 引擎) ORDER BY col(MergeTree 引擎) ✅ 独立(但分布键在 Distributed 层,排序键在 MergeTree 层——物理上分层而非同层解耦)
Snowflake 自动 clustering 自动 micro-partition clustering ✅ 自动但不可见

(来源:C-Store 7 Years §3.1 Figure 1;Redshift 文档 DISTKEY/SORTKEY;ClickHouse 文档 MergeTree ORDER BY + Distributed sharding_key)

分段与分区:两个正交的概念【通用】

这是最常见的混淆点。一张表可以同时有分段和分区:

CREATE TABLE trade (
    tdate DATE, tsymbol VARCHAR(8), ttime TIME
)
PARTITION BY EXTRACT(year FROM tdate);  -- 每个节点内按年份分区

CREATE PROJECTION trade_proj
AS SELECT * FROM trade
ORDER BY tdate, tsymbol
SEGMENTED BY HASH(tsymbol) ALL NODES;  -- 节点间按股票代码分段
维度 分段(Segmentation) 分区(Partitioning)
作用域 跨节点(决定数据在哪个节点) 单节点内部(决定数据在哪个 ROS container)
目的 MPP 并行计算、本地 JOIN/聚合 减少 I/O(存储裁剪)、数据生命周期管理
定义位置 Projection 级别(SEGMENTED BY Table 级别(PARTITION BY
对查询的影响 决定并行度和网络传输 决定存储裁剪效率

所有 MPP 系统都有类似的分层设计——区别在于具体实体的名称。Greenplum 叫「distribution + partition」,ClickHouse 叫「sharding + partition key」,Doris 叫「distribution + partition」。无论叫什么,核心逻辑不变:跨节点分布负责并行度,节点内分区负责 I/O 效率,两者各司其职

一张大表通常同时需要分段和分区。分段保证查询在所有节点上并行执行,分区保证每个节点内部能快速裁剪到目标数据。

(来源:理解 Vertica 的分区 §1.2;Vertica 表分区策略选择指南 §1.2)


3. 设计决策与 Trade-off

3.1 Table-level 分布 vs Projection-level 分段

这是本文最核心的设计对比。

设计 代表系统 优势 代价
Table-level 分布 Greenplum, Teradata, Redshift 简单:一张表一个分布键,无需在投影间选择 只有一种分布方式,无法针对多维度 JOIN 优化
Projection-level 分段 Vertica 每种 JOIN 模式可用不同分段键的投影,优化器自动选择 存储冗余(每份投影是完整副本)、ROS 写入与投影数成正比、优化器搜索空间更大

Table-level 分布更简单,projection-level 分段更灵活。选择哪一种取决于你的查询模式

  • 如果事实表 90% 的 JOIN 都用一个键(如 customer_id),table-level 分布就够了,额外的投影是浪费
  • 如果事实表高频 JOIN 多个维度(客户、产品、日期、地区),projection-level 分段的价值才会显现

在实践中,大多数 Vertica 客户只有 0-3 个额外的窄投影外加 1 个 super projection(来源:C-Store 7 Years §3.1)。因为投影的存储和运维成本是真实的——每个额外投影都是一份完整的物理副本,占用额外磁盘空间,且 mergeout 和 rebalance 操作需要在每个投影上独立执行。

3.2 全副本 vs 哈希分布:所有 MPP 的共同选择

MPP 共性: 每个 MPP 系统都需要回答「小维度表怎么办」——如果维度表也按某个键分段,那么 JOIN 事实表的非分段键维度时仍需要网络传输。所有 MPP 系统都提供了同一种解决方案:让小表在每个节点上存一份全量副本。大表则始终用哈希分布——只在各节点存 1/N 的数据。

系统 全副本策略 语法 哈希分布策略 语法
Vertica UNSEGMENTED ALL NODES CREATE PROJECTION ... UNSEGMENTED ALL NODES SEGMENTED BY HASH(col) CREATE PROJECTION ... SEGMENTED BY HASH(col) ALL NODES
Greenplum DISTRIBUTED REPLICATED CREATE TABLE ... DISTRIBUTED REPLICATED DISTRIBUTED BY (col) CREATE TABLE ... DISTRIBUTED BY (col)
Redshift DISTSTYLE ALL CREATE TABLE ... DISTSTYLE ALL DISTSTYLE KEY + DISTKEY(col) CREATE TABLE ... DISTSTYLE KEY DISTKEY(col)
ClickHouse 无内置「全节点副本」分布策略 —(通过 GLOBAL JOIN 在查询时将小表广播到所有节点) Distributed 引擎 + sharding_key ENGINE = Distributed(cluster, db, table, rand())
Doris / StarRocks 无内置「全节点副本」分布策略 —(推测:通过 Broadcast Join 在查询时广播小表) DISTRIBUTED BY HASH(col) BUCKETS N CREATE TABLE ... DISTRIBUTED BY HASH(col) BUCKETS 32
Snowflake 自动(小表自动缓存到各节点本地 SSD) 无用户语法 自动 clustering 无用户语法

关键差异:全副本是存储层还是查询层实现? Vertica、Greenplum、Redshift 将全副本做在存储层——数据常驻每个节点,查询时直接用,I/O 不随查询次数变化。ClickHouse 和 Doris 没有存储层的全副本分布,而是通过 GLOBAL JOIN / Broadcast Join查询层将小表广播——每次查询都需要网络传输小表数据,但换来更简单的存储管理。Snowflake 走的第三条路——自动缓存热数据到本地 SSD,对用户完全透明。

以 Vertica 为例

策略 原理 存储成本 JOIN 行为 适用场景
UNSEGMENTED 每节点存全量副本 N 倍(N = 节点数) 任意 JOIN 都能本地完成 极小维度表(< 10 万行)
SEGMENTED 每节点存 1/N 数据 1 倍 只有分段键匹配的 JOIN 能本地完成 所有中大型表

UNSEGMENTED 是一把双刃剑。小维度表用 UNSEGMENTED ALL NODES 是生产环境的正确选择——每个节点都有全量副本,任意 JOIN 都能本地完成,无需 Resegment。但大表这样用是灾难——每节点要构建完整哈希表,内存压力极大。这一约束适用于所有 MPP:Greenplum 的 DISTRIBUTED REPLICATED、Redshift 的 DISTSTYLE ALL 都只能用于小表——只不过「小」的定义因系统而异(Vertica 推荐 < 10 万行,Greenplum 和 Redshift 通常容忍到百万行级别,因为它们的存储引擎不涉及 ROS container 的额外维护开销)。

语法上存在 UNSEGMENTED NODE <n> 将数据限定到单节点,但这仅用于临时表、测试环境或特定合规需求等极特例场景。生产环境始终用 UNSEGMENTED ALL NODES

📋 真实案例:某运营商 93 节点集群,大表按 statis_date 分段,查询过滤 statis_date = '20200819' → 分段键与过滤键重合导致当天全部数据集中在一台节点 → 6.5 小时未完成。改为按 JOIN 键分段后恢复正常。来源:某运营商 Vertica 数据仓库性能问题分析报告(2020-08-21)。

3.3 分段键选择的 trade-off:基数、均匀性、JOIN 对齐

分段键的选择需要同时权衡三个维度:

因素 要求 反例
基数 足够高(N 个节点至少需要 > 100 × N 个 distinct 值) gender(2 个 distinct 值)→ 数据最多分布在 2 个节点,其余节点空转
均匀性 值分布接近均匀 customer_id 中某大客户占 30% 数据 → 30% 数据集中在一个节点
JOIN 对齐 与最高频 JOIN 键一致 order_id 分段但 90% JOIN 用 customer_id → 每次 JOIN 都需要 Resegment

这三个要求有时互相冲突。例如,date_id 基数高且与时间维度 JOIN 对齐,但数据加载通常是按时间顺序的——所有今天的数据哈希到同一个节点 → 加载时单节点成为热点。此时需要权衡:是优先查询性能(分段时间维度)还是加载均匀性(分段主键)。

(来源:C-Store 7 Years §3.6:「The most common choice is HASH(col₁..colₙ), where colᵢ is some suitably high cardinality column with relatively even value distributions, commonly a primary key column.」)

3.4 Hash 分段 vs 一致性哈希:动态扩容的工程权衡

Vertica 使用的 ring-style 分段(§2.2)并非标准的一致性哈希。它没有虚拟节点(virtual nodes),也不支持节点在环上的任意位置插入。新节点的插入位置由系统计算——选择数据迁移量最小的位置(来源:Vertica 集群 Rebalance 完全指南 §1.2)。

如果 Vertica 使用简单的 hash % N,从 4 节点扩到 6 节点时,66.7% 的数据需要迁移(hash%N 下保留概率 = min(N,M)/lcm(N,M) = 4/12 ≈ 33.3%,仅 1/3 数据留在原节点)。而 ring-style 的连续模型只需迁移 $M/(N+M) = 2/6$ = 33.3%——比取模方案少 33.3 个百分点。加上 local segments 预分段(见 §2.3.2),设 S=3 使 segment 边界对齐切分点,实际迁移等于理想值。

方案 4→6 节点数据迁移比例 实现复杂度
hash % N 66.7%(需全表重算哈希) 极简
Ring-style 连续模型 33.3%($M/(N+M)$)
Ring-style + local segments(scaling factor=3) 33.3%(整除,边界对齐) 中等
Ring-style + local segments(scaling factor=2) 倾斜超阈值→系统自动重划边界 中等
完整一致性哈希(virtual nodes) 33.3%(virtual node 数量大,等效能整除)

核心差异:hash%N 的 66.7% vs ring-style 的 33.3% 来自算法本身——ring-style 每个原节点只需给出 $M/(N+M)$ 比例的范围,且旧节点间无需换位;hash%N 除了 33.3% 迁往新节点外,还有 33.3% 在旧节点间无意义换位。S=3 vs S=2 的 33.3% vs 37.5% 差异则来自 S 是否让 segment 总数被新老节点数整除(见 §2.3.2)——整除时边界对齐,不整除时切分点落在 segment 内部,产生额外搬迁。

Vertica 选择 ring-style + local segments 的组合:以中等复杂度换取了比取模方案少 1/2 的迁移量,同时通过选择整除的 S 使 segment 边界对齐切分点,迁移无浪费;segment 整块搬迁又避免了逐行重算哈希和重排序。

3.5 多投影 vs 单投影 + Resegment:空间换时间

策略 查询性能 存储成本 ROS 写入成本 运维复杂度
单投影 + Resegment 依赖网络(JOIN 时 Resegment) 1 倍
2 个投影(不同分段键) 两个 JOIN 维度都能本地执行 2 倍
4 个投影(全面覆盖) 所有 JOIN 都本地执行 4 倍 高(rebalance × 4)

关于「加载成本」的精确区分:COPY 操作的数据解析和网络分发是一次性的——无论有多少投影,源数据只被读取和解析一次。随投影数量线性增长的是 ROS 容器写入成本(每个投影各自生成 ROS container 并写入磁盘),以及后续的 mergeout / rebalance 成本。在 Vertica 10.0 之前,WOS(Write Optimized Store)作为内存缓冲层,COPY 数据先写入单份 WOS,再由 Tuple Mover 的 moveout 操作分发到各投影的 ROS——进一步将写入开销延后到了后台。10.0 起 WOS 被移除,COPY 直写 ROS,但执行引擎仍以单次 pass 同时为所有投影生成 ROS container,因此 COPY 本身的 CPU/网络开销仍不随投影数线性增长。详见 §4.2。

这个 trade-off 由 DBD(Database Designer)的三个设计策略量化控制:

  • Load-optimized:最少投影(K+1 个 super projection),优先加载速度
  • Query-optimized:为所有查询创建最优投影,优先查询性能
  • Balanced:在收益递减点停止(默认优化 75% 的查询)

(来源:DBDesigner §IV 设计策略)


4. 设计对实际使用的影响

4.1 查询维度

自动生效的收益

  • 分段键 = JOIN 键时,优化器自动选择 Local Join——完全本地执行,零网络开销。这是「零干预」的最高性能
  • 分段键 = GROUP BY 键时,优化器自动选择 Fully Distributed Group-by——每个节点独立聚合,最后汇总
  • 小表 UNSEGMENTED + 大表 SEGMENTED:小表在每个节点都有全量副本 → 任意 JOIN 都是 Local Join

需要手动干预的场景

  • EXPLAIN 中出现 RESEGMENT(涉及大表) → 分段键与 JOIN 键不匹配。需要创建按 JOIN 键分段的新投影,或接受网络开销
  • projection_storagerow_count 节点间差异 > 100% → 分段键存在数据倾斜。需要选择更高基数或更均匀的分段键
  • 单节点执行(Execute on: v_xxx_node0001 → 常见原因是分段键与查询过滤键重合:分段键决定数据去哪个节点,查询恰好过滤分段键的某个具体值,过滤后所有数据集中在一台节点

4.2 加载维度

分段设计对数据加载的影响需要区分不同环节来理解——并非「投影数量 N = 加载时间 N 倍」那么简单。

COPY 解析与网络分发:一次性操作。 无论一张表有多少投影,COPY 对源数据的读取、解析、哈希计算和跨节点分发只执行一次。执行引擎在单次 pass 中同时为所有投影生成数据流——不会因为 3 个投影就把源数据解析 3 遍。

ROS 容器写入:随投影数量线性增长。 每个投影各自将分段后的数据写入独立的 ROS container 文件。这是投影数量影响写入性能的主要环节。

WOS 缓冲的演进(10.0 是关键分水岭)

版本 写入路径 对多投影的影响
≤ 9.2 COPY → 单份 WOS 内存缓冲 → Tuple Mover moveout → 各投影 ROS WOS 是共享的,COPY 写 WOS 不受投影数影响。ROS 写入被延后到后台 moveout
9.3 新建数据库默认跳过 WOS,直写 ROS(DMLTargetDirect 控制) 过渡期
10.0+ WOS 完全移除,COPY 直写 ROS。Tuple Mover 只做 mergeout 执行引擎单次 pass 同时生成所有投影 ROS,解析/分发不翻倍,磁盘 I/O 翻倍
Eon Mode(所有版本) 从未有过 WOS,始终直写公共存储 同 10.0+

(来源:C-Store 7 Years §3.7 WOS 架构;Vertica 9.3/10.0 文档 WOS 废弃说明;Eon Mode SIGMOD 论文)

Prejoin projection 的特殊加载成本:与普通投影不同,prejoin projection 在加载时执行 JOIN——维度表越大,加载越慢。实践中大多数客户不愿为查询性能牺牲加载速度(C-Store 7 Years §3.3)。

实际影响总结

  • 存储成本与投影数成正比(每个投影是完整物理副本)
  • ROS 写入 I/O 与投影数成正比
  • COPY 的解析/网络/CPU 成本基本不受投影数影响(单次 pass)
  • mergeout 和 rebalance 成本与投影数成正比(每个投影独立处理)

4.3 运维维度

Rebalance 成本:集群扩缩容时,每个投影都要独立完成数据重分段。投影数量 × 分段复杂度 = rebalance 总耗时。在 rebalance 的四个阶段中,分段投影的 resegment 和 ROS split 可占总耗时的 80%(来源:Vertica 集群 Rebalance 完全指南 §1.4)。

恢复成本:节点故障恢复时,如果 buddy projection 有相同的排序和分段,可以直接复制物理文件(最快)。如果排序不同,则需要类似 INSERT...SELECT 的执行计划来重建数据(较慢)。(来源:C-Store 7 Years §5.2)

监控要点

-- 检查分段倾斜:各节点间数据量是否均匀
SELECT node_name, projection_name, row_count
  FROM projection_storage
 WHERE anchor_table_name = 'orders'
 ORDER BY projection_name, node_name;

4.4 常见误解与澄清

误解 事实
「UNSEGMENTED 能避免所有网络传输」 UNSEGMENTED ALL NODES 确实让每节点有全量副本(可并行),但大表会导致每节点都需构建完整哈希表 → 内存爆炸
「投影越多越好」 每个投影 = 额外的磁盘空间 + rebalance/mergeout 独立处理。COPY 解析不翻倍但 ROS 写入翻倍。0-3 个窄投影是最佳实践范围(C-Store 7 Years §3.1)
「分段键必须是主键」 主键是常见的推荐选择(基数高、分布均匀),但不是唯一的。分段键的真正要求是基数高 + 分布均匀 + JOIN 对齐
「分区能替代分段」 分区是节点内的组织方式,分段是跨节点的分布方式——两者正交,各司其职
「REBALANCE 就是简单的重新哈希」 Rebalance 涉及 ROS split(最耗时的阶段)、数据传输、Tuple Mover 合并等多个阶段,是一个 CPU/磁盘/网络密集型操作

5. 案例验证

5.1 虚构案例:分段键不对齐导致 RESEGMENT

📝 虚构案例

场景:某零售企业 Vertica 集群(6 节点),事实表 sales_fact(按 sale_id 分段,5 亿行/天),维度表 product_dim(按 product_category 分段,2000 万行)。高频日报查询:

SELECT p.product_name, SUM(s.quantity * s.unit_price) AS revenue
FROM sales_fact s
JOIN product_dim p ON s.product_id = p.product_id
WHERE s.sale_date = CURRENT_DATE - 1
GROUP BY p.product_name;

EXPLAIN 关键输出

+-JOIN HASH [Cost: 45M, Rows: 500M] (PATH ID: 2)
|  Join Cond: (s.product_id = p.product_id)
| +-- Outer -> STORAGE ACCESS for s [Rows: 500M] ← 按 sale_id 分段
| +-- Inner -> STORAGE ACCESS for p [Rows: 20M] (RESEGMENT) ← 按 product_category 分段!

根因sales_fact 分段键 = sale_idproduct_dim 分段键 = product_category,JOIN 键 = product_id——三个键各不相同。优化器以 product_dim 为 inner 表并对其做 RESEGMENT,2000 万行 × 6 节点的网络传输。

更致命的是,product_dim 重分段后发现 product_id 有热点值(大品牌产品),导致某节点接收的网络数据量是其他节点的 4 倍。

修复:为 product_dim 创建按 product_id 分段的投影:

CREATE PROJECTION product_dim_by_id
AS SELECT * FROM product_dim
ORDER BY product_id, product_name
SEGMENTED BY HASH(product_id) ALL NODES KSAFE 1;

SELECT REFRESH('product_dim_by_id');

效果

指标 优化前 优化后
执行时间 12 分钟 1.5 分钟
RESEGMENT 网络量 ~18 GB 0(本地 JOIN)
节点间 CPU 差异 4×(热点节点 98%) < 15%

原理回溯:这个案例体现了 §3.3 的 trade-off——product_dim 原来按 product_category 分段虽然能优化按品类过滤的查询,但与最高频 JOIN(product_id)不匹配。分段键选择必须优先对齐最高频的 JOIN 键——这是投入产出比最高的优化。

5.2 虚构案例:DBD 如何为多查询自动选择分段策略

📝 虚构案例

场景:某电商平台有两类高频查询:

  • Q1:sales JOIN customers ON customer_id — 客户行为分析(每天 500 次)
  • Q2:sales JOIN products ON product_id — 商品销售分析(每天 300 次)

DBD 在 query-optimized 策略下的决策过程:

  1. 从 Q1 提取 Join Cols:{customer_id},从 Q2 提取 Join Cols:{product_id}
  2. 枚举分段候选:HASH(customer_id)(收益:Q1 本地 JOIN)、HASH(product_id)(收益:Q2 本地 JOIN)、HASH(sale_id)(无特殊收益但加载均匀)
  3. 优化器对每个候选评估 cost:Q1 用 HASH(customer_id) → Network cost = 0,Q2 用 HASH(product_id) → Network cost = 0
  4. 两个候选的 Benefit 都足够高 → DBD 为 sales 表创建两个窄投影,分别按 customer_idproduct_id 分段

效果

查询 优化前(单投影 HASH(sale_id) 优化后(双投影)
Q1(customer JOIN) 12s(Resegment) 2s(Local Join)
Q2(product JOIN) 15s(Resegment) 2.5s(Local Join)
存储增量 +40%(两个窄投影比一个 super projection 多存储的列数据)

但这是有代价的:两个窄投影意味着加载后需要写入三份 ROS(super + 两个窄投影),存储和 mergeout 成本约为单投影的 3 倍。如果业务要求 15 分钟内的数据新鲜度,这个运维开销可能不可接受——此时需要用 DBD 的 balanced 策略只保留收益最高的那个投影。

(来源:DBDesigner §V 候选枚举算法;§VI 性能评估)

5.3 跨系统模拟案例:同一场景下各 MPP 的分布策略表现

📝 虚构案例 · 跨系统对比

场景:某大型零售商,事实表 sales(100 亿行,日增 5 亿行),6 节点集群(每节点 256GB 内存、8 块 NVMe SSD)。两类高频查询:

  • Q1:sales JOIN customers ON customer_id(客户行为分析,日均 500 次)
  • Q2:sales JOIN products ON product_id(商品销售分析,日均 300 次)

对比问题:在不同的 MPP 系统中,如何为这张表设计分布策略?代价分别付在哪里?

各系统的分布设计与预估行为:

维度 Vertica Greenplum ClickHouse Snowflake
分布设计 2 个窄投影:分别按 customer_idproduct_id 分段 DISTRIBUTED BY HASH(customer_id)(只能选一个) 本地表 ORDER BY (sale_date, ...) + Distributed 表按 customer_id 分片 自动 clustering(用户不指定分布键)
Q1 执行方式 Local Join(选 orders_by_cust 投影) Local Join(分布键对齐 ✅) 本地 JOIN(分片键对齐 ✅) 自动优化(推测:数据自动按查询模式重组 clustering)
Q1 预估耗时 2s 2s 1.8s(推测:基于 MergeTree 的预排序优势) 3-5s(推测:无用户控制的分布键,依赖自动 clustering 质量)
Q2 执行方式 Local Join(选 orders_by_prod 投影) Redistribute Motion(分布键不对齐 ❌) GLOBAL JOIN(分片键不对齐 ❌) 自动优化(同上)
Q2 预估耗时 2.5s 45s(推测:100 亿行重分布,网络传输 ~400GB) 35s(推测:分布式子查询 + 各节点拉取全量维度表) 5-10s(推测:自动 clustering 可能已按 product_id 重组部分数据)
存储成本 3×(super + 2 个窄投影) 1×(本地表) 自动压缩 + 历史数据自动归档
加载吞吐 单日解析约 2 TB/h(解析 1 次,写 3 份 ROS) 单日解析约 3.5 TB/h(解析 1 次,写 1 份堆表) 单日解析约 4 TB/h(推测:MergeTree 写入效率高) 自动弹性
扩缩容迁移 ELASTIC_CLUSTER,迁移比例由 $M/(N+M)$ 和 S 决定(见 §2.2),+50% 容量时约 33.3% gpexpand,全表扫描重分布(I/O 100%,网络迁移 ~50%) 手动 RESHARD(全量重分布) 零运维
运维复杂度 中(需管理多投影的 mergeout/rebalance) 低(单表单分布键) 中(需维护 Distributed 表 + 本地表双层结构) 极低(全托管)

关键洞察:代价付在哪里?

  • Vertica 把代价付在存储和写入上——用 3 倍存储换所有 JOIN 本地执行。这是「空间换时间」的极端案例。
  • Greenplum 把代价付在Q2 的网络上——一旦分布键选错,每次 Q2 查询都要在网络上传输 400GB+ 数据。这是「简单性换灵活性」的 trade-off。
  • ClickHouse 同样面临分布键二选一的困境,但本地表的 MergeTree 排序可以部分缓解——如果 Q2 维度表很小,GLOBAL JOIN 的代价可能比预期低。
  • Snowflake 把代价付在用户无法控制上——如果自动 clustering 质量高,两个查询都很快;如果 clustering 没跟上工作负载变化,两个都可能慢。用户失去了手动优化的能力,但换来了零运维。

跨系统启示:这个案例揭示了一个通用原则——分布策略的灵活性越高,存储和运维的代价就越大;分布策略越简单,查询的不可预测性就越高。没有哪个系统在这两个维度上同时做到最优,选择取决于你的查询模式是「多维度且可预测」(选 Vertica)还是「单维度且稳定」(选 Greenplum/Redshift)还是「愿意牺牲控制权换运维简单」(选 Snowflake)。

5.4 真实案例:分段键与过滤键重合 → 数据全部落在一个节点 → 单节点瓶颈

📋 真实案例 · 来源:某运营商 Vertica 数据仓库性能问题分析报告(2020-08-21)

背景:某运营商,93 节点集群。凌晨 3 点数据库性能严重下降。

故障现象:一条 SQL 从 22:00 持续执行到次日 04:44(超过 6.5 小时仍未完成),网络传输总量达 11 GB

诊断过程

  1. 执行计划显示绝大部分步骤在单节点上执行(Execute on: v_edw_node0039
  2. 执行计划中出现 Outer (RESEGMENT)(LOCAL ROUND ROBIN) + Inner (RESEGMENT)——两个 RESEGMENT 都是节点内的本地哈希重组,用于 Hash Join 并行分桶
  3. 两个表的统计信息均已过期(NO STATISTICS
  4. 关键发现:表 1 按 statis_date 分段;查询 SQL 中恰好有 WHERE statis_date = '20200819'——因为所有 statis_date = 20200819 的数据按哈希分段落在同一节点上,过滤后的结果集全部集中在一个节点
  5. 表 2 采用 UNSEGMENTED 方式分布

这是「分段键 = 谓词键」导致的经典陷阱。分段键决定了每行数据去哪个节点。当查询的过滤条件恰好等于分段键的某个具体值时,所有符合条件的数据都在哈希空间的同一个区间——也就是同一台节点。查询无法并行,只能在那台节点上执行。

在这个案例中,表 1 在同一查询中出现了 3 次(SQL 中分别作为子查询 b、c 和直接 JOIN 的表 d),每次都带 statis_date = '20200819' 过滤——三份中间结果全部落在同一台节点上,单节点执行无可避免。

根因链条:表 1 按 statis_date 分段 → 查询谓词 WHERE statis_date = '20200819' 将数据过滤到哈希空间中的同一个分段区间 → 过滤后的所有数据集中在一台节点上 → 表 1 在查询中出现 3 次(子查询 b、c 和直接 JOIN 的表 d),全部落在这台节点 → 整个 JOIN 被迫在单节点执行 → 单节点承担全部哈希表构建 + JOIN 计算 → 6.5 小时仍未完成。

修复

  • 表 1:投影改为 SEGMENTED BY HASH(user_id_zk, user_id_fk)(对齐 JOIN 键)
  • 表 2:从 UNSEGMENTED 改为 SEGMENTED BY HASH(JR_USER_ID, KD_USER_ID)(对齐 JOIN 键)

效果:单次查询耗时从 11 分 48 秒降至 5 分 37 秒,且不再有单节点瓶颈。

原理回溯:这个案例揭示了分段设计中一个容易被忽视的反模式——分段键与高频过滤键重合

分段键 statis_date 使得某一天的数据全部集中在一台节点上。查询恰好按某一天过滤——这本该让存储裁剪高效工作,但因为过滤后的数据全在一台节点,并行度直接降为 1。更致命的是,表 1 在同一查询中被引用了 3 次,三份中间结果全部落在同一台节点,完全没有利用 93 个节点的并行能力。

修复的核心思路:将表 1 的分段键从 statis_date 改为 JOIN 键 (user_id_zk, user_id_fk)——数据按 JOIN 键均匀分布到所有节点,过滤后的结果也分散在各节点上,JOIN 可以在 93 个节点上并行执行。

5.5 真实案例:统计信息缺失导致分段设计的性能潜力无法释放

📋 真实案例 · 来源:某运营商 Vertica 数据库性能问题报告(2021-05-17)

背景:某运营商,50 节点集群,Vertica v7.2.3,每节点 256GB 内存。

故障现象:多表关联查询执行超过 1 小时。

诊断:5 张 JOIN 表缺失统计信息。优化器在盲猜每张表的大小 → 做出的 JOIN 顺序和 Broadcast/Resegment 决策完全错误。

修复:对这 5 张表执行 SELECT ANALYZE_STATISTICS('schema.table');

效果:执行时间从 1 小时降至 3.5 秒,提升超过 1000 倍。

与数据分布的关系:统计信息是分段设计能发挥作用的前提条件。即使投影的分段键完美对齐了 JOIN 键,如果统计信息缺失,优化器可能根本意识不到 Co-location 的优势——它不知道数据已经分布在对的节点上,也就不会选择 Local Join。统计信息缺失 = 优化器在盲猜 = 再好的分段设计也白费。


6. 设计原则总结

原则 1【通用】:分段键 = 最高频 JOIN 键。这是投入产出比最高的优化

为什么:分段键决定数据在哪个节点。当两张表的分段键相同且等于 JOIN 键时,JOIN 完全本地执行——零网络传输、零 resegment。这是「一次设计,永久受益」的优化。反例:用 order_id 分段但 90% JOIN 用 customer_id——每个 JOIN 都需要 Resegment 或 Broadcast。在 Greenplum 中违反此原则会导致每次查询触发 Redistribute Motion;在 ClickHouse 中会导致 GLOBAL JOIN 而非本地 JOIN。

原则 2【通用】:分段键必须是高基数、均匀分布的列

为什么:基数低 → 数据只能分布在少数节点,其余节点空转。分布不均 → 热点节点成为集群瓶颈。Vertica 推荐 HASH(col₁..colₙ),其中 colᵢ 通常为主键列。反例:用低基数列如 gender(2 个 distinct 值)分段 → 数据最多分布在 2 个节点,6 节点集群 4 个节点空转。此约束对所有 hash 分布的系统均适用——Greenplum 的 DISTRIBUTED BY、ClickHouse 的 sharding_key、Doris 的 DISTRIBUTED BY HASH 都面临同样的均匀性要求。

原则 3【通用】:分段键不应与高频过滤键重合——这是「谓词裁剪」和「并行度」的 trade-off

为什么:用日期列分段看似能优化按日期的查询(存储裁剪),但一旦查询过滤到某一天,所有数据落在一台节点上,并行度降为 1。如果每日数据量巨大,单节点瓶颈远大于存储裁剪的收益。优先用 JOIN 键分段而非过滤键分段。反例:93 节点集群,大表按 statis_date 分段,查询过滤 statis_date = '20200819' → 当天全部数据集中在一个节点 → 6.5 小时未完成(来自 某运营商 Vertica 数据仓库性能问题分析报告(2020-08-21))。在 Greenplum 中同样成立——DISTRIBUTED BY (date_col) + WHERE date_col = '2024-01-01' → 所有数据集中在一个 segment。

原则 4【Vertica】:UNSEGMENTED 只用于真正小的维度表(< 10 万行)

为什么:UNSEGMENTED 每节点存全量副本 → 任意 JOIN 都能本地完成(好)。但大表 UNSEGMENTED = 每节点都要构建完整哈希表 → 内存压力巨大(坏)。反例:3 亿行表 UNSEGMENTED → 每节点 3 亿行哈希表 → 200GB+ 内存占用 → JOIN 必然 spill。Greenplum 对应概念为 DISTRIBUTED REPLICATED,ClickHouse 对应概念为 ReplicatedMergeTree——虽然机制不同,但同样适用「小维度表才用全副本」的原则。

原则 5【Vertica】:投影数量与性能不是线性关系——超过 3 个窄投影后收益骤减

为什么:每个额外投影 = 额外的磁盘空间 + ROS 写入 I/O + Tuple Mover 多一份维护 + rebalance 多一份处理。C-Store 7 Years §3.1 指出大多数客户只有 0-3 个窄投影。反例:为每种 JOIN 模式都建投影 → 10 个投影 → mergeout 和 rebalance 成本约为单投影的 10 倍。

原则 6【通用】:分段与分区的选择各司其职——两者正交,一张大表通常两者都需要

为什么:分段决定数据在哪个节点(并行度),分区决定节点内数据如何组织(I/O 裁剪)。分段不足以优化时间范围查询的 I/O,分区不足以让 JOIN 在多个节点并行。反例:大表只分段不分区 → 按日期查询时需要扫描所有 ROS container;大表只分区不分段 → 数据在一台节点,无法并行。在 Greenplum 中同样存在 DISTRIBUTED BY + PARTITION BY 的组合使用;在 ClickHouse 中为 sharding_key + PARTITION BY

原则 7【通用】:统计信息是分段设计发挥作用的必要条件

为什么:优化器依赖统计信息来评估 Co-location 和 Network cost。缺失统计信息 = 优化器不知道数据已在正确的节点上 = 可能错误选择 Resegment。反例:分段键完美对齐 JOIN 键,但 5 张表缺失统计信息 → 优化器盲猜 → 查询 1 小时降到 3.5 秒(统计信息修复后,来自 某运营商 Vertica 数据库性能问题报告(2021-05-17))。所有基于 CBO 的 MPP 系统都适用——Greenplum 缺失统计信息同样会导致错误的 Motion 选择,ClickHouse 缺失统计信息可能导致错误的 JOIN 策略。

原则 8【Vertica】:扩容前做好 Rebalance 准备,投影越少越快

为什么:Rebalance 80% 的时间花在分段投影的 ROS split 上,每个投影独立处理。清理无用投影 = 直接减少 rebalance 的处理量。反例:集群有 30,000 个投影,rebalance 杂项操作就耗时约 2 小时(每个投影约 250ms catalog 更新,单线程串行)。Greenplum 的 gpexpand 和 ClickHouse 的 RESHARD 虽然也涉及数据重分布,但各自有不同的性能瓶颈——Greenplum 瓶颈在整表全量重写,ClickHouse 瓶颈在分布式表的协调开销。

原则 9【通用】:分段键设计要考虑未来的加载模式

为什么:按时间相关列(如 date_id)分段虽然优化了时间维度 JOIN,但会导致数据加载热点——同一时间段的数据全部哈希到同一个节点,该节点成为加载瓶颈。反例:按 trade_date 分段 → 每天的数据都写入同一个节点 → 加载吞吐受限于单节点 I/O。在 Greenplum 中同样成立——DISTRIBUTED BY (date_col) + 按时间顺序加载 → 单 segment 写入热点。


7. 延伸阅读

Vault 内笔记(按推荐阅读顺序)

  1. C-Store 7 Years Later — §3.6 是本文的核心理论来源,包含 ring-style 分段和 local segments 的完整描述。§3.1 提供了多投影使用频率的真实数据;§5.2 描述了 buddy projection 在分段恢复中的作用
  2. DBDesigner — §V-A 是分段候选枚举和 cost-based 选择的完整算法。理解 DBD 如何自动选择分段键,就能理解手动优化时应该考虑什么
  3. MPP JOIN 策略全解析与优化器决策逻辑 — 分段键如何直接影响 JOIN 策略(Local / Broadcast / Resegment)。本文 §3 和 JOIN 文章的 §1.3、§3.3 重叠,建议对照阅读
  4. Vertica 集群 Rebalance 完全指南 — 分段投影在集群扩缩容时的完整数据迁移流程。§1.2(数据如何移动)和 §1.4(四个阶段)是分段设计的运维视角
  5. Vertica Join 重分段倾斜诊断与修复 — 分段设计不当导致 JOIN 倾斜的完整诊断闭环。§1.3 的四种倾斜来源(Value Skew / Design Skew / 统计信息缺失 / 节点不一致)非常重要
  6. Vertica 表分区策略选择指南 — §1.2 清晰对比了分段与分区的正交关系。§1.4 的分区粒度成本模型有助于理解「分段解决不了的问题」

论文章节引用

  • C-Store 7 Years §3.6 — ring-style 分段方案的完整定义,local segments 用于弹性扩缩容的机制
  • C-Store 7 Years §3.1 — projection 级别的分段设计,super projection 与 narrow projection 的关系
  • C-Store 7 Years §5.2 — buddy projection 的恢复机制,分段在容错中的作用
  • DBDesigner §V-A — 分段候选枚举:从 group-by 和 join 列生成候选分段键,segmentation-friendliness 评估
  • DBDesigner §IV — 设计策略(load-optimized / query-optimized / balanced)如何控制投影数量和存储成本
  • verticaOptimizer §III-A — 分段作为物理属性在优化器搜索空间中的角色(Sorted-Segmented property)

Vertica 的 projection-level 分段设计在 MPP 领域做了一个激进但理性的工程选择:将数据分布从建表时的一次性决策,变成了运维中可渐进调整的策略组合。ring-style 的分层设计(哈希空间等分 → local segments 预分段 → ELASTIC_CLUSTER 整块搬迁)让数据重分布不再等同于全量重写,这是 Vertica 区别于 Greenplum(gpexpand 全表重分布)和 ClickHouse(手动 RESHARD)的关键架构差异。但它的真正洞察不在于技术实现,而在于一个架构哲学的选择——当一张表必须回答多种维度的 JOIN 需求时,你是让优化器去猜测(依赖统计信息),还是让存储去适应(多份物理副本)?Vertica 选择了后者,代价是真实的存储和运维开销。这种 trade-off 没有绝对的对错——它取决于你的查询模式是否足够多维度、是否足够可预测,来让那份额外的存储投资物有所值。