Day 18|分布式执行引擎:Fragment、Instance、Exchange、Pipeline、向量化与 Spill

Day 16 沿着一条 SQL 的生命周期走过连接、解析、绑定、优化、调度和结果返回,Day 17 深入了 Nereids 如何选择 Join 顺序、分布方式和 Runtime Filter。今天继续向下进入 BE 执行层,观察一份物理计划怎样变成真正占用 CPU、内存、磁盘和网络的运行任务。
理解执行引擎以后,很多性能现象会变得清晰:
EXPLAIN中为什么出现多个 Plan Fragment;- 同一个 Fragment 为什么会在多个 BE 上出现多个 Instance;
- Broadcast、Shuffle、Bucket Shuffle 和 Colocate 分别会产生怎样的数据流;
- Hash Join 为什么会拆成 Build 与 Probe 两条 Pipeline;
- 某个算子的
ExecTime很短,查询仍然可能长时间等待; - PipelineTask 的
max明显高于avg时,为什么要检查数据倾斜; - 向量化怎样提高单核效率,Pipeline 怎样提高多核利用率;
- 内存不足时,Spill 怎样让查询继续完成,又会带来哪些代价。
生产排障经常从“SQL 很慢”开始,最终需要把问题落到某个 Fragment、某条 Pipeline、某个 Operator 或某组 PipelineTask。今天建立的层级模型,就是完成这次落点转换的基础。
今天完成后,你应当能够:
- 解释
Plan → Fragment → Instance → Pipeline → PipelineTask → Operator的层级关系; - 根据
EXPLAIN VERBOSE识别 Exchange 边界和数据分布策略; - 说明 Hash Join、Aggregation、Sort 为什么会形成多个 Pipeline;
- 理解 Dependency、LocalState、Local Exchange 和 Scanner 并行化;
- 解释列式 Block、NullMap、连续内存与 SIMD 的性能价值;
- 使用 MergedProfile 与 DetailProfile 判断扫描、网络、等待、内存、倾斜和 Spill;
- 为大查询设计一套可验证、可回滚的执行侧调优方案。

一、从物理计划到运行任务:先建立六层执行模型

一条 SQL 在 FE 中完成优化后,会得到物理计划。物理计划仍然是一棵描述“做什么”的算子树,真正执行还要继续经历分布式切分、实例化和本地调度。
可以用六层结构理解当前 Doris 4.x 的执行体系:
| 层级 | 作用 | 典型对象 |
|---|---|---|
| Physical Plan | 描述物理算子与物理属性 | Scan、Hash Join、Aggregate、Sort |
| PlanFragment DAG | 按跨节点数据边界拆分计划 | Fragment 0、Fragment 1、Fragment 2 |
| Fragment Instance | 某个 Fragment 在某个 BE 上的运行副本 | F1-I0、F1-I1 |
| Pipeline | 由非阻塞 Operator 连接成的执行链 | Scan → Filter → Join Probe |
| PipelineTask | 可提交到线程池的调度实体 | Task-0、Task-1 |
| Operator | 执行具体逻辑的最小单元 | OlapScan、HashJoin、Agg、Exchange |
这六层经常被混用,排障时需要严格区分。
Fragment 解决跨节点分工。 一个查询可能包含扫描阶段、Join 阶段、聚合阶段和最终汇总阶段。数据需要重新分布的位置会插入 Exchange,完整计划随之被切成多个 Fragment。
Instance 解决水平并行。 同一个 Fragment 可以在多个 BE 上运行,每个运行副本处理一部分 Tablet、Bucket、文件 Split 或上游数据流。
Pipeline 解决节点内算子编排。 BE 收到 Fragment 后,会把其中能够连续处理 Block 的 Operator 串起来;遇到 Hash Table 构建、全量聚合状态、排序缓冲等阻塞边界时,再拆成不同 Pipeline。
PipelineTask 解决 CPU 调度。 Pipeline 属于逻辑结构,实例化后才形成多个 Task。每个 Task 拥有自己的 LocalState,进入共享线程池,在 Ready、Running、Blocked 和 Yield 等状态之间切换。
Operator 负责真实计算。 Scan 读取数据,Filter 计算谓词,HashJoin 进行 Build 或 Probe,Aggregation 更新聚合状态,Exchange 发送或接收 Block。
因此,看到 HASH_JOIN 耗时高时,还需要继续确认:问题位于 Build 还是 Probe,哪个 Fragment 承担它,多少个 Task 在运行,是否只有一个 Task 处理了大部分数据,以及上游是否已经发生等待或倾斜。
二、PlanFragment:分布式计划的阶段边界
官方 MPP 文档将 PlanFragment 定义为 FE 下发给 BE 的最小分布式工作单元。一条查询会形成一个 Fragment DAG,Fragment 之间通过 DataStreamSink 与 ExchangeNode 连接。[1]

