MPP架构详解:从Shared-Nothing到分布键,剖析大规模并行处理的核心原理与工程实践

发布时间:2026/10/8 0:29:13
MPP架构详解:从Shared-Nothing到分布键,剖析大规模并行处理的核心原理与工程实践
MPP 这个词很多人第一次接触到它是在面试或者在处理一个跑了几小时都没出结果的报表时被有经验的同事一句“这表得走 MPP 引擎”给点醒。说实话MPP 不是一个新概念从数据仓库时代它就存在了但这些年随着大数据和云数仓的普及它又一次成了架构选型里绕不开的坎。如果只看定义MPP 就是 Massively Parallel Processing大规模并行处理听起来很直白但真正理解它并不是在说“把任务分成多个并行执行”这么简单而是要先搞明白它解决的是哪一类问题、为什么分布式系统绕了这么多年却依然以 MPP 为骨干。这篇内容我打算从一个实际问题切入当单机计算撑不住的时候MPP 是如何通过一整套分工协作机制扛下来的。我会从架构拆解、平台支持、以及我实际踩过的几个坑出发把 MPP 的“家底”一次说清楚。适合刚接触分布式数据库的读者也适合那些正在做技术选型、想搞清楚 Greenplum、ClickHouse、Redshift 这些引擎底层逻辑的同学。1. 内容整体设计与思路拆解1.1 核心需求解析单机瓶颈到底卡在哪要理解 MPP 存在的意义得先看单机架构走到尽头时暴露的三个瓶颈。硬件层面单台服务器的 CPU 核数和内存带宽是有上限的即使一台 128 核、2TB 内存的机器在处理 TB 级数据的关联聚合时内存带宽和 I/O 也会迅速饱和。软件层面传统单机数据库的查询优化器基于的是集中式执行模型所有数据都要经过一个执行引擎节点数据量上去之后这个节点的 CPU 和网络协议栈会成为绝对的瓶颈。运维层面单机扩容是纵向扩展换一台更强的机器往往意味着停机、迁移、重新压测成本非常高。这三个瓶颈分别对应了 MPP 架构的核心设计目标通过将数据分布到多个计算节点上让每个节点只处理一部分数据并且节点之间通过网络连接协同完成一个查询任务。如果把单机数据库比作一家全能型的小店店员既要收银又要理货还要做售后那么 MPP 就是一家连锁超市每家门店只负责自己片区里的商品总部通过一套调度系统把顾客的需求拆解到对应的门店去执行。这种“拆解”就是 MPP 的精髓也是它和普通分布式系统最本质的区别。1.2 MPP 的定位不是某一个数据库而是一类计算范式很多人容易把 MPP 理解为某款产品比如 Greenplum、Teradata或者干脆有人说 ClickHouse 就是 MPP其实这里有个概念混淆。MPP 是一个计算范式的统称它描述的是“无共享架构下多个节点并行处理同一任务”的软件设计模式。在这个范式之下各家系统有不同的实现方式有基于 PostgreSQL 扩展而来的 Greenplum有纯列式存储的 Vertica有基于 PostgreSQL 的 Citus也有 ClickHouse 这种兼顾列式存储和分布式查询的引擎。理解了这一点再去讨论平台支持就有意义了。因为 MPP 不是一个孤立的软件而是一种架构选型这意味着你在选型时更多是在选“哪套平台对 MPP 的实现更贴合我的业务”而不是在选“哪个数据库支持 MPP”。从这个角度看MPP 的架构拆解反而是更值得花时间的地方因为无论底层是哪家系统它们的核心模型基本一致理解了共性再看差异就是分分钟的事。2. MPP 的核心架构拆解2.1 三个关键角色协调节点、计算节点、存储层一套标准的 MPP 架构里无论外观怎么变内部基本都跑不了这三个角色。协调节点Coordinator有的系统叫 Master、Leader、Query Planner负责接收用户的 SQL做语法解析、逻辑优化、生成执行计划然后把计划分发到各个计算节点上执行。它不存储真正的业务数据或者说只有元数据、统计信息、分布策略这类轻量数据。这里很多人会有个误区以为协调节点就是瓶颈其实在成熟的 MPP 系统里协调节点的作用更接近“指挥中枢”而不是“数据搬运工”。它负责拆任务但不负责算数据数据还是在各个计算节点本地完成的这样协调节点的压力就能控制在一定范围内。计算节点Segment、Worker、Executor是最核心的角色。它们各自持有数据的一部分通常是以分片Shard或分区Partition的形式存储。当一个查询被协调节点拆分后每个计算节点只需要处理自己本地的那份数据这个过程叫本地计算Local Compute也是 MPP 能获得线性扩展能力的基础。存储层在 MPP 里有两种形态一种是存储和计算耦合的本地盘模式像 Greenplum、ClickHouse 集群数据直接落在各节点的本地存储上数据分布策略决定了数据落在哪些节点另一种是存储和计算分离的模式像云数仓 Redshift Spectrum、Snowflake数据放在对象存储上计算节点按需加载数据到本地缓存。后一种在弹性扩缩容上更有优势但数据传输的效率往往成为新的瓶颈这也是为什么云厂商都在搞缓存亲和性和数据本地性优化。2.2 数据分布策略分布键是 MPP 的灵魂如果说协调节点是 MPP 的大脑那分布键就是血液。MPP 的数据分布策略直接决定了后续所有查询的性能表现。以 Greenplum 为例建表时通过 DISTRIBUTED BY 指定分布键数据会根据该键的哈希值散列到各个 Segment 节点上如果不指定系统会默认按第一个字段的哈希分布。这个分布键的选择极其讲究因为它决定了关联查询时数据能否在本地完成。举个例子一个订单表和一个订单明细表如果两张表都用“订单ID”作为分布键那么在关联查询时每个计算节点上都能找到自己本地对应的订单和明细数据整条关联链路完全走本地执行不需要任何跨节点数据传输。反之如果订单表按订单ID分布明细表按商品ID分布那关联时系统就得把明细表的数据按订单ID重新洗牌Redistribute到对应节点上这一下会带来巨大的网络开销查询可能从秒级直接变成分钟级。这种“按分布键对齐”的设计在 MPP 术语里叫 co-located join。不仅是关联分组聚合、去重这类操作同样受分布键影响。比如要按用户维度做聚合用户ID作为分布键时每个节点只需要聚合本地的那部分用户完全不需要 shuffle 数据。所以分布键的选择绝不是建表时随便填一个字段那么简单它需要结合业务最常见的查询模式来做设计。这也是我后面要重点讲的实操要点之一。2.3 为什么 Shared-Nothing 架构是 MPP 的主流MPP 领域最经典的架构分类是 Shared-Nothing 和 Shared-Disk。Shared-Disk 指的是所有计算节点共享同一套存储系统比如 Oracle RAC数据只有一份所有节点都能访问但为了保证一致性需要额外的锁机制和缓存融合机制这在高并发下非常容易出现阻塞。Shared-Nothing 指的是每个节点独享自己的 CPU、内存和磁盘节点之间只通过网络交换数据不存在共享存储的竞争问题。MPP 数据库绝大多数都采用了 Shared-Nothing原因很直接第一扩展性更好加一个节点数据就多一份分布落点计算能力跟着线性增长不需要考虑存储层面的共享瓶颈第二数据本地性更好每个节点只需要处理本地数据不需要远程读取I/O 延迟大幅降低第三故障隔离更好某个节点宕机其他节点还可以继续工作配合副本机制可以实现高可用。当然 Shared-Nothing 也有代价。数据需要冗余多份存储成本高一些节点间通信依赖网络质量万兆网卡和低延迟交换机几乎是标配数据重分布的操作代价也比较大比如集群扩容时要重新平衡数据如果节点间传输的数据量巨大这个过程可能持续数小时。但整体上看Shared-Nothing 在大规模并行计算场景下依然是性价比最高的选择这也是从 Teradata 到 Greenplum再到云上 Snowflake 都坚持这一架构的根本原因。2.4 控制平面与数据平面的分离成熟的 MPP 系统在设计上会把控制平面和数据平面分开。控制平面处理元数据、锁管理、会话管理、调度决策数据平面负责数据的存储、传输和计算。这种分离带来的好处非常实际控制平面负载很轻即使集群规模很大协调节点也能轻松应对数据平面则可以充分水平扩展每个计算节点独立处理自己的数据分片互不干扰。以 Greenplum 为例控制平面由 Master 节点承担数据平面由多个 Segment 节点构成。Master 节点的硬件要求反而没有那么夸张因为它的职责只是接收 SQL、生成计划、汇总结果真正耗资源的计算都发生在 Segment 上。这个架构设计的另一个好处是在做性能调优时定位问题会清晰很多如果查询计划阶段耗时高问题大概率在控制平面如果执行阶段耗时高再去看数据平面的节点负载和网络传输。3. 平台支持全景从传统数仓到大数据库生态再到云上数仓3.1 传统数仓阵营里MPP 是怎么站稳脚跟的MPP 在传统数据仓库领域的代表有 Teradata、Vertica、Greenplum、Netezza 等。Teradata 是鼻祖级的存在1990 年代就开始大规模应用在金融、电信、零售行业它的节点间通信机制和优化器设计影响了后来很多产品。Teradata 的架构是典型的 Shared-Nothing每个节点也叫 AMP数据按主索引分布在各个 AMP 上这个设计和 Greenplum 的分布键几乎是一个逻辑。Vertica 很有意思它是列式存储和 MPP 结合的典范最早源自 C-Store 研究项目后来被 HP 收购现在是 Micro Focus 旗下的产品。Vertica 最大的特点是压缩率极高列式存储配合高级压缩算法相同的数据量比行式存储省 5-10 倍空间在这个基础上再跑 MPP 查询I/O 量就明显小下去。Greenplum 则是开源领域最有代表性的 MPP 数据库它基于 PostgreSQL 改造兼容 PostgreSQL 语法这大概是它能在互联网公司大规模落地的重要原因。Greenplum 在 6.x 版本后引入了增强的 ORCA 优化器对复杂查询的优化能力大幅提升。但要注意Greenplum 更像是一个分析型数仓OLTP 类的高并发点查并不是它的强项。Netezza 则被 IBM 收购它独特的地方在于使用了 FPGA 加速在硬件层面做数据过滤和聚合这个设计理念至今仍有参考意义。3.2 大数据生态里的 MPPClickHouse、Trino、Impala大数据生态里挂 MPP 名字的系统更多ClickHouse、Trino前身 Presto、Impala 都常被归入 MPP 阵营但它们的实现风格差别很大。ClickHouse 是俄罗斯公司 Yandex 开源的列式数据库它的 MPP 实现方式非常激进。ClickHouse 默认不做跨节点的数据关联它更鼓励通过预先设计好的分布式表和大宽表来规避 shuffle。ClickHouse 的分布式查询走的是“每个分片独立执行完再由协调节点合并”的模式如果查询涉及跨分片的数据关联性能会退化得非常明显。所以 ClickHouse 在多数场景下是被当作“单机性能极强的分布式聚合引擎”来使用而不是一个完整的分布式数据仓库。Trino 则走的是另一条路线它定位是分布式 SQL 查询引擎本身不存储数据数据源可以接 Hive、对象存储、MySQL、PostgreSQL 等。它的 MPP 能力体现在查询执行阶段数据从数据源并行读入在内存里完成 join 和 aggregation。Trino 的优势是灵活能跨多种数据源做联邦查询劣势则是它完全依赖网络传输数据如果数据源到引擎的网络带宽不够查询性能就上不去。Impala 是 Cloudera 推出的查询引擎和 Trino 类似也依赖底层存储如 HDFS 或 S3。Impala 的无共享架构在查询优化和并行执行上做得不错尤其在搭配 Kudu 时可以实现分钟级数据更新的分析场景。但它对内存的要求很高深度分页或者大聚合时内存溢出是常见的事故源头。3.3 云数仓对 MPP 的重塑Snowflake、Redshift、BigQuery云数仓把 MPP 从“物理集群”变成了“虚拟资源池”这是对传统架构的一次重大修正。Snowflake 的做法最有代表性存储层用对象存储计算层用虚拟 warehouse多个计算集群可以同时挂载在同一个存储层上互不共享计算资源但共享同一份数据。这种架构下扩缩容只是启停虚拟计算集群的事按需计费而且读写隔离做得很好不会出现一个慢查询拖垮其他查询的情况。Redshift 是 AWS 的老牌数仓它的核心用的是基于 PostgreSQL 改造的 MPP 引擎。Redshift 早期是典型的 Shared-Nothing节点挂本地存储后来推出 RA3 节点类型把数据落地到 S3 并引入本地缓存走向了存储计算分离的方向。Redshift 的分布键和排序键设计非常关键和 Greenplum 的分布键是同一个逻辑但 Redshift 还额外增加了分布方式的选择比如 ALL 分布适合小表广播EVEN 分布适合按顺序轮询KEY 分布则是哈希分布。BigQuery 则是 Google 的云数仓它的 MPP 能力藏在柱状存储和分布式执行引擎 Dremel 的背后。BigQuery 对用户屏蔽了分片和节点概念用户只管写 SQL系统自动调度执行。它的优化器依赖统计信息自动选择执行策略因此用户不需要手动指定分布键但反过来也意味着如果表数据的统计信息陈旧执行计划可能不理想。云数仓的共性趋势是存储和计算分离、弹性扩缩容、按量计费这些本质上都是 MPP 架构的延伸和优化只是把资源管理的维度从物理节点提升到了虚拟资源池。对用户来说运维复杂度降低了很多但架构理解和查询优化的工作量并没有减少反而因为屏蔽了底层细节更考验工程师对执行计划的把握。3.4 一个真实例子MPP 引擎处理一条查询的全流程为了更直观地理解 MPP 的工作流程我用 Greenplum 处理一条 SQL 来串一遍整个过程。假设有两张表用户表包含用户ID、注册城市、注册时间订单表包含订单ID、用户ID、订单金额、下单时间。查询需求是统计每个城市的订单总额。这条 SQL 到达 Master 节点后先是解析和验证生成语法树然后走优化器。优化器会做两件事代价估算和执行计划生成。这里关键的点是优化器需要知道两张表的分布键是什么。如果用户表按用户ID分布订单表也按用户ID分布那优化器会认为这种 join 可以走本地关联于是采用 co-located join 策略每个 Segment 节点只处理本地那部分用户的订单。之后按注册城市做分组聚合每个 Segment 节点本地做一次聚合得到部分结果然后把结果发回 MasterMaster 再做最终合并返回给客户端。如果分布键不匹配执行计划里就会多出一步 Redistribute 或者 Broadcast。Redistribute 是把一张表的数据按目标字段重新哈希后发送到对应节点Broadcast 则是把小表复制到所有节点。这两种数据移动都会增大网络开销也是 MPP 查询变慢最常见的可视化原因。通过 EXPLAIN 查看执行计划时如果看到运动节点Motion非常多就该怀疑分布键设计是否有问题了。另外MPP 执行还有一个值得注意的细节——物化中间结果。有些 MPP 引擎会把 shuffle 的中间结果写磁盘Greenplum 在内存不足时也会做磁盘溢出这会进一步放大 I/O 开销。所以控制中间结果集的大小比如尽早做过滤、避免 SELECT 全字段、尽量在聚合前压缩数据量是 MPP 查询优化的重要思路。4. 实操过程与核心环节实现4.1 环境选型与节点规划在规划一套 MPP 集群时第一个要决定的是“要不要上 MPP”而不是“上哪个 MPP”。当数据量在几百 GB 到几个 TB 级别、查询模式以 SQL 聚合分析为主、对实时写入没有极端要求时MPP 是一个非常合理的选择。但如果数据量只有几十 GB且频繁的并发点查占比很高传统的关系型数据库或者单机 PostgreSQL 反而是更好的选择。选型确定之后节点规划有一些实战经验可以分享。计算节点数量和 CPU 核数的配比要结合查询复杂度来看。一般建议单个 Segment 节点分配 8 核到 16 核内存和 CPU 核心数按 8GB 到 16GB 每核心来配。比如一个 4 节点的 Greenplum 集群每个节点 16 核、128GB 内存总共 64 核、512GB 内存这样的规模应对 5TB 左右的数仓数据是够用的。另外要考虑网络。MPP 集群的节点间通信对网络延迟和带宽非常敏感强烈建议上万兆网络。我见过一个集群因为用了千兆网络结果在数据重分布阶段网络成了瓶颈整个集群跑一个跨节点的 join 都要十几分钟后来换成万兆网络同一个查询只需要不到一分钟这个差距非常夸张。4.2 建表语句里的分布键设计现场在建表时分布键的选择需要结合业务的实际查询模式。一个来自生产环境的经验做法是找出查询频率最高的三张表和它们最常用的关联条件把这些关联字段作为分布键的首选。比如前面的用户表和订单表如果订单表是事实表用户表是维度表那么订单表按用户ID分布用户表也按用户ID分布就能让最常见的关联查询走本地关联。另一个关键技巧是处理数据倾斜。如果分布键的取值分布不均比如用户ID集中在少数几个值上例如某个头部用户贡献了绝大多数订单那么这些热点数据会全部落到同一个节点上导致该节点的负载远高于其他节点形成“木桶效应”。解决方法是使用复合分布键或者在前缀字段上增加一个随机因子也可以在建模时把大用户的数据单独分桶处理。简单的检查方法是执行一个“SELECT 分布键, COUNT(*) FROM 表 GROUP BY 分布键”的查询观察各个取值的数据量如果发现有明显的长尾分布就要考虑调整分布策略。在 Redshift 里还有一个 ALL 分布的选择。对于数据量较小的维度表比如几百 MB 的国家表、城市表直接用 ALL 分布在每个节点都放一份全量数据join 时就可以完全避免广播操作。这个技巧虽然简单但在实际优化中的效果非常显著。4.3 查询优化从执行计划里找运动节点在实际做性能调优时第一件事永远是看执行计划。以 Greenplum 为例EXPLAIN ANALYZE 输出中的关键信息有三个每步操作的行数估算和实际行数、Motion 的类型和行数、以及每个节点的执行耗时。Motion 在 Greenplum 执行计划里就是数据移动的标志如果在计划里看到很多 Redistribute Motion 或者 Broadcast Motion且涉及的行数很大那这个查询必然快不了。一个经常被忽略的问题是统计信息的时效性。MPP 优化器依赖统计信息来做代价估算如果一张千万级的表从创建后就没跑过 ANALYZE优化器可能认为它是空表或者在过滤条件上严重低估行数导致选错 join 顺序。比如一张大事实表和维度表关联维度表过滤后的数据量被低估了 10 倍优化器就可能选择把维度表广播到所有节点产生大量无效数据传输。解决办法是定期对关键表执行 ANALYZE以及在大批量数据导入后立刻做一次统计信息更新。另外MPP 查询中常见的高成本操作是“带 ORDER BY 的聚合”和“窗口函数”。窗口函数在 MPP 下的执行往往需要把数据按窗口字段重分布这是一个全量数据洗牌的过程。如果窗口函数用得太随意比如 OVER (PARTITION BY 某个高基数字段)那代价极高。优化思路是提前过滤数据、只在所需的数据范围上做窗口计算或者拆分成更小的数据集分别处理再合并结果。4.4 资源管理队列与并发控制MPP 集群的资源不是一个无限平摊的资源池每个查询都会消耗一定配额。Greenplum 里的资源队列Resource Queue就是用来控制并发的它定义了每个队列的 CPU、内存、并发数量限制。生产环境里常犯的错误是把所有查询都放进默认队列也不限制并发数结果几个大查询同时跑内存直接耗尽集群状态瞬间变得不可用。合适的做法是把查询按业务优先级分成几个队列实时报表一个队列并发数限制在 5 个以内内存配额给足批量任务一个队列并发数 2-3 个执行时间长也没关系临时查询一个队列优先级最低。这样即使有人提交了一个跑 2 小时的大查询也不会把实时报表的通道堵死。Redshift 里也有对应的 WLMWorkload Management队列配置原理完全一致。还有一个实践中很有效的做法是设置 statement_mem 或类似的查询内存限制防止单个查询吃掉整个节点内存。大查询如果内存不足可以接受它落磁盘来换稳定性但绝不能让一个失控查询把整个集群搞挂。4.5 常见问题与排查技巧实录MPP 集群的常见故障我按大类整理了一份速查表现象可能原因排查手段解决思路查询整体变慢但无报错数据分布倾斜查看各节点 CPU/IO 负载调整分布键增加随机前缀执行计划中的 Motion 行数异常大分布键不匹配或统计信息过期EXPLAIN ANALYZE 对比估算与真实行数更新统计信息拆分子查询单节点内存溢出OOM并发过高或查询中间结果过大查看资源队列活跃查询限制并发拆分大查询集群扩容后数据长期不均衡扩缩容后数据重分布未完成查看系统表 rebalance 状态手动执行数据重分布任务高并发点查性能差未遵循 MPP 设计规则查看是否触发了全表扫描在维度表和主键上建立索引或更换为 OLTP 类数据库跨库关联频繁出现 Broadcast维度表未用 ALL 分布查看执行计划中的 Broadcast Motion小维度表改为 ALL 分布这里面最值得强调的是第一个问题——数据倾斜。倾斜问题的隐蔽性强集群看起来所有节点都在工作但只有一个节点负载接近 100%其他节点闲得发呆。排查起来也简单Ganglia、Prometheus 这类监控工具看节点负载图对比各节点 CPU 曲线的差异如果一条线竖得老高而两侧都是平线基本就是倾斜无疑了。之前在生产环境处理过一个问题一个按商品维度统计销量的查询某个头部商品的数据占了全表 35%这一个值让单节点负载比平均高出 4 倍查询从预期的 30 秒拖到了 5 分钟。后来把分布键改成了“商品ID 一个城市字段”的复合键让大数据量的商品也能在不同节点上散开查询恢复到了 40 秒以内。第二个要提的坑是“MPP 未必比单机快”。当数据量不大、单机完全能放下时MPP 的网络通信开销反而会让性能更差。我有一次在同一套数据上对比测试数据量 200GBMPP 集群是 4 节点单机是 64 核大内存。一个中等复杂度的 join 聚合查询MPP 跑了 45 秒单机 PostgreSQL 只跑了 28 秒。原因很简单MPP 执行计划里有两处 Redistribute Motion数据传输占了大半时间而这个查询在单机上完全不需要移动数据。所以MPP 的“快”是有前提的数据量足够大到单机无法合理承载或者查询本身能通过并行化受益。数据量小的时候不要迷信分布式。第三个是关于并发和慢查询互相干扰的典型事故。某个周五晚上一条数据清洗 SQL 和线上报表查询同时运行清洗任务一次性 UPDATE 了全表 80% 的行这个 UPDATE 在 MPP 里会生成巨大的中间结果并且锁定大量行导致线上报表查询被阻塞。最终的结果是清洗任务运行了 3 小时期间线上报表一直超时。事后复盘下来核心问题是没做资源隔离和锁管理。后来把清洗任务放到专门的维护窗口并且提交前加上超时控制和锁等待时间限制这类问题就再没出现过。5. 继续深挖的两个方向其实聊到这里MPP 的基本盘已经说完了架构上理解 Shared-Nothing 和分布键平台上理解各家的差异实操上理解执行计划和资源管理。但如果真想在这一块继续深入还有两个方向值得花时间。第一个方向是“MPP 和 NewSQL”的边界。MPP 在分析型场景里有统治力但在高并发事务场景里表现不理想。TiDB、OceanBase 这类 NewSQL 数据库虽然也用了分布式存储和计算分离的思路但它们的目标是兼顾 OLTP 和 OLAP执行引擎的设计和 MPP 有本质差异。理解这个边界对做技术选型的人非常重要——不是说分布式数据库就能解决所有问题选型的关键前提要厘清“这个系统处理的是什么类型的查询”。第二个方向是“湖仓一体”里的 MPP 角色。数据湖和数据仓库的界限在模糊像 Trino、Hive 这类引擎已经可以把数据湖文件当作源来做 MPP 查询Doris、StarRocks 也把数据湖联邦查询做成了常态化能力。在这个背景下MPP 的定位不再局限于数仓内部它正在变成一个统一查询计算层。这个演进非常快值得持续跟踪。从我个人实际应用的角度来说最好的学习方式不是在文档里背概念而是找一个真实的业务场景拿一套小规模的 MPP 环境亲手建表、设计分布键、跑 EXPLAIN ANALYZE 去观察 Motion再故意选错一次分布键感受一下性能回退。这个过程走一遍你对 MPP 的理解会远超读十篇文章。最后再分享一个小经验MPP 集群里执行计划里的 Motion 数量永远和性能成反比优化目标就是尽可能减少数据移动所有的分布键设计、SQL 改写、统计信息更新本质上都是围绕这一件事在转。