资讯动态

Presto 查询引擎内核详解:分布式规划机制

发布时间:2026/8/11 4:46:26 来源:尧图企业网站定制
从 Physical Plan 到 Distributed Plan一、引言为什么需要分布式规划在前面的文章中我们了解到 Presto 如何完成SQL -- AST -- Analysis -- LogicalPlan -- PhysicalPlan优化器最终生成的是一棵 Physical Plan TreePlanNode Tree。例如Output | Exchange | Join / \ Scan Orders Exchange \ Scan Customer从单机执行角度看这棵树已经描述了数据如何读取算子如何连接数据如何流动但是对于一个 MPP 查询引擎而言集群中的 Worker 节点无法直接执行整棵计划树。原因在于数据需要跨节点重新分布不同计算阶段需要分布到不同 Worker查询需要拆分为多个独立调度单元不同阶段之间需要建立依赖关系。因此Presto 需要完成一次关键转换PhysicalPlan -- Distributed Plan这个过程就被称之为Distributed Planning分布式规划完成这一转换的核心组件是 PlanFragmenter它负责将 Physical Plan 切分为多个 PlanFragment并构建后续分布式调度所需的 Distributed Plan。二、整体认知Presto 的分布式规划及调度执行模型在Presto中一个SQL查询从优化后的物理计划转化为可执行的分布式任务需要经历几个核心阶段将物理计划切分为计划片段、规划各个片段的分布式执行特性、为分布式计划构造调度器并实际调度和执行。在深入细节之前我们需要先建立一个分层的全局认知。Presto 是一个以 Pipeline Streaming 执行为核心的 MPP 查询引擎其默认执行模型强调数据流式传输和算子流水执行。同时为支持 Exchange Materialization、CTE Materialization 等高级能力现代 Presto 引入了 Section 等更高层调度抽象用于统一描述纯 Streaming 执行和带物化依赖的混合执行模型。随着物化执行等能力的发展Presto 的调度执行模型逐渐形成 Section、Stage、Task 等多层抽象用于描述调度执行的边界、生命周期以及依赖关系默认情况下Stage 之间通过 Remote Streaming Exchange 进行流式数据交换避免中间结果完全物化。在该模型下整个 Query 的所有 Stage 通常位于同一个 Section 内Stage 之间通过 Remote Streaming Exchange 形成流水依赖而不存在跨 Section 的物化依赖 DAG。在 Exchange Materialization 或 Cte Materialization 等模式下查询会根据物化依赖边界被切分为多个 SectionSection 之间可以形成显式 DAG 调度关系。调度以 Section 作为高层组织单元以 Stage 作为生命周期管理和依赖调度实体以 Task / Split 作为执行层实际调度对象。关于 Section 这一层级概念的引入有必要再专门介绍一下。Presto 最初的分布式执行模型主要围绕 Stage 和 Streaming Exchange 构建Stage Tree 足以描述Producer-Consumer 之间的流水依赖关系。随着 Exchange Materialization 等能力的引入Presto 开始支持具有持久化边界的阶段化执行模式查询执行不再始终表现为单一 Streaming Pipeline而可能被切分为多个具有独立生命周期的执行区域。因此Section 作为位于 Stage 之上的执行组织抽象被引入用于表示一组共享执行生命周期的 Stage并描述这些独立执行区域之间基于物化数据形成的 DAG 依赖关系。三、从 PlanNode 到 SubPlanPlanFragmenter 的核心职责3.1 PlanFragmenter 做什么在 Presto 中负责完成分布式规划的核心组件为PlanFragmenter。PlanFragmenter 不仅仅是一个“切 Stage 的工具”它实际上承担的是分布式执行计划构造Distributed Plan Construction包括执行 Stage 边界划分构造 PlanFragment推导 Fragment 的数据分布属性连接上下游 Stage 的数据流关系反向推导并调整数据分布需求协调 Source 与输出 Partitioning标记 Grouped Execution 能力用一句话来描述PlanFragmenter 将优化器生成的 Physical Plan Tree 转换为面向分布式执行的 SubPlan Tree。在这个过程中它不仅根据 Remote Exchange 边界切分 Stage还负责确定每个PlanFragment 的分布式执行属性包括内部数据分布PartitioningHandle、输出数据分布PartitioningScheme以及相关执行能力。输入PlanNode Tree | PlanFragmenter v 输出SubPlan TreeSubPlan 是 Presto 分布式规划阶段生成的逻辑执行结构它表示一个 Stage 及其依赖的子 Stage 组成的树形执行计划。一个 SubPlan 主要由两部分组成SubPlan -- PlanFragment -- children(SubPlan)PlanFragment 描述一个 Stage 的计算语义以及其在分布式执行中的数据交换语义它包含fragmentId: 当前 Fragment 的唯一标识root: 当前 Stage 内部执行的 PlanNode 子树根节点partitioningHandle: 描述当前 Stage 内部数据分布方式partitioningScheme: 描述当前 Stage 输出数据的分布方式variables: 描述当前 Stage 输出的变量集合即该 Stage 对外暴露的数据 schematableScanSchedulingOrder: 描述当前 Stage 中包含的按调度顺序排列的 TableScan 节点列表remoteSourceNodes: 描述当前 Stage 中包含的远程数据源的集合execution properties: 描述该 Stage 的执行属性例如 grouped execution 等等。因此PlanFragment 不仅描述“这个 Stage 内部执行什么计算”还定义了该 Stage 在分布式执行过程中的计算逻辑、数据分布语义、输出接口以及执行能力等等信息。children 描述当前 Stage 依赖的上游 Stageschildren(SubPlan) 表示当前 Stage 所依赖的上游 Stage这些 Stage 通常对应当前 Fragment 中 RemoteSourceNode 背后的 Producer Stage。后续阶段中SqlQueryScheduler 会根据 SubPlan Tree 创建对应的 StageExecution并按照其中描述的依赖关系划分并推进各个 Section / Stage 的调度执行。3.2 PlanFragmenter 的整体流程PlanFragmenter 的整体处理过程可以抽象为createSubPlans() | v FragmentProperties 初始化 | v Visitor 遍历 PlanNode | -------------------------------- | | v v 普通 PlanNode 处理 ExchangeNode 处理 | | | ---------------------------- | | | | v v | Streaming Exchange Materialized Exchange | | | | v v v 切断 Fragment 边界 创建物化数据边界构建临时表写入计划 设置分布属性 | | | v v | 生成 RemoteSourceNode 生成 TableScanNode | | | | ---------------------------- | | | v | 生成 PlanFragment | | ------buildRootFragment--------- v SubPlan | v Grouped Execution 分析 | v Partitioning 修正 | vPlanFragmenter 是 Presto 从单节点 Physical Plan 到 Distributed Execution Plan 的关键转换层。它通过识别 Exchange 边界将连续的计算逻辑转换为由多个 Fragment 组成的 Distributed Plan并同时确定 Fragment 间的数据交换方式、Partitioning 语义以及执行属性。四、FragmentPropertiesStage 构建过程中的状态模型PlanFragmenter 在遍历 PlanNode 时需要维护当前 Stage 的各种属性。这些状态封装在FragmentProperties 中。从更加工程实现的角度来说FragmentProperties可以看作 Fragmenter Visitor 在遍历 PlanNode 过程中维护的“当前 Stage 构建上下文”。其核心字段包括FragmentProperties --ListSubPlan children; --PartitioningScheme partitioningScheme; --OptionalPartitioningHandle partitioningHandle; --SetPlanNodeId partitionedSources;各个字段的说明如下字段含义children遍历过程中推导出来的当前 Stage 依赖的子 Stage 列表partitioningHandle遍历过程中推导出来的当前 Stage 内部数据分布方式partitioningScheme遍历过程中推导出来的当前 Stage 输出数据分布方式partitionedSources遍历过程中推导出来的当前 Stage 内的数据源其中最容易混淆的是partitioningHandle 和 partitioningScheme 两个概念。4.1 PartitioningHandlePartitioningHandle 描述了当前 Stage 内部计算应该采用什么数据布局。其作用为决定Task 如何划分数据Source 如何提供数据Connector、Exchange 以及执行调度阶段都会基于该分布语义进行数据布局和调度决策。例如Hive Bucket(customer_id, 8)表示数据在源端已按 Hive 桶算法在 customer_id 上分为 8 组Stage 内部要求沿用此布局进行计算。4.2 PartitioningSchemePartitioningScheme 描述了当前 Stage 输出数据如何提供给消费者。其被用于在 Remote Exchange 中与下游 Stage 要求的数据布局保持兼容。例如HASH(order_id, 4)表示输出的数据会按 order_id 哈希到 4 个分区以匹配下游节点的读取方式。两者方向不同可以简单理解为partitioningHandle 决定“我怎么接收/处理数据”partitioningScheme 决定“我怎么输出数据”五、ExchangeStage 切分的核心机制Presto 分布式规划阶段最重要的概念为Exchange 决定 Stage 边界。需要说明的是“Exchange 决定 Stage 边界”描述的是分布式规划阶段的职责——即根据已经存在的 Remote Exchange 节点进行 Stage 切分。而在更早的 AddExchanges 优化阶段优化器会基于不同算子的语义特征决定是否需要以及在何处插入 Remote Exchange 节点。换言之AddExchanges 阶段根据算子特点决定要不要加 Exchange原因层分布式规划阶段根据已有的 Exchange 执行 Stage 切分结果层本文主要聚焦于后者。理解了这一分层就不难理解为什么会有如下结论不是Join Stage Aggregation Stage而是Remote Exchange Stage BoundaryStage 的边界由数据交换节点定义而交换节点的插入则由上游算子的计算需求驱动。5.1 Local ExchangeLocal Exchange 表示同一个 Stage 内的数据交换也即是在同一个 Worker 节点上的同一个 Task 内不同 Pipeline 之间的数据交换。例如在一个 Task 内部Driver A | Local Exchange | Driver B不会产生新的 Stage。详细机制请参见另一篇文章《Presto 查询引擎内核详解AddLocalExchanges——基于数据流属性约束的本地执行规划》。5.2 Remote Streaming ExchangeRemote Streaming Exchange 表示跨 Worker 的流式数据交换。例如原始的物理执行计划为Output | Exchange | Join / \ Scan Orders Exchange \ Scan CustomerStage 切分转换之后Stage 0 Stage 1 Stage2 Output Join Scan Customer | / \ RemoteSource Scan Orders RemoteSource在 Stage 切分过程中Exchange 被替换为RemoteSourceNode。同时创建新的Child SubPlan。最终形成如下形状的一个 Stage TreeStage0 | Stage1 | Stage2六、Materialized Exchange 与 Section前文提到的 Streaming Exchange 和 Stage 切分构成了经典的 Stage Tree 执行模型。但当引入物化边界后调度执行模型可以从树演进为有向无环图DAG这便是引入 Section 这一层级的核心价值所在。6.1 从 Streaming 到 Materialized物化依赖的引入Streaming Exchange 的语义是流水线式的Producer Stage |Network Stream不落盘 ▼ Consumer StageProducer 和 Consumer 之间通过网络直接传输数据不存在中间状态持久化。Materialized Exchange 则引入了物化依赖Producer Stage |写入临时表落盘 ▼ Temporary Table |从临时表读取 ▼ Consumer Stage此时 Producer Stage 必须完全执行完毕并物化结果后Consumer Stage 才能开始执行。这种先写后读的依赖关系在物化交换边界上引入了同步 barrier使 Producer 和 Consumer 生命周期解耦但也牺牲了一部分端到端流水执行能力。6.2 Section物化边界定义的执行切分Materialized Exchange 的出现在 Stage Tree 中切分出了SectionStage Tree Section DAG ┌─────────────────┐ ┌─────────────────┐ │ Stage 0, 1 │ │ Section A │ │ │ │ │ (Stages 0-1) │ │ ▼ │ └────────┬────────┘ │ Materialized │ │ 物化依赖 │ Exchange │ ▼ │ │ │ ┌─────────────────┐ │ ▼ │ │ Section B │ │ Stage 2, 3 │ │ (Stages 2-3) │ └─────────────────┘ └─────────────────┘Section 的定义Section 是由 Materialized Exchange 边界划分出的调度执行组织单元。同一个 Section 内的 Stage 通过 Streaming Exchange 连接形成流水线不同 Section 之间通过 Materialized Exchange 连接形成物化依赖。相比早期 Presto 的单纯流式 Stage Tree 模型Section DAG 带来了几个关键能力容错恢复Materialized Exchange 物化的中间结果可作为检查点某个 Section 失败后无需重跑整个查询只需重跑该 Section 及其下游此外结合分组执行grouped execution可以支持分组级任务级的失败重试。运行时自适应优化下游 Section 的计划可以基于上游 Section 物化后的真实数据统计信息如行数、数据分布、NDV、Null 比例等进行动态调整引入类似于 Spark AQE 的能力。资源解耦不同 Section 可独立调度无需同时持有所有 Stage 的资源对大规模 ETL 场景尤其友好。这一演进扩展了 Presto 的执行模型使其在保持交互式查询能力的同时具备更强的长链路查询、大规模数据处理以及复杂工作负载支持能力。七、Fragment 创建与 SubPlan 组装在 PlanFragmenter 遍历 PlanNode 的过程中随着发现不同类型的 Exchange 边界会逐步完成 Fragment 切分和创建。因此PlanFragment 并不是在整个 Plan 遍历完成后统一生成而是在 Fragment 边界识别过程中逐步构建Visitor 遍历 PlanNode | v 发现 Exchange 边界 | v 递归构建 Exchange Source 对应的 Child Fragment | v 当前 Fragment 替换 Exchange 为 RemoteSourceNode 或建立物化读取节点 | v 继续处理父节点 | v 构建 Root Fragment | v 组装 SubPlan Tree在 Fragment / SubPlan 构建阶段主要完成以下工作1. 确定 Fragment 的执行结构PlanFragmenter 根据当前 Stage 内部的 PlanNode 拓扑关系确定当前 Fragment 的根节点Source 节点关系数据交换关系。确保该 Fragment 描述的是一个完整、可执行的计算区域。2. 汇总执行属性在前面的遍历过程中FragmentProperties 已经逐步收集了当前 Stage 的执行信息包括数据分布方式PartitioningHandle输出分区方案PartitioningScheme直属子 Fragment 对应的 SubPlan 依赖关系children相关执行能力属性。这些信息最终都会被封装到 PlanFragment / SubPlan 中。3. SubPlan Tree 组装当一个 Fragment 构建完成后PlanFragmenter 会根据 FragmentProperties 中维护的子 Fragment 关系构造对应的 SubPlanParent SubPlan | ---------------- | | Child SubPlan Child SubPlan | ......其中Parent Fragment 通过 RemoteSource 或物化读取节点消费 Child Fragment 输出Child Fragment 是 Producer StageSubPlan Tree 描述 Fragment 之间的数据依赖关系。至此PlanFragmenter 完成了 Distributed Plan 的初步构建。生成的 SubPlan Tree 作为后续分布式执行优化和调度规划的基础接下来的阶段还将进一步分析 Fragment 的执行能力并调整其数据分布策略。八、Grouped Execution 能力分析除了构建 Stage 拓扑和确定数据分布之外PlanFragmenter 还会对部分执行能力进行分析其中包括 Grouped Execution。Grouped Execution 主要用于降低大规模 Join 或 Aggregation 场景下的内存压力。传统执行模式下一个 Task 通常需要一次性加载完整的数据分组Task | v 读取全部 Bucket | v Build Hash Table | v 执行 Join当 Build Side 数据规模较大时会导致较高的峰值内存消耗。Grouped Execution 则将执行过程拆分为多个独立的数据组通常对应 Bucket 或 LifespanBucket n | v 执行 Join | v 释放资源通过限制单次参与计算的数据规模降低 Hash Build 等算子的内存压力。在分布式规划阶段PlanFragmenter 并不负责执行 Grouped Execution而是通过 GroupedExecutionTagger 分析当前 Fragment 是否具备该能力并将分析结果记录到StageExecutionDescriptor 中以供后续 Stage 调度和 Task 执行阶段使用。其核心关注点包括当前 PlanNode 是否支持 Grouped Execution子树是否适合采用 Grouped Execution哪些 TableScan 节点可以提供分组数据分组数量等执行信息。关于 Grouped Execution 的完整机制我们将在后续专门的文章中单独展开。九、计算需求反向影响数据布局Partitioning Reassignment在分布式规划过程中Presto 不仅需要决定计算如何分布执行有时还需要根据计算阶段的需求反向调整数据源提供的数据布局。这一机制称为Partitioning Reassignment它体现了 Presto 分布式规划中的一个重要思想执行计划中的数据分布需求可以反向影响 Source 层的数据读取方式。9.1 为什么需要 Partitioning Reassignment在默认情况下TableScan 通常只提供数据源自身的数据分布方式。例如Hive Connector 的某一个表上已经存在的数据分布格式为bucket(customer_id, 16)但是在某些执行场景中Stage 本身已经确定需要特定的数据布局。例如Stage Requirement Hive Bucket Partitioning bucket(customer_id, 8)这意味着数据需要以 Hive 的 Bucket Hash 算法按照 customer_id 的值被划分到 8 个 bucket 中而不是该表原始划分的 16 个 bucket 中。如果 Source 提供的数据布局无法满足该要求则可能导致额外的数据 Shuffle增加 Network Exchange 开销降低 Join 等算子的执行效率。因此Presto 会尝试让数据源提供与 Stage 计算需求更加匹配的数据布局。9.2 计算需求如何影响 SourcePartitioning Reassignment 的核心过程可以抽象为Stage Execution Requirement | v PartitioningHandle Reassignment | v Table Layout 调整 | v Connector 提供匹配的数据分布也就是说分布式规划阶段首先遍历整个计划树切分 Stage 并自底向上的确定各个 Stage 希望数据以什么方式组织。然后再自顶向下的将数据输出需求传递给 Source 层使 Connector 最终调整并提供满足条件的数据布局。例如Stage: HASH(customer_id, 8) ↓ TableScan: Hive Bucket(customer_id, 8)需要说明的是在切分阶段甚至是更早的 AddExchanges 优化阶段已经确保了一个 Stage 与其内部持有的 Sources 之间的数据布局的兼容性。在优化阶段优化器已经通过ConnectorMetadata.getCommonPartitioningHandle等接口询问并确认过 Source Connector 可以调整并提供兼容的数据布局。否则会通过添加 Remote Exchange节点来调整数据布局因此也会将其切分到不同的 Stage 中。因此在 Partitioning Reassignment 阶段不会出现 Stage 要求的数据布局无法被 Source Connector满足而导致的规划失败。9.3 体现的数据计算协同思想传统查询优化通常是Data Source | v Execution Plan即数据源提供数据执行计划适配数据。而 Partitioning Reassignment 体现的是另一种方向Execution Requirement | v Data Layout Selection | v Source Optimization这是现代湖仓查询引擎中的重要设计趋势查询执行不再只是被动消费数据执行计划可以影响数据读取方式存储布局与计算模型形成协同优化。在 Presto 中这种机制连接了Optimizer ↓ Distributed Planning ↓ Connector Layer ↓ Physical Data Layout使查询引擎能够根据执行目标选择更加高效的数据访问路径。这体现了一种“计算分布需求反向约束数据源布局”的设计思想。十、从 SubPlan 到 Distributed Scheduler分布式规划阶段生成的 SubPlan 是一种具有 Stage 拓扑、数据分布策略和执行能力描述的 Distributed Plan其为后续 Query Scheduler 创建 Task、分配 Split 和驱动 Worker 执行奠定基础。Optimizer | v PlanNode Tree | v PlanFragmenter | v SubPlan Tree | v SqlQueryScheduler | v StageExecution关于 Query Scheduler 分布式调度部分将在后续的文章中专门讲解说明此处不再展开。总结Presto 分布式规划的核心思想Presto 的分布式规划过程本质上是将优化后的 Physical Plan 转换为包含执行边界、数据交换语义和调度拓扑的 Distributed Execution Model。这一过程围绕三个核心抽象展开1. 通过 Exchange 定义计算边界Presto 并不是简单按照算子类型切分执行阶段而是通过识别 Exchange 边界将连续的计算逻辑划分为多个独立执行区域。Remote Streaming Exchange 定义流式数据交换边界Materialized Exchange 定义物化数据交换边界。基于这些边界PlanFragmenter 将 Physical Plan 切分为多个相互关联的 PlanFragment。2. 通过 Partitioning 描述分布式数据语义在 Fragment 切分之后Presto 需要进一步确定数据如何在 Worker 之间分布Fragment 之间如何交换数据下游计算如何消费上游结果。其中PartitioningHandle 描述 Fragment 数据分布策略PartitioningScheme 描述 Fragment 输出数据的分区方案以及数据交换布局。它们共同定义了分布式执行中的数据流转语义。3. 通过分层执行模型连接规划与调度基于 Fragment 依赖关系Presto 构建 SubPlan并进一步将其映射为运行时执行模型Streaming Exchange 保留 Pipeline Streaming 的低延迟执行能力Materialized Exchange 和 Section 引入持久化执行边界使查询能够形成更复杂的执行拓扑Scheduler 根据 SubPlan 管理 Stage 生命周期并将计算任务调度到 Worker 执行。最终Presto 分布式规划的核心目标是在保留 Pipeline Streaming 低延迟执行特性的同时将优化后的计算计划转换为具有明确执行边界、数据交换语义和调度拓扑的分布式执行模型从而连接查询优化与集群调度执行。作者王冬PrestoDB Committer | Presto Iceberg Code OwnerGitHub: https://github.com/hantangwangdEmail: mingwbdgmail.com本文章同步发表于https://hantangwangd.github.io/zh/posts/2026-08-10-distributed-planning-mechanism.html

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价