2.1 Fragment 为什么会出现
以下操作通常要求数据改变分布:
- 两张大表按 Join Key 重新分区;
- 局部聚合结果按 Group Key 汇总到下一阶段;
- 多个节点的 TopN、Sort 或聚合结果集中合并;
- Broadcast 小表发送到多个 Probe Instance;
- 最终结果集中到 Root Fragment 返回客户端。
每出现一次跨节点分布变化,计划中就会形成一对发送端和接收端:
上游 Fragment
DataStreamSink
│ BRPC / Block
▼
下游 Fragment
ExchangeNode
Fragment 形成 DAG,而非简单链表。一个 Join Fragment 可能同时接收事实表扫描和维表广播,一个 Root Fragment 也可能同时接收聚合结果与其他分支结果。
2.2 Root Fragment 的职责
Root Fragment 位于 DAG 顶部,常见职责包括:
- 合并多个 BE 的局部聚合状态;
- 完成最终 Sort、TopN、Limit;
- 格式化最终结果;
- 通过 Result Sink 把数据返回给 FE,再由 MySQL 协议流式返回客户端。
Root Fragment 经常只有一个 Instance。若上游输出仍然很大,Root 会成为集中瓶颈。下面几种写法尤其需要关注:
- 没有合理 Limit 的全局排序;
- 高基数 Group By 后仍输出大量结果;
- 将数百万行明细直接返回 BI 工具;
- 上游没有完成局部聚合,所有原始行集中到 Root。
优化 Root 的核心方法包括局部预聚合、TopN 下推、减少返回列、限制结果规模,以及让客户端采用合理 Fetch Size 持续消费。
2.3 Fragment 数量多代表什么
Fragment 多并不自动说明计划差。复杂 SQL、外部 Catalog、多阶段聚合和多个 Join 都可能产生更多 Fragment。真正需要观察的是:
- 每个 Fragment 的输入规模;
- Fragment 之间发送了多少字节;
- Distribution 是否符合数据规模;
- Root 是否承担了过多集中计算;
- 某个 Exchange 是否形成长时间等待;
- Fragment DAG 中是否存在明显串行关键路径。
三、Fragment Instance:同一份工作怎样铺到多个 BE
PlanFragment 是模板,Fragment Instance 是运行副本。Coordinator 根据数据位置、Replica 状态、BE 负载、并行度和分布策略,将同一个 Fragment 实例化到多个 BE。[1]
假设事实表共有 96 个 Tablet,分布在 6 个 BE 上。扫描 Fragment 可能在 6 个 BE 上各创建若干 Instance,每个 Instance 负责本地 Tablet 的一部分扫描范围。上层 Join Fragment 也可能在相同或不同 BE 集合上创建多个 Instance。
Instance 数量增大可以提升并行能力,同时会增加以下成本:
- PipelineTask 数量;
- Hash Table、聚合状态和排序缓冲的副本数量;
- Exchange Sender 与 Receiver 数量;
- RPC、队列和调度开销;
- 高并发场景下的 CPU 争用。
因此,并行度存在合理区间。大扫描在资源空闲时可以利用更高并行度,高并发小查询通常需要更克制的 Task 数量。
当前 Profile 文档使用 parallel_pipeline_task_num 控制单个 Pipeline 的并行 Task 数,并以它演示 MergedProfile 中 instance_num 的变化。[2]
SHOW VARIABLES LIKE 'parallel_pipeline_task_num';
SET parallel_pipeline_task_num = 2;
生产调优建议从较低值逐步增加,例如 1、2、4、8,并同时记录:
- 查询 P50、P95、P99;
- BE CPU 利用率;
- Context Switch 与线程池等待;
- Scan 吞吐;
- Exchange 字节;
- 单个 Task 的
max/avg/min; - 同时运行的其他查询数量。
单条查询加速后,如果整体集群吞吐下降,说明并行度已经侵占了其他工作负载的资源。
四、Exchange:分布式执行最昂贵的边界之一

