十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

# Presto 查询引擎内核详解:AddExchanges——基于物理属性的全局数据分布规划

# Presto 查询引擎内核详解:AddExchanges——基于物理属性的全局数据分布规划 AddExchanges — Global Data Distribution Planning Based on Physical Properties引言在 Presto 的分布式执行引擎中查询优化器在将逻辑计划转换为物理执行计划时面临一个核心问题如何确保每个算子都能获得符合其执行要求的数据分布数据分布指的是数据在分布式环境中的组织和路由方式——数据是否按某列分区、相同 Key 是否路由到同一分区、是否需要汇聚到单节点或复制到多节点。不同算子对数据分布有截然不同的要求例如Aggregation 要求相同 Group Key 的数据进入同一个 partition从而保证同一个 Group 的数据能够由同一个执行实例完成聚合Join 根据执行策略可能要求两侧数据具有兼容的 Partitioning也可能通过 Broadcast 将一侧数据复制到参与执行的 WorkerSingle-node execution 则要求所有输入数据最终汇聚到同一个 Worker。然而上游算子产生的数据分布几乎不可能天然满足下游所有算子的需求这种数据分布不匹配正是 AddExchanges 要解决的核心问题。本文将从问题动机、属性机制、属性推导和具体实现四个层面系统分析 Presto AddExchanges 如何基于物理属性完成全局数据分布规划。一、数据分布不匹配AddExchanges 要解决什么问题1.1 典型场景聚合查询考虑如下的聚合查询SELECT a, b, COUNT(*) FROM t GROUP BY a, b;对于 Aggregation 来说它在分布式执行层面的核心要求是所有具有相同 (a,b) 的数据行必须路由到同一个分区保证同一 Group 由同一个执行实例完整处理Need: PARTITIONED BY (a, b)然而上游 TableScan 提供的数据分布很可能是任意的例如Have: Unknown distribution Worker 1: (1,1), (2,2), (1,2) Worker 2: (1,1), (2,1), (3,3) Worker 3: (2,2), (3,3), (1,2)相同 (a,b) 的数据分散在不同 Worker无法直接进行最终聚合。因此需要在中间插入 Remote Exchange 重新组织数据分布Aggregation (GROUP BY a, b) ▲ │ 需要: PARTITIONED BY (a, b) │ Remote Partitioned Exchange ←── 在此插入 ▲ │ 实际: Unknown Distribution │ TableScanRemote Exchange 通过网络将数据重新路由到目标 Worker使相同 (a,b) 的数据进入同一个 partition。因此Remote Exchange 改变的不是数据的逻辑内容而是数据在分布式执行环境中的物理分布和路由方式。1.2 AddExchanges 的职责定位AddExchanges 是 Presto 物理规划阶段负责全局数据分布规划的核心组件具体职责包括根据算子执行语义确定其对输入的数据分布需求规划子节点并判断其实际具备的数据分布是否满足需求在必要时插入合适类型的 Remote Exchange建立所需的数据分布基于新生成的计划推导相应物理属性使后续 Optimizer 能够基于新的属性继续进行规划。核心原则与 AddLocalExchanges 一样AddExchanges 的主要职责是根据算子的物理属性需求建立满足执行语义的数据分布而不是直接进行 Exchange 的成本优化。Exchange 是否可以进一步简化或消除则由后续优化阶段处理。1.3 Exchange 规划的层次划分在 Presto 的物理计划优化阶段Exchange 相关处理由多个优化规则共同完成各阶段关注的问题不同组件关注层次Exchange 类型核心问题AddExchangesWorker 之间Remote Exchange数据如何在 Worker / Partition 之间重新分布AddLocalExchangesWorker 内部Local ExchangeWorker 内部不同数据流之间如何分发Exchange Planning │ ┌─────────────┴─────────────┐ │ │ AddExchanges AddLocalExchanges │ │ ▼ ▼ Worker 之间的数据分布 Worker 内部的数据流组织 │ │ ▼ ▼ Remote Exchange Local Exchange这种分层设计使每个组件聚焦于自己层次的问题逻辑保持清晰专注。但仅仅知道需要重新分布还不够。Optimizer 必须回答两个问题当前算子需要什么样的数据分布Child 当前实际具备什么样的数据分布Presto 通过 PreferredProperties 和 ActualProperties 分别描述这两类信息——这正是下一章的主题。二、属性驱动的数据分布规划机制AddExchanges 的核心并不是针对某个 Operator 固定插入某一种 Exchange而是一套数据分布约束机制Data Distribution Property Enforcement分析 Operator 对输入数据分布的要求通过这些要求自顶向下规划 ChildChild 规划完成后再结合其 ActualProperties 判断需求是否已经满足必要时插入 Remote Exchange。2.1 核心抽象PreferredProperties 与 ActualProperties这套机制建立在两个互补的物理属性概念之上概念方向含义示例PreferredPropertiesTop-down需求Parent 希望 Child 具有的分布PARTITIONED BY (a,b)ActualPropertiesBottom-up实际Child 规划完成后实际具备的分布PARTITIONED BY (a,b) 或 SINGLE需要注意PreferredProperties 更准确的理解是 Parent 对 Child 的物理属性偏好它首先用于引导 Child 的自顶向下规划Child 会尽可能生成满足偏好的计划但并不保证成功规划完成后AddExchanges 再依据 ActualProperties 判断偏好是否被满足并在必要时进行 Enforcement。整个机制形成一个闭环Parent Operator │ │ ① 构造 PreferredProperties表达需求 ▼ planChild(preferred) │ │ ② 自顶向下规划 Child ▼ Child Plan │ │ ③ PropertyDerivations 自底向上推导属性 ▼ ActualProperties描述实际 │ ▼ Property Matching需求 vs 实际 │ ┌────┴────┐ │ │ 满足 不满足 │ │ ▼ ▼ 直接复用 插入 Remote Exchange │ ▼ 建立所需分布 │ ▼ 推导新的 ActualPropertiesExchange 并非预先固定在某个算子前而是当且仅当需求无法被现有计划满足时作为强制转换机制动态产生。2.2 PreferredProperties需求如何表达PreferredProperties 是 Parent 算子向优化器表达其对 Child 数据分布期望的载体PreferredProperties ├── Global Properties跨节点 / 分区分布 │ ├── Partitioned on columns K按 K 分区 │ ├── Single单节点汇聚 │ └── Any无特定要求 └── Local Properties分区内分布 ├── Grouped on columns K按 K 分组 └── Sorted on columns K按 K 排序以 Aggregation 为例Aggregation (GROUP BY a, b) │ ▼ PreferredProperties: Global: PARTITIONED BY (a, b) Local: GROUPED BY (a, b)同时该 PreferredProperties 会与 Parent 的 PreferredProperties 合并使相邻算子的属性要求尽可能共同满足避免不必要的 Exchange。2.3 ActualProperties实际如何描述ActualProperties 描述算子在完成规划后实际具备的物理属性ActualProperties ├── Global Properties │ ├── nodePartitioning 数据在多个节点之间是如何分布的 │ └── streamPartitioning 数据在多个流 / 切片之间是如何分布的 ├── Local Properties 数据在单个流 / 切片内部的特性 │ ├── Grouped on columns K │ └── Sorted on columns K └── Constants对于 AddExchangesGlobal Properties 是最核心的属性因为它描述了数据在分布式执行环境中的节点级和数据流级分布是判断是否需要 Remote Exchange 的主要依据。2.4 Property Matching满足而非相等Presto 的属性匹配并不是两个属性之间的简单相等判断而是判断 Child 的 ActualProperties 是否能够满足Parent 的需求。例如匹配场景示例结果完全匹配Need(a,b)/ Have(a,b)匹配可满足的属性Have 的 partitioning 能够满足 Need匹配实际分布单节点NeedPARTITIONED/ HaveSINGLE匹配不兼容Need(a)/ Have(b)不匹配其中SINGLE 满足 PARTITIONED 需求最能体现这种语义所有数据已经在同一个 partition 中无论要求按什么列分区需求都天然成立。另一个需要了解的实现事实是AddExchanges 中并不存在一个统一的属性匹配入口。实际的匹配与 Enforcement 决策分散在AddExchanges.Rewriter针对不同 PlanNode 类型的visitXxx()方法中由各类算子的具体执行语义决定。因此理解 AddExchanges 的关键不是寻找一个统一的match()方法而是结合具体的visitXxx()方法看它如何构造 PreferredProperties、规划 Child以及如何依据 ActualProperties 决定是否插入 Remote Exchange。三、PropertyDerivations实际属性的自底向上推导如果无法准确知道 Child 实际具备什么属性Property Matching 就无从谈起。PropertyDerivations 是 Presto 用于自底向上推导 PlanNode 实际物理属性ActualProperties的核心机制。3.1 递归推导框架PropertyDerivations 的主体入口是derivePropertiesRecursively( PlanNode node, Metadata metadata, Session session)该方法自底向上遍历 Plan Tree对于当前 node先递归处理所有 source 得到各自的实际属性再调用deriveProperties(node, inputProperties, ...)根据当前节点语义与 Child 属性推导当前节点的 ActualPropertiesParent │ ┌────┴────┐ ▼ ▼ Child1 Child2 │ │ ▼ ▼ ActualProps ActualProps │ │ └────┬────┘ ▼ deriveProperties() │ ▼ Parent ActualProps这里需要区分两个层次derivePropertiesRecursively()负责整个 Plan Tree 的递归遍历PropertyDerivations.Visitor针对具体 PlanNode根据已推导出的 Child ActualProperties 计算当前节点的属性本身不负责递归。3.2 针对不同 PlanNode 的推导语义deriveProperties()创建 PropertyDerivations.Visitor通过node.accept(visitor, inputProperties)分派到对应 PlanNode 的visitXxx()方法。不同 PlanNode 对物理属性具有不同的语义TableScan根据 Connector 提供的 TableLayout 推导数据的初始分布Exchange根据 Exchange 自身的 partitioning scheme 推导其输出分布Aggregation、Join 等算子根据自身执行阶段与 Child 的 ActualProperties 推导输出属性。Visitor 并不是一套独立的属性系统而是 ActualProperties 推导框架针对不同 PlanNode 的具体实现。3.3 与 AddExchanges 的关系规划与推导的分工PropertyDerivations 与 AddExchanges 不是两个并列的规划组件而是规划机制与属性推导能力的分工AddExchanges 负责整个数据分布规划过程——构造需求、自顶向下规划 Child、匹配需求与实际、必要时插入 ExchangePropertyDerivations 则是通用的物理属性推导机制被 AddExchanges 在规划过程中调用以获取 Child 的 ActualProperties。PropertyDerivations 本身不决定是否插入 Remote Exchange。四、实例AggregationNode 的完整规划流程4.1 场景与需求AggregationNode 是理解 AddExchanges 工作机制的典型例子。对于如下的聚合查询SELECT customer, SUM(amount) FROM orders GROUP BY customer;在分布式执行环境中相同 customer 的数据可能分散在不同 Worker 上Worker 1: (A,100), (B,200), (C,150) Worker 2: (A,300), (B,250), (D,100) Worker 3: (B,50), (C,400), (D,200)如果直接执行 Aggregation相同 customer 会产生多个局部聚合结果。因此Aggregation 的核心分布要求是相同 grouping key 的数据必须进入同一个 partition。需要强调的是本章讨论的是 AddExchanges 阶段的物理属性规划。此时 AggregationNode 仍然作为一个完整的 PlanNode 参与 Property Planning并不存在 Partial / Final Aggregation 的拆分。本节关注的是 Aggregation 对 Child 提出的 partitioning requirement以及 AddExchanges 如何根据 Child 的 ActualProperties 判断是否需要插入 Exchange。4.2 构造需求partitioned by customer在visitAggregation()中首先根据 grouping keys 构造 preferred propertiesSetVariableReferenceExpression groupingKeys node.getGroupingKeys(); // groupingKeys {customer} PreferredProperties preferred PreferredProperties.partitionedWithLocal( partitioningRequirement.forKeys(groupingKeys), grouped(groupingKeys) );这意味着 Aggregation 希望 childGlobal: partitioned by customer Local: grouped by customer其中partitioned by customer 是 AddExchanges 阶段重点处理的分布属性grouped by customer 属于 Local Property主要由后续的 AddLocalExchanges 处理。同时该偏好会与 Parent 的 PreferredProperties 合并使相邻算子的属性要求尽可能共同满足避免不必要的 Exchange。4.3 递归规划与几种实际情况构造好 preferredProperties 后执行PlanWithProperties child planChild(node, preferredProperties); ActualProperties childProperties child.getProperties();此时得到了明确的 Preferred 与 Actual 信息随后根据 Child 的 ActualProperties 判断是否满足 partitioning requirement。可能呈现如下几种情况情况Child ActualProperties是否满足 Property Requirement1PARTITIONED BY customer满足2SINGLE满足3PARTITIONED BY order_date通常无法满足4UNKNOWN / ARBITRARY不满足情况 2 值得再次强调SINGLE 分布意味着所有数据已经位于同一个 partition 中因此天然满足PARTITIONED BY customer——不管 customer 是什么所有数据都在同一个地方。这正是 2.4 节满足而非相等的最好例证。4.4 Enforcement插入 Remote Exchange如果 Child 当前的实际分布无法满足需求例如!isStreamPartitionedOn(customer) 并且 !isNodePartitionedOn(customer)则 AddExchanges 会插入 Remote Partitioned Exchange从而形成Aggregation ▲ │ partitioned by customer │ Remote Partitioned Exchange hash(customer) ▲ ┌────────────┼────────────┐ │ │ │ W1 W2 W3Exchange 改变数据在 Worker 之间的分布使相同 customer 的数据进入同一个 partition建立 Aggregation 所需的 Global Property。插入 Exchange 后Exchange 产生的数据分布会被重新推导为相应的 ActualProperties其上层 Aggregation 由此能够看到满足要求的 partitioning property。4.5 机制回顾Aggregation 的完整规划路径就是 2.1 节那个闭环在具体算子上的展开构造partitioned(customer)需求 →planChild()自顶向下传递 → 自底向上推导 ActualProperties → Property Matching → 满足则直接复用不满足则插入 Remote Partitioned Exchange 并推导新的物理属性。其中最重要的是三个相互配合的机制自顶向下Top-downOperator 根据执行需求构造 PreferredProperties通过 planChild() 传递给 child引导其尽可能产生满足要求的物理属性自底向上Bottom-upchild 规划完成后根据实际生成的 Plan 结构推导其 ActualProperties属性匹配Property Matching比较 ActualProperties 与物理需求满足则复用否则通过 Exchange 强制建立缺失的物理属性。这也是 AddExchanges 与简单的规则式 Exchange 插入机制之间最重要的区别Exchange 并不是预先固定在某个 Operator 前面而是作为物理属性需求无法通过现有 Plan 满足时的 enforcement mechanism 动态产生。五、总结Exchange 是手段Property 是核心AddExchanges 表面上负责 Exchange 的规划与插入但其真正解决的问题不是什么时候插入 Exchange而是如何让分布式执行计划满足 Operator 所需要的物理属性。在这一过程中PreferredProperties 表达 Operator 对 Child 的物理属性期望PropertyDerivations 从实际 Plan 结构中推导 Child 的 ActualPropertiesProperty Matching 判断两者是否满足而 Exchange 则作为 Property Enforcement 的手段在必要时改变数据分布。几个机制共同构成一个完整的 Property Planning 闭环。这套设计的关键价值在于它将“Operator 需要什么”与“Child 已经有什么”解耦Optimizer 不再针对每一种 Operator 固定定义 Exchange 插入规则而是通过统一的 Property System 描述、推导和匹配物理执行约束。当 Child 已具备满足要求的属性时直接复用已有分布无需再次 Shuffle当属性无法满足时才通过 Exchange 等 Enforcement 机制进行调整。从数据库系统设计的更宏观视角看这不是 Presto 的独创而是一脉相承的经典设计从 Volcano/Cascades 优化器的 Physical Property 与 Enforcer 模式到 SQL Server、Orca 以及 Calcite 等系统中的物理属性框架现代优化器通常都将物理属性纳入优化过程用统一的 Property System 描述和满足物理执行约束。Presto 的 AddExchanges则是这一思想在分布式执行场景下的具体落地。因此AddExchanges 真正重要的地方并不在于它“插入了多少个 Exchange”而在于它完成了一个更深层次的抽象将分布式执行所隐含的数据分布约束从具体 Operator 的执行逻辑中提取出来转化为 Optimizer 可以显式描述、推导、匹配和优化的 Physical Properties。理解 AddExchanges 最重要的也不是记住某个 visitXxx() 在什么情况下插入什么 Exchange而是理解它背后的规划思想Exchange 是手段Property 是核心。AddExchanges 并不是一个简单的 Exchange 插入器而是 Presto 将逻辑计算需求映射为物理执行约束的关键规划机制Operator 提出 Property 需求Optimizer 推导并匹配已有 PropertyExchange 则在必要时承担 Enforcement。最终抽象的逻辑计划被逐步落实为一个满足数据分布约束、能够真正并行执行的分布式物理计划。作者王冬PrestoDB Committer | Presto Iceberg Code OwnerGitHub: https://github.com/hantangwangdEmail: mingwbdgmail.com本文章同步发表于https://hantangwangd.github.io/zh/posts/2026-08-28-add-exchanges.html
返回列表