MPP 数据分布策略 —— 从 Hash 分段到一致性哈希的工程实践¶
作者:JiangChong | 撰写时间:2026年04月
适用场景框:当你需要理解为什么 Vertica 在 projection 级别而非表级别定义分段键、为什么一张表可以有多份不同分段的副本、以及如何为不同查询模式选择最合适的分布策略时,这篇文章适合你。
开篇声明:本文以 Vertica 为主要剖析对象,但在所有环节对比其他 MPP 系统(Greenplum、ClickHouse、Doris/StarRocks、Snowflake、Redshift)的不同实现。「MPP 共性」与「Vertica 专属」在文中明确区分——共性机制适用于所有类似架构的分布式数据库,Vertica 专属机制会标注适用范围。
关联文章:
- C-Store 7 Years Later — §3.6 Segmentation 是本文的核心理论来源
- DBDesigner — 分段候选枚举和 cost-based 选择算法
- MPP JOIN 策略全解析与优化器决策逻辑 — 分段键如何直接影响 JOIN 策略(Local/Broadcast/Resegment)
- Vertica 集群 Rebalance 完全指南 — 分段投影在集群扩缩容时的数据迁移机制
- Vertica 表分区策略选择指南 — 分段(跨节点)与分区(节点内)的正交关系
- Vertica Join 重分段倾斜诊断与修复 — 分段键选择不当导致的倾斜问题
理解全文脉络¶
本文从「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 个节点。
如果你将 orders 按 order_id 哈希分布到 4 个节点,customers 按 customer_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 节点集群为例,基本流程:
- 对每行的分段键计算哈希值,得到一个整数(通常 32 或 64 位)
- 将哈希空间映射到各节点——具体的映射算法因系统而异
- 哈希值落在哪个区间,数据就存入哪个节点
比喻:把一副扑克牌按花色分给 4 个人——红心给第一个人,黑桃给第二个人,以此类推。每张牌的去向是确定的、互不重叠的。
Hash 分段有两个关键优点:确定性(每行数据的位置由哈希函数唯一确定,不需要查询目录)和均匀性(只要分段键基数足够高且分布均匀,数据在各节点间接近平均分配)。
各系统的分布算法差异:
| 系统 | 分布算法 | 哈希函数 | 关键特征 |
|---|---|---|---|
| Vertica | Ring-style hash | 内部哈希(64 位空间等分) | 映射规则与节点数解耦,扩容时只需调整区间边界(详见 §2.3) |
| Greenplum | 哈希映射表 | hash(col) → segment mapping |
简洁直观,依赖节点总数 |
| Redshift | 哈希映射表(推测) | 内部哈希 | AWS 托管,用户无需关心算法细节 |
| ClickHouse | half-md5 / murmur | halfMD5(col) % N 或 murmurHash2_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 节点时:
每行的归属由「哈希值落在哪个区间」唯一确定——不需要知道总共有几个节点。当集群从 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) 被切分为:
切分点 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 倍 | 1× | 低 |
| 2 个投影(不同分段键) | 两个 JOIN 维度都能本地执行 | 2 倍 | 2× | 中 |
| 4 个投影(全面覆盖) | 所有 JOIN 都本地执行 | 4 倍 | 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_storage中row_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_id,product_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 策略下的决策过程:
- 从 Q1 提取 Join Cols:
{customer_id},从 Q2 提取 Join Cols:{product_id} - 枚举分段候选:
HASH(customer_id)(收益:Q1 本地 JOIN)、HASH(product_id)(收益:Q2 本地 JOIN)、HASH(sale_id)(无特殊收益但加载均匀) - 优化器对每个候选评估 cost:Q1 用
HASH(customer_id)→ Network cost = 0,Q2 用HASH(product_id)→ Network cost = 0 - 两个候选的 Benefit 都足够高 → DBD 为
sales表创建两个窄投影,分别按customer_id和product_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_id 和 product_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× | 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。
诊断过程:
- 执行计划显示绝大部分步骤在单节点上执行(
Execute on: v_edw_node0039) - 执行计划中出现
Outer (RESEGMENT)(LOCAL ROUND ROBIN)+Inner (RESEGMENT)——两个 RESEGMENT 都是节点内的本地哈希重组,用于 Hash Join 并行分桶 - 两个表的统计信息均已过期(
NO STATISTICS) - 关键发现:表 1 按
statis_date分段;查询 SQL 中恰好有WHERE statis_date = '20200819'——因为所有statis_date = 20200819的数据按哈希分段落在同一节点上,过滤后的结果集全部集中在一个节点 - 表 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 内笔记(按推荐阅读顺序)¶
- C-Store 7 Years Later — §3.6 是本文的核心理论来源,包含 ring-style 分段和 local segments 的完整描述。§3.1 提供了多投影使用频率的真实数据;§5.2 描述了 buddy projection 在分段恢复中的作用
- DBDesigner — §V-A 是分段候选枚举和 cost-based 选择的完整算法。理解 DBD 如何自动选择分段键,就能理解手动优化时应该考虑什么
- MPP JOIN 策略全解析与优化器决策逻辑 — 分段键如何直接影响 JOIN 策略(Local / Broadcast / Resegment)。本文 §3 和 JOIN 文章的 §1.3、§3.3 重叠,建议对照阅读
- Vertica 集群 Rebalance 完全指南 — 分段投影在集群扩缩容时的完整数据迁移流程。§1.2(数据如何移动)和 §1.4(四个阶段)是分段设计的运维视角
- Vertica Join 重分段倾斜诊断与修复 — 分段设计不当导致 JOIN 倾斜的完整诊断闭环。§1.3 的四种倾斜来源(Value Skew / Design Skew / 统计信息缺失 / 节点不一致)非常重要
- 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 没有绝对的对错——它取决于你的查询模式是否足够多维度、是否足够可预测,来让那份额外的存储投资物有所值。