Doris 在 Fragment 边界使用 DataStreamSink 发送 Block,使用 ExchangeNode 接收 Block。数据经过序列化后通过 BRPC 在 BE 之间传输。[1]
常见分布模式可以理解为四类。
4.1 UNPARTITIONED:集中汇总
所有上游数据进入一个下游 Instance。Root Merge、全局聚合和最终结果返回经常使用这种方式。
风险很直观:上游并行度再高,最终仍然要经过单点汇总。观察 Profile 时应重点查看 Root Exchange 的输入行数、远端接收字节、等待时间和下游 Sort/Aggregation 的峰值内存。
4.2 RANDOM:均匀扩散
数据以轮询或随机方式分给多个下游 Instance,目标是维持负载均衡。它不保证相同 Key 落到同一节点,因此不适合需要 Key 对齐的最终 Join 或聚合阶段。
4.3 HASH_PARTITIONED:按 Key 重分布
上下游按照 Join Key 或 Group Key 计算 Hash,相同 Key 进入同一目标 Instance。大表 Join 与全局 Group By 经常使用这种方式。
其代价包含:
网络成本
≈ 参与 Shuffle 的数据字节数
+ 序列化 / 反序列化 CPU
+ Sender / Receiver 缓冲
+ 网络等待与背压
如果过滤和局部聚合没有提前完成,Shuffle 会携带大量无效行。Runtime Filter、谓词下推、列裁剪和两阶段聚合都可以在 Exchange 前缩小数据。
4.4 Bucket Shuffle:复用一侧分桶
当一侧表已经按 Join Key 分桶,Bucket Shuffle 可以保持这一侧数据不动,只把另一侧重新发送到对应 Bucket 的 BE。它减少一侧网络移动,在大事实表与中等规模表 Join 时很有价值。
Colocate 条件满足时,两侧 Tablet 已经稳定同位,Join 可以在本地完成,网络 Shuffle 进一步消失。Day 17 已经讲过优化器如何选择这些策略,Day 18 更关注它们在 Exchange 指标中的真实代价。
4.5 Exchange 的背压
下游处理速度跟不上时,Sender 缓冲区会逐渐占满,上游 Task 随后进入等待。此时上游算子的 ExecTime 可能不高,查询仍然很慢,Profile 中会出现:
WaitForRpcBufferQueue;WaitForData;RemoteBytesReceived;SendersBlockedTotalTimer;- 某些 Pipeline 长时间
WaitForDependency。
排查网络瓶颈需要同时检查发送端和接收端。接收端的 Sort、Hash Join 或聚合处理缓慢,也会通过背压反映到发送端。
4.6 Block 在网络中的完整路径
Exchange 传输的基本单位仍然是列式 Block。发送端完成上游算子计算后,会根据目标分布方式把 Block 中的行路由到一个或多个下游通道。Hash Shuffle 通常要先计算分区表达式,再按目标 Instance 重组数据;Broadcast 会把同一份 Build 数据复制到多个下游;本地发送可以绕过部分远程协议开销,远程发送则需要经过序列化、BRPC、接收队列和反序列化。
一条 Block 的路径可以写成:
Upstream Operator
→ DataStreamSink 分区
→ Sender Buffer
→ BRPC Channel
→ Receiver Queue
→ ExchangeOperator
→ Downstream Pipeline
每个环节都有容量和速度差异。Sender 生产太快时,缓冲区会形成背压;Receiver 迟迟拿不到数据时,会出现首批数据等待;Block 太小会增加 RPC 和调度次数,Block 太大又会提高瞬时内存与网络抖动。执行引擎需要在吞吐、延迟和内存之间保持平衡。
排查时可以把问题拆成三段:
- 发送前:上游计算是否过慢,分区表达式是否昂贵;
- 传输中:远程字节、网络带宽、RPC 队列是否异常;
- 接收后:下游 Join、聚合或排序是否消费不及时。
仅看网络监控无法判断根因。下游算子阻塞会反向造成发送端等待,发送端计算缓慢也会表现为接收端长期无数据。
五、Pipeline:把一个 Fragment 拆成可调度的算子链
官方 Pipeline 文档说明,Doris 自 3.0 起已经由 Pipeline 执行模型完整替代原火山模型。BE 收到 PlanFragment 后,会继续切分为多个 Pipeline,再实例化为 PipelineTask。[3]

