[译] Velox Primer 1、2、3
本文合译 Velox 官方博客的 A Velox Primer 三篇文章,作者为 Orri Erling 和 Pedro Pedreira。三篇文章围绕同一个分组聚合查询,依次介绍分布式执行、扫描与分区输出,以及 shuffle 消费和本地聚合。
译文保留原文的章节顺序、示例和 2025 年的论述背景。配图按原文重新绘制为 SVG,另有两张标注为“译者补图”的说明图。少量 API 拼写和概念边界说明集中放在文末译注中。
| 原文 | 发布日期 | 主题 |
|---|---|---|
| A Velox Primer, Part 1 | 2025-02-17 | 分布式查询执行基础 |
| A Velox Primer, Part 2 | 2025-03-25 | 从计划到 Driver 和 Operator |
| A Velox Primer, Part 3 | 2025-05-12 | Shuffle 消费、字典与聚合屏障 |
第一篇:分布式查询执行基础
原文 Part 1 · 2025-02-17 · Orri Erling、Pedro Pedreira
这是一个介绍 Velox 内部结构与基本概念的短文系列的第一篇。本篇将讨论分布式查询如何执行、数据如何在不同阶段之间进行 shuffle,以及支持这些功能的 Velox 概念:Task、Split、Pipeline、Driver 和 Operator。
分布式查询执行
Velox 是一个库,为分布式计算引擎中的查询片段(query fragment)提供执行功能。Presto、Spark 这样的分布式计算引擎运行的是所谓的 exchange parallel plan,即通过 exchange 实现并行的执行计划。Exchange 也称为数据 shuffle,用于让数据从一个阶段流向下一个阶段。查询片段由 shuffle 连接,包含在单个 worker 节点内执行的处理逻辑。Shuffle 从一组片段接收输入,并根据数据的某个特征——也就是分区键——将输入行路由到特定的消费者。消费者从 shuffle 读取数据;无论数据来自哪个生产者,只要这些行的分区键与该消费者负责的分区匹配,就会由它接收。
以下面的查询为例:
SELECT key, count(*) FROM T GROUP BY key
假设有 n 个叶子片段分别扫描表 T 的不同部分,也就是图中阶段 1 的绿色圆圈。假设叶子片段的末端有一个 shuffle,按 key 对数据行进行分区。接下来有 m 个消费者片段,即阶段 2 的黄色圆圈;每个消费者根据 key 列接收互不重叠的一部分数据行。随后,每个消费者构建一张以 key 为键的哈希表,记录每个 key 值已经出现了多少次。
如果 key 有 1000 亿个不同的值,那么这张哈希表就很难方便地放进一台机器。为了提高效率,可以将其划分为 100 张各有 10 亿个条目的哈希表。这就是借助 exchange 实现横向扩展的意义。可以把 shuffle 理解为这样一种机制:面对多个 worker 共同产生的一条大数据流,每个消费者只消费其中属于自己的、与其他消费者不同的一片数据。
Presto 和 Spark 这样的分布式查询引擎都符合上述描述。它们的区别之一是 shuffle 的实现方式,后文会再讨论这一点。
Task 与 Split
在一个 worker 节点内部,查询片段在 Velox 中的表示称为 Task(velox::exec::Task)。决定 Task 行为的主要有两类信息:
- Plan——
velox::core::PlanNode:规定 Task 要做什么。 - Split——
velox::exec::Split:规定 Task 要处理什么数据。
Split 与 Task 正在执行的计划相对应。对于第一阶段的 Task,也就是表扫描,Split 指定要扫描的文件片段。对于第二阶段的 Task,也就是 group by,Split 则标识它要从哪些表扫描 Task 读取输入。Split 分为文件 split(velox::connector::ConnectorSplit)和远程 split(velox::exec::RemoteConnectorSplit):前者标识要读取的数据,后者标识一个正在运行的 Task。
分布式引擎创建 PlanNode 和 Split。Velox 接收这些信息并创建 Task。Task 再将统计信息、错误以及其他状态信息反馈给分布式引擎。
Pipeline、Driver 与 Operator
Task 内部包含 Pipeline。每条 pipeline 都是一串线性排列的算子(velox::exec::Operator),算子是实现关系运算逻辑的对象。在上面的 group by 示例中,第一个 Task 有一条 pipeline,包含 TableScan(velox::exec::TableScan)和 PartitionedOutput(velox::exec::PartitionedOutput)。第二个 Task 也有一条 pipeline,其中依次有 Exchange、LocalExchange、HashAggregation 和 PartitionedOutput。
每条 pipeline 有一个或多个 Driver。一个 Driver 是一串线性排列的 Operator 的容器,通常运行在自己的线程上。Pipeline 则是具有相同算子序列的 Driver 的集合。每个 Operator 实例属于相应的 Driver,Driver 又属于 Task。它们通过智能指针相互关联,从而在需要的时间内保持存活。
Hash join 是一个 Task 包含两条 pipeline 的例子:build 侧和 probe 侧分别对应一条 pipeline。这种划分很自然,因为 build 必须完成之后,probe 才能继续执行。后文会进一步讨论它。
同一个 Driver 中的 Operator 通过 Driver 相互传递数据。Driver 从一个 Operator 取出输出,维护统计信息,再把数据交给下一个 Operator。Operator 之间传递的数据由向量组成。具体而言,一个 Operator 产生或消费一个 RowVector;RowVector 为关系中的每一列包含一个子向量,与 Arrow 中的 RecordBatch 对应。所有向量都是 velox::BaseVector 的子类。
算子的 Source、Sink 与状态
Driver 中的第一个 Operator 称为 source。Source 算子根据 Split 确定提供输入数据的文件或远程 Task,并产生一系列 RowVector。Pipeline 中的最后一个算子称为 sink。Sink 算子不产生输出 RowVector,而是把数据放到某个地方,供消费者取走。典型的 sink 是 PartitionedOutput,它的消费者是另一个 Task 中位于 source 位置的 Exchange。既不是 source 也不是 sink 的算子包括 FilterProject、HashProbe、HashAggregation 等,后文会继续介绍。
Operator 也具有状态:它可能被阻塞,可能可以接收输入,可能有或没有输出可供产生,也可能收到不再有输入的通知。Operator 不会直接调用其他 Operator 的 API,而是由 Driver 决定接下来推进哪个 Operator。这样做的好处是,Driver 的任何状态都不需要保存在嵌套的函数调用中。Driver 具有平坦的调用栈,因此可以随时获得执行线程或离开执行线程,而不需要展开并恢复嵌套的函数栈帧。
回顾
我们已经看到,分布式查询由 shuffle 连接的多个片段组成。每个片段对应一个 Task,Task 中又包含 pipeline、Driver 和 Operator。PlanNode 表示 Task 应该做什么,并告诉 Task 应该如何配置 Driver 和 Operator。Split 告诉 Task 去哪里寻找输入,例如文件或其他 Task。一个 Driver 对应一条执行线程:它可能正在运行,也可能被阻塞,例如等待数据可用,或者等待消费者消费已经产生的数据。Operator 之间通过向量传递数据。
下一篇文章将讨论一个查询生命周期中的不同阶段。
第二篇:从计划到 Driver 和 Operator
原文 Part 2 · 2025-03-25 · Orri Erling、Pedro Pedreira
本文将讨论分布式计算引擎如何执行一个与第一篇文章中类似的查询:
SELECT l_partkey, count(*) FROM lineitem GROUP BY l_partkey;
我们使用 TPC-H 的 schema 来说明这个例子,并使用 Prestissimo 作为编排分布式查询执行的计算引擎。Prestissimo 负责查询引擎的前端工作,包括解析、元数据解析、计划生成与优化,以及分布式执行,包括分配资源和发送查询片段;Velox 则负责在单个 worker 节点内部执行计划片段。本文会逐步说明哪些功能由 Velox 完成,哪些功能由分布式引擎完成——在这个例子中,后者就是 Prestissimo。
查询的建立
Prestissimo 首先通过 coordinator 节点接收查询,由 coordinator 负责解析查询并生成计划。对于我们的示例查询,会创建一个包含三个查询片段的分布式查询计划:
- 第一个片段读取 lineitem 表的 l_partkey 列,并根据 l_partkey 的哈希值划分输出。
- 第二个片段读取第一个片段的输出,并更新一张以 l_partkey 为键的哈希表,记录每个 l_partkey 值出现的次数;这就是
count(*)聚合函数的实现。 - 最后一个片段在第二个片段收齐第一个片段的全部数据行之后,读取这些哈希表的内容。
前两个片段之间的 shuffle 按 l_partkey 分区。假设第二个片段有 100 个实例:如果 l_partkey 的哈希值对 100 取模为 0,这一行就发送给第二阶段的第一个 Task;如果余数为 1,就发送给第二个 Task,依此类推。这样,第二阶段的每个 Task 都会得到互不重叠的一部分数据行。第二阶段和第三阶段之间的 shuffle 是 gather,也就是说,第三阶段只有一个 Task,它读取第二阶段全部 100 个 Task 的输出。
Stage 是共享同一个计划片段的一组 Task。Task 是 Prestissimo 与 Velox 之间的主要集成点:它是 Velox 中真正执行处理的实例,负责处理流经该 stage 的全部或部分数据。
为了建立分布式执行,Prestissimo 首先从自己管理的 Prestissimo 服务进程池中选择 worker。假设第一阶段的宽度为 10 个 worker,它就选择 10 个服务进程,并向它们发送第一阶段的计划。随后,为第二阶段选择 100 个 worker,并发送第二阶段的计划。汇集结果的最后一个阶段只有一个 worker,因此最终计划只发送给一个 worker。不同阶段的 worker 集合可以重叠,所以一个 worker 进程可能同时承载同一个查询中的多个 stage。
接下来,我们更仔细地看看每个 worker 在查询建立时做了什么。
Task 的建立
在 Prestissimo 中,用于在 worker 中建立 Task 的消息称为 Task Update。一条 Task Update 包含计划、配置项,以及一个可选的 split 列表。Split 还会附带进一步的信息:它发给哪个 plan node,以及这个 plan node 和 split group 后续是否还会收到更多 split。
生成 split 需要枚举存储系统中的文件,因此可能需要一些时间。Presto 允许异步向 worker 发送 split,使 split 的生成能够与首批 split 的执行同时进行。因此,第一条 Task Update 指定计划和配置,后续的消息只继续添加 split。
除了计划,coordinator 还以字符串键到字符串值的映射提供配置,包括顶层配置和 connector 级配置。Connector 配置包含各个 connector 的设置;TableScan 和 TableWriter 通过 connector 处理存储系统及文件格式。这些配置,以及线程池、顶层内存池等其他信息,通过 QueryCtx 对象传给 Task。细节可参阅 Velox 仓库中的 velox/core/QueryCtx.h。
从计划到 Driver 和 Operator
Velox Task 创建之后,TaskManager 会向它提供需要处理的 split。这通过 Task::addSplit() 完成,也可以在 Task 开始执行之后继续调用。细节可参阅 velox/exec/Task.h。
现在放大来看 Task 创建时发生的事情:描述 Task 工作内容的 PlanNode 树,作为 PlanFragment 的一部分交给 Task。Task 创建过程中最重要的一步,是把计划树拆分成 pipeline。每条 pipeline 随后获得一个 DriverFactory,即为该 pipeline 创建 Driver 的工厂类。Driver 又包含实际执行查询工作的 Operator。DriverFactory 在 LocalPlanner.cpp 中创建,细节可参阅 LocalPlanner::plan。
按照称为 Volcano 的执行模型,计划表示为一棵算子树,其中每个节点消费子算子的输出,并向父算子返回输出。根节点通常是 PartitionedOutputNode 或 TableWriteNode。叶节点则是 TableScanNode、ExchangeNode,或者用于查询字面量的 ValuesNode。Velox 的完整 PlanNode 集合可见 velox/core/PlanNode.h。
PlanNode 大多与 Operator 对应。PlanNode 本身不可执行,它只是一种结构,用于描述应该怎样创建真正执行工作的 Driver 和 Operator。如果节点树只有一条分支,计划就只有一条 pipeline。如果存在包含多个子节点,也就是多个输入的节点,那么这个节点的第二个输入就会成为一条独立的 pipeline。
Task::start() 创建 DriverFactory,DriverFactory 再创建 Driver。为了开始执行,Driver 被加入线程池 executor 的队列。运行 Operator 的主要函数是 Driver::runInternal();可以查看这个函数了解 Operator 与 Driver 如何协作。Operator::isBlocked() 用于判断 Driver 能否继续推进;如果不能,Driver 就离开执行线程,直到相应的 future 就绪,再被放回 executor。
getOutput() 从一个 Operator 获取数据,addInput() 则把数据交给另一个 Operator。执行顺序是先推进最靠后的、能够产生输出的 Operator,再把输出传给下一个 Operator。如果一个 Operator 无法产生输出,就调用它前一个 Operator 的 getOutput(),不断向前寻找能够产生数据的 Operator。如果没有算子被阻塞,也没有算子能够继续产生输出,那么计划就到达末尾。此时会对各个 Operator 调用 noMoreInput()。这个调用可能使算子开始产生结果,例如,OrderBy 只有在知道自己已经拿到全部输入之后,才可以产生输出。
最小的 Pipeline:表扫描与重新分区
表扫描。 在示例查询的表扫描阶段,我们有一条包含两个算子的 pipeline:TableScan 和 PartitionedOutput。假设这条 pipeline 有五个 Driver,并且这五个 Driver 都进入线程池 executor。PartitionedOutput 由于没有输入,暂时什么也做不了,于是会调用 TableScan::getOutput(),相关实现见 velox/exec/TableScan.cpp。TableScan 的第一个动作是通过 Task::getSplitOrFuture() 查找 Split。如果没有可用的 split,该调用就返回一个 future。Driver 随后离开执行线程,并在 future 上安装回调,在 split 可用时重新调度这个 Driver。
也可能既没有 split,Task 又已经收到通知,知道后续不会再有 split。在这种情况下,TableScan 就到达结束状态。最后,如果有可用的 Split,TableScan 就解释它。根据计划中提供的 TableHandle 描述,包括列的列表与过滤条件,由 Split 指定的 Connector 创建 DataSource。DataSource 负责处理 I/O、文件格式以及表格式的细节。
接着把 split 交给 DataSource。此后,可以反复调用 DataSource::next(),从这个 Split 指定的文件或文件片段中获得输出向量,也就是批次。如果 DataSource 已到达末尾,TableScan 就去寻找下一个 split。Connector 和 DataSource 的接口见 Connector.h。
重新分区。 到这里,我们已经追踪到 TableScan 返回第一批输出。Driver 把这批数据交给 PartitionedOutput::addInput(),相关实现见 PartitionedOutput.cpp。PartitionedOutput 首先计算分区键的哈希值;在本例中,分区键是 l_partkey。这会为批次 RowVectorPtr input_ 中的每一行产生一个目的地编号。
每个目的地,也就是每个目标 worker,都有一个尚未填满的序列化缓冲区 VectorStreamGroup。如果第二阶段有 100 个 Task,每个 PartitionedOutput 就有 100 个目的地,每个目的地都有一个 VectorStreamGroup。VectorStreamGroup 的主要函数是 append(),它接收一个 RowVectorPtr 和其中一组行号,将这些行号标识的每个值序列化,追加到正在构建的序列化内容中。当 VectorStreamGroup 中积累了足够多的数据行之后,它就产生一个 SerializedPage。相关过程见 PartitionedOutput.cpp 中的 flush()。
SerializedPage 是一个自包含的序列化数据包,可以经由网络传输到下一个阶段。每个 page 只包含发给同一个接收方的数据行。这些 page 随后进入 worker 进程的 OutputBufferManager 队列。可以留意 flush() 中与 BlockingReason 有关的代码。Buffer manager 为所有 Task 的所有消费者维护各自的队列。如果队列已满,继续添加输出就可能被阻塞;这时返回一个 future,在队列有空间接收更多数据时,该 future 变为就绪。这取决于该 Task 的消费者 Task 何时取走数据。
Prestissimo 的 shuffle 由生产者一端的 PartitionedOutput 和消费者一端的 Exchange 实现。OutputBufferManager 保存已经准备好的序列化数据,等待消费者取走。它们与 Presto 线上传输协议的衔接,生产者一侧位于 TaskManager.cpp,消费者一侧位于 PrestoExchangeSource.cpp。
回顾
我们介绍了计划如何变成可执行对象,以及数据如何在 Operator 之间和内部流动。我们讨论了 Driver 如何通过阻塞、离开执行线程来等待 Split 可用,或者等待自己的输出被消费。到这里,我们才刚刚触及分布式查询叶子阶段执行过程的表面,Operator 和向量还有很多内容值得讨论。下一篇 Velox Primer 将继续研究这个最小示例查询的第二阶段。
第三篇:Shuffle 消费、字典与聚合屏障
原文 Part 3 · 2025-05-12 · Orri Erling、Pedro Pedreira
在上一篇文章结束时,我们已经走过了第一个分布式查询执行过程的一半:
SELECT l_partkey, count(*) FROM lineitem GROUP BY l_partkey;
我们讨论了查询如何启动、Task 如何建立,以及计划、Operator 与 Driver 之间如何交互。我们还介绍了查询第一阶段的执行过程:从表扫描到分区输出,也就是 shuffle 的生产者一侧。
本文将讨论查询的第二阶段,也就是 shuffle 的消费者一侧。
Shuffle 的消费者
如上一篇文章所述,在查询的第一阶段,每个 worker 读取表,然后产生一系列发给第二阶段不同 worker 的数据包,即 SerializedPage。在本例中,lineitem 表没有特定的物理分区键或聚簇键。这意味着,组成这张表的任何文件中的任意一行,都可能具有任意的 l_partkey 值。因此,为了按 l_partkey 对数据分组,我们首先要保证:l_partkey 值相同的数据行由同一个 worker 处理。这正是第一阶段末尾进行数据 shuffle 的目的。
查询的整体结构如下:
查询 coordinator 不按特定顺序向第一阶段的 worker 分发表扫描 split。Worker 处理这些 split,并在处理过程中填充目的地缓冲区,供第二阶段的 worker 消费。假设第二阶段有 100 个 worker,那么第一阶段的每个 Driver 都有自己的 PartitionedOutput,其中包含 100 个目的地。缓冲的序列化数据达到足够规模之后,就会交给 worker 的 OutputBufferManager。
现在把注意力转向第二阶段的查询片段。每个第二阶段 worker 的计划片段自上而下是:PartitionedOutput、Aggregation、LocalExchange、Exchange。
每个第二阶段 Task 对应第一阶段各个 worker 的 OutputBufferManager 中的一个目的地。第二阶段的第一个 Task 从第一阶段所有 Task 读取目的地 0 的数据;第二个 Task 读取第一阶段各 Task 的目的地 1,依此类推。所有生产者与所有消费者之间都有联系,shuffle 过程本身不需要一个中心节点进行协调。
接下来,放大看看第二阶段启动时具体发生了什么。
计划在 Exchange 之后有一个 LocalExchange 节点。这会形成两条 pipeline:一条包含 Exchange 和 LocalPartition;另一条包含 LocalExchange、HashAggregation 和 PartitionedOutput。
Velox Task 按多线程执行来设计,通常有 5 到 30 个 Driver。一个 stage 可能有数百个 Task,因此一个 stage 中可能有数千条线程。所以,第二阶段的 100 个 worker 各自消费第一阶段总输出的 1/100,但每个 worker 内部又以多线程方式运行,由多条线程从 ExchangeSource 消费数据。后面会进一步解释这一点。
为了执行多线程 group by,可以使用一张线程安全的哈希表,也可以把数据划分为 n 条互不重叠的数据流,再分别在各自线程上完成聚合。在 CPU 上,我们几乎总是希望每个线程操作自己的内存,因此会利用 local exchange,按 l_partkey 在本地对数据分区。CPU 具有复杂的缓存一致性协议,用于让观察者看到一致且有序的内存视图;当多个核心写入同一条 cache line 后,为此协调状态既不可避免,代价又高。严格的内存与线程亲和关系,是实现多核可扩展性的关键。
LocalPartition 与 LocalExchange
为了形成高效且相互独立的内存访问模式,第二阶段使用 local exchange 再次对数据进行 shuffle。在概念上,这与 Task 之间的 remote exchange 类似,但作用范围限制在 Task 内部。生产者一侧的 LocalPartition 对分区列 l_partkey 计算哈希,将数据划分到不同目的地,每个目的地对应消费者 pipeline 的一个 Driver。消费者 pipeline 的 source 算子是 LocalExchange,它从 LocalPartition 填充的队列中读取数据。细节可参阅 LocalPartition.h;有关如何在 local exchange 的生产者与消费者之间建立队列,也可参阅 Task.cpp。
Remote shuffle 操作的是序列化数据,local exchange 则在线程之间传递内存中的向量。这也是我们第一次遇到利用列式编码加速向量化执行的概念。Velox 因大量使用这类技术而为人所知,我们把它称为 compressed execution(压缩执行)。在这里,我们使用 dictionary 将向量切分给多个目的地;下面就来讨论它。
Dictionary 的妙用
查询执行经常需要改变数据的基数(cardinality),也就是结果中的行数。过滤和连接本质上就在做这件事:过滤器或选择性较高的连接会减少行数,而同一个键匹配多个结果的连接则可能增加行数。
LocalPartition 中的重新分区会根据分区键,为每个输入行分配一个目的地。随后,它为每个目的地构造一个向量,其中只包含发往该目的地的行。假设当前输入中的第 2、8、100 行都被哈希到目的地 1,那么目的地 1 的向量就只需要原始输入中的第 2、8、100 行。我们可以复制数据来构造一个三行向量,但这里采用另一种方式来省掉复制:为原始输入中的每一列套上一层 DictionaryVector,长度为 3,索引为 2、8、100。与复制相比,这样效率高得多,尤其是在处理宽数据和嵌套数据时。
之后,负责目的地 1 的 LocalExchange 消费者线程会看到一个三行的 DictionaryVector。当 HashAggregation 算子访问它时,聚合逻辑识别出这是一个 dictionary,于是通过间接寻址访问实际数据:输出的第 0 行访问底层 base 向量的第 2 个值,第 1 行访问第 8 个值,依此类推。目的地 0 的消费者也进行同样的处理,只是它访问的可能是第 4、9、50 行。
所有消费者线程的 dictionary 都指向同一个 base。每个消费者线程只查看其中不同的子集。各个核心会读取相同的 cache line,但由于 base 不会被修改,也就是只读的,因此不会产生缓存一致性方面的开销。
概括来说,**DictionaryVector<T>** 是任意类型为 T 的向量的外层包装。DictionaryVector 保存索引,这些索引指向 base 向量中的位置。Dictionary 编码通常用于一列中不同值较少的情况。以字符串 “experiment” 和 “baseline” 为例,如果一列只包含这两个值,那么更高效的表示方式是:在一个向量的位置 0 保存 “experiment”,位置 1 保存 “baseline”,再用一个具有例如 1000 个元素的 DictionaryVector 表示这一列,其中每个元素都是索引 0 或 1。
除此之外,DictionaryVector 也可以表示 base 向量元素的一个子集,或者它们的一种重排。由于所有接受向量的地方也接受 DictionaryVector,DictionaryVector 就成为表示选择操作的通用方式。这是 Velox 以及其他现代向量化引擎的一项核心原则,后面我们还会经常遇到这个概念。
聚合与 Pipeline 屏障
现在来到第二阶段的第二条 pipeline。它从 LocalExchange 开始,随后向 HashAggregation 提供输入。LocalExchange 取出 Task 输入中属于自己本地目的地的那一部分,比例大约是 Task 输入的 1 / Driver 数量。
我们会在后续文章中讨论哈希表的具体布局和其他技巧。这里先将 HashAggregation 看作一个黑盒。这个例子中的聚合是最终聚合,它形成一道完整的 pipeline 屏障,只有在接收了全部输入之后才产生输出。
聚合如何知道自己已经收到了全部输入?我们沿着 shuffle 追踪一下完成信号的传播。叶子 worker 收到 coordinator 在最后一条 Task Update 中发送的“不会再有 split”消息后,Task 就知道后续不会再收到新的 split。
如果某个 TableScan 内的 DataSource 已经读到末尾,并且不会再有 split,那么这个 TableScan 既没有被阻塞,也已经到达结束状态。这会使 Driver 对该 TableScan 后面的 PartitionedOutput 调用 PartitionedOutput::noMoreInput()。于是,面向各个目的地的缓冲数据都会被 flush 并转交给 OutputBufferManager,同时附带后续不会再有数据的通知。OutputBufferManager 知道 TableScan pipeline 中有多少个 Driver。当它收到同样数量的“不会再有数据”通知后,就可以告诉所有目的地:这个 Task 不会再为它们产生数据。
接下来,第二阶段的 Task 向第一阶段的生产者请求数据。当所有生产者都发出不再有数据的信号后,第二阶段就知道输入已经结束。获取某个目的地数据的请求,其响应中包含一个标志,用于标识最后一批数据。第二阶段 Task 中的 ExchangeSource 据此设置 no-more-data 标志。所有 Driver 都会查询这一状态,各个 Exchange 算子也都会观察到它。随后就会对 LocalPartition 调用 noMoreInput(),把“不再有数据”的信号放入 local exchange 队列。如果第二阶段第二条 pipeline 开头的 LocalExchange 从自己的每一个来源都收到了这个信号,它就到达结束状态,并对 HashAggregation 调用 noMoreInput()。
数据结束信号就是这样传播的。在此之前,HashAggregation 一直没有产生输出,因为只有接收全部输入之后,才能知道最终的计数。现在,HashAggregation 开始产生输出批次,其中包含 l_partkey 值以及它出现的次数。这些输出到达最后一个 PartitionedOutput;在这个例子中,它只有一个目的地,就是负责输出结果集的最终 worker。当全部 100 个来源都报告输入结束时,最终 worker 也就到达结束状态。
回顾
我们终于完整走过了一个简单查询的分布式执行过程。我们介绍了数据如何先在集群中的 worker 之间分区,然后在每个 worker 内部再分区一次。
Velox 和 Presto 的设计目标是积极利用并行执行,也就是为每条线程创建彼此不同、互不重叠的数据集合。线程越多,吞吐量越高。同时也要记住,要让 CPU 线程有效工作,任务必须足够大——通常需要超过 100 微秒的 CPU 工作量——并且不能过多地与其他线程通信,也不应该写入其他线程正在写入的内存。Local exchange 正是实现这一点的手段。
本文另一个值得记住的要点是,列式编码,尤其是 DictionaryVector,可以用零拷贝的方式表示选择、重排和重复。我们还会在过滤、连接以及其他关系运算中再次看到这种模式。
接下来我们会继续讨论连接、过滤和哈希聚合,敬请期待!
译注与来源
1. Pipeline 与线程。第一篇把第二阶段概括为一条算子链;第三篇明确展开:LocalExchange 会形成生产者与消费者两条 pipeline。原文用“线程”解释 Driver 的执行并发,但 Driver 不永久绑定一条操作系统线程,实际由 executor 调度,等待 future 时可以离开线程。5–30 个 Driver、100 微秒等都是原文用于说明的规模和经验值。
2. API 拼写与结束状态。第二篇原文中的 addInputs() 和 task::getSplitOrFuture(),译文分别按实际拼写写为 addInput() 和 Task::getSplitOrFuture()。原文对推进顺序和结束状态作了概括;实际阅读实现时,不能把一次 getOutput() 返回空值直接理解为永久结束,还需要结合阻塞状态、isFinished() 和 noMoreInput() 协议。
3. 只读共享与零拷贝。第三篇的“无缓存一致性开销”强调只读共享不会引发多个写者之间的缓存行失效争用;共享读取仍然会消耗缓存和内存带宽。“零拷贝”指不复制底层列值,dictionary 包装和索引本身仍有构造及管理成本。“线程越多,吞吐量越高”是原文对有效并行的概括,实际还受 CPU、带宽、调度和通信成本限制。
4. 聚合示例。正文沿用原文这个完整聚合屏障的示例与术语,并不表示所有聚合模式都要等完整输入才输出,也不规定每个真实 Presto 计划都必须恰好包含图中的三个 stage。
原文与源码:Part 1、Part 2、Part 3;对应 MDX 文件位于 Velox 官方仓库。原文按仓库的 Apache License 2.0 发布。本页为中文翻译,图为重绘或译者补充,译注已单独标明。