Dong Wang

Presto 查询引擎内核详解:Worker 本地执行模型——从物理计划到可调度执行单元

本文档深入探讨 Presto Worker 节点内部的核心执行机制。我们将追踪一个 Stage 的 PlanFragment 如何从物理计划片段一步步转化为可并发执行的本地任务,并最终被调度器高效执行。

涉及到的核心概念包括 Driver、SplitRunner、TaskExecutor 以及多级优先级队列。

1. 核心概念:从物理计划片段到本地执行单元

1.1 Driver:执行片段与数据分组的二维切分

当一个Stage的PlanFragment被发送到Worker节点并规划成本地可执行计划后,它会被切分成多个Pipeline。每个Pipeline是组可以流水线化执行的物理算子(Operator)。 随后,每个Pipeline会根据执行策略被实例化为多个 Driver 实例。

Driver 是 Presto 中的核心执行概念,它被通过两个维度进行定义:

Pipeline 决定执行逻辑,Lifespan 决定数据隔离边界,而 Driver 是二者组合后产生的实际执行实例。

关键规则:切分后,每个 Driver 最多只有一个 SourceOperator 类型的输入节点(如 TableScanOperator 或 ExchangeOperator)。

1.2 DriverSplitRunner & PrioritizedSplitRunner:可调度的最小执行单元

Driver本身不会直接被调度执行,它会被封装到DriverSplitRunner中:

2. 核心调度器:TaskExecutor

TaskExecutor 是每个 Worker 节点内全局唯一的调度器,负责管理和调度本节点上所有 SQL 任务(SqlTask)产生的所有 PrioritizedSplitRunners。

2.1 核心组件(属性)

2.2 调度模型:时间片轮转

PrioritizedSplitRunner的执行采用类似操作系统的分时调度策略。

2.3 多级优先级队列 (MultiLevelSplitQueue)

为了实现公平调度,TaskExecutor使用了一个多级队列来管理等待中的Split。

note: 多级优先级队列涉及到的算法较为精细与复杂,后面会有专门的一篇文章详解

3. 执行流程:从创建到运行

3.1 创建阶段

SqlTaskManager 启动:Worker 节点的 Presto Server 进程启动时,通过依赖注入框架(Guice)初始化全局唯一的 SqlTaskManager,它内部持有 TaskExecutor。

接收 Task 请求:当 Worker 接收到一个 Task 请求(TaskRequest)时,SqlTaskManager 会创建或获取对应的 SqlTask 和 QueryContext。

构建 SqlTaskExecution:SqlTask 的核心是 SqlTaskExecution,它负责管理一个 Task 的生命周期。在构造时,它会:

创建 Driver 实例:SqlTaskExecution 会根据数据来源(如 Split)调用 DriverSplitRunnerFactory.createDriverRunner() 来创建实际的 DriverSplitRunner。

3.2 调度与执行阶段

入队 (enqueueSplits):SqlTaskExecution 调用 TaskExecutor.enqueueSplits(),将新创建的 PrioritizedSplitRunner 放入 TaskHandle 的队列中,并最终进入 MultiLevelSplitQueue。

工作线程 (TaskRunner) 轮询:一个空闲的 TaskRunner 线程会调用 MultiLevelSplitQueue.take(),根据多级优先级算法获取下一个要执行的 PrioritizedSplitRunner。

执行时间片 (PrioritizedSplitRunner.process()):

Driver执行 (Driver.processFor(duration)):

核心数据处理循环 (Driver.processInternal()):

调度器响应:PrioritizedSplitRunner.process() 会根据Driver.processFor()的返回值决定下一步:

说明:Presto Worker 内部并不是简单的“一个 Task 一个线程”,而是通过 TaskExecutor 将大量 Pipeline + Lifespan 对应的执行实例抽象为支持协作式抢占的执行单元,在 Worker 级别完成类似操作系统 CPU 调度的资源复用。

4. 并发控制与优先级管理

4.1 TaskHandle 与并发控制器 (SplitConcurrencyController)

每个 TaskHandle 管理其所属 Task 的所有 Split,并通过 SplitConcurrencyController 动态调整该Task的目标并发度。

调整依据:

Leaf Split vs. Intermediate Split:

4.2 优先级追踪 (TaskPriorityTracker)

TaskPriorityTracker与MultiLevelSplitQueue协同工作。

5. 执行路径优化:Fragment Result Cache

除了正常Driver执行路径之外,对于特定模式的查询(如聚合的 Partial 阶段),Presto 还支持将中间结果缓存起来,从而进一步加速后续相同查询。

适用场景:Partial AggregationNode,且其子孙节点是 TableScan、Filter、Project 等。

工作流程:

6. 总结

Presto Worker端的执行模型是一个精妙设计的分层异步系统:

其核心思想是将静态的物理计划片段转换为动态执行流水线:通过 Pipeline 和 Driver 将执行片段拆解为大量可并行推进的细粒度执行实体,再由 TaskExecutor 在有限计算资源上进行统一调度,从而同时获得节点内流水线并行能力、高吞吐执行效率以及多查询资源公平性。


欢迎交流

本文基于作者当前的理解与实践经验整理而成,难免存在疏漏或值得进一步探讨之处。

如果您对文中的观点有不同看法,发现任何问题,或有相关实践经验,欢迎通过 GitHub Issue 与作者交流讨论。

期待与更多同行围绕数据基础设施相关技术展开交流,分享实践经验,共同学习、共同进步。