5.1 一条 Pipeline 的结构
一条 Pipeline 通常包含:
一个 SourceOperator
→ 若干中间 Operator
→ 一个 SinkOperator
Source 可以来自本地表扫描、外部文件扫描或 Exchange Buffer;中间 Operator 可以包含 Filter、Project、Join Probe;Sink 可以把数据发往网络、写入 Hash Table、写入聚合状态或写入排序缓冲。
可以连续消费一个 Block 并立即产生下一个 Block 的 Operator,适合放在同一 Pipeline 中。这样能够减少中间物化、内存复制和线程切换。
5.2 阻塞算子为什么要拆开
Hash Join、Aggregation 和 Sort 都包含“先积累状态,再向下输出”的阶段。
官方文档给出了以下拆分关系:[3]
| PlanNode | 执行层 Operator |
|---|---|
| JoinNode | JoinBuildOperator + JoinProbeOperator |
| AggNode | AggSinkOperator + AggSourceOperator |
| SortNode | SortSinkOperator + SortSourceOperator |
以 Hash Join 为例:
Build Pipeline
Scan dim → HashJoin Build Sink → HashTable
Probe Pipeline
Scan fact → HashJoin Probe → downstream operators
Probe 需要等待 Hash Table Ready。两条 Pipeline 之间通过 Dependency 表达这项约束。Build 完成后调用 set_ready,Probe Task 才进入可运行状态。[3]
5.3 ExecTime 与等待时间需要分开理解
某个 Probe Operator 的 ExecTime 可能只有 500 ms,查询总时间却达到 8 秒。差额可能来自 7.5 秒的 Dependency 等待。常见等待来源包括:
- Hash Join Build 尚未完成;
- Exchange 数据尚未到达;
- Local Exchange Buffer 暂无数据;
- Sender 受到下游背压;
- Scanner 尚未生产足够 Block;
- Spill 正在写盘或读盘。
Profile 的 WaitForDependency 经常比 ExecTime 更能解释关键路径。阅读顺序建议采用:
Total Time
→ 最慢 Fragment
→ 最慢 Pipeline
→ WaitForDependency
→ 上游依赖
→ 对应 Operator 的行数、字节、内存和网络
六、PipelineTask、LocalState 与共享线程池
Pipeline 是逻辑结构,PipelineTask 才能进入线程池运行。同一 Pipeline 的多个 Task 拥有相同 Operator 结构,处理的数据范围和 LocalState 不同。[3]
LocalState 可能包含:
- 当前 Scanner 分配到的 Split;
- 当前 Task 的局部聚合状态;
- Join Build 或 Probe 的局部状态;
- Exchange Sender / Receiver 通道;
- Operator 的计时器、行数、内存计数器;
- Spill 分区和文件状态。
Pipeline 调度减少了“一查询一组固定线程”带来的线程膨胀。Task 在遇到依赖未满足、数据暂缺或缓冲区满时,可以让出 CPU;依赖就绪后再回到队列。固定规模的线程池可以被多个查询共享,CPU 核心由当前 Ready Task 使用。
这套模型的工程价值包括:
- 控制活跃线程数量;
- 降低大量查询并发时的上下文切换;
- 提升多核利用率;
- 让等待网络、等待 Scan 和等待 Build 的任务及时让出 CPU;
- 为 Workload Group、队列和资源治理提供调度基础。
高并发场景仍然需要限制查询数量和单查询并行度。Pipeline 解决线程模型问题,无法消除 CPU、内存和磁盘总量约束。
6.1 Task 状态怎样影响 CPU 利用率
PipelineTask 在一次调度时间片内会持续拉取和推送 Block,直到遇到完成条件、依赖未就绪、缓冲区受限或主动让出。调度器随后把 CPU 分配给其他 Ready Task。这个过程让网络等待、磁盘等待和计算任务可以交错推进。
可以用四种状态理解运行过程:
| 状态 | 含义 | 常见原因 |
|---|---|---|
| Ready | 已满足依赖,等待线程池 | CPU 繁忙、队列排队 |
| Running | 正在执行 Operator 链 | Scan、Hash、聚合、序列化 |
| Blocked | 当前无法继续 | 等 Build、等 Exchange、等 Scanner、等内存 |
| Finished | 数据和清理均完成 | PipelineTask 结束 |
CPU 利用率低并不总是线程数量不足。大量 Task 处于 Blocked 时,提高并行度可能增加等待者数量,无法增加有效计算。此时应查看 Dependency、I/O、网络或内存申请。CPU 已经接近饱和时继续提高 Task 数,常见结果是上下文切换增加、缓存局部性下降和 P99 波动扩大。
6.2 取消、错误与资源清理
查询被用户 KILL、超时、Workload Policy 熔断或某个 Operator 报错后,取消信号会沿查询上下文传播到相关 Fragment Instance 和 PipelineTask。执行引擎需要停止 Scanner、关闭网络通道、释放 Hash Table 与聚合状态、删除临时 Spill 文件,并向 Coordinator 汇报最终状态。
生产环境中,频繁取消大查询仍会产生资源尾部:已经发出的 RPC、正在进行的 I/O、异步 Profile 上报和临时文件清理都需要时间。故障演练应观察取消后的 BE 内存、活动 Task、Spill 目录和磁盘占用是否恢复,避免把“客户端收到取消结果”直接等同于“所有资源已经释放”。
