预计阅读时间: 45 分钟 前置阅读: doc-01 §7 (Pipeline概览), doc-03 (计划来源) 下一次阅读: doc-05 (向量化计算), doc-06 (存储)
1. Pipeline 引擎 vs 旧 ExecModel
Doris 从 2.0 起重新设计了 BE 的执行引擎: 从传统的"线程专有"模型 (每个 Fragment Instance 独占一个线程) 演化为 "Pipeline" 模型 (多个查询的 Task 交织执行)。
| 维度 | 旧 ExecModel (be/src/exec/) | Pipeline 引擎 (be/src/exec/pipeline/) |
|---|---|---|
| 线程分配 | 每个 Fragment Instance 1 个线程 | 多个 Task 竞争线程池 |
| 资源利用 | 单个慢查询独占线程 → 阻塞 | 线程数恒定, 慢查询 YIELD 让出线程 |
| 内存管理 | 查询全量数据在内存 (易 OOM) | Pipeline 分批处理, 控制峰值内存 |
| 并发能力 | 线程数 = 并发 Fragment 数 | 线程数固定, 支持更多并发查询 |
| 反压 | 无 | 通过 Dependency 实现反压 |
| 代码位置 | be/src/exec/ (大部分保留) | be/src/exec/pipeline/ (新路径) |
源码导航
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
be/src/exec/pipeline/pipeline.h | Pipeline | Operator 链定义 |
be/src/exec/pipeline/pipeline_task.h | PipelineTask::execute() | Task 执行循环 |
be/src/exec/pipeline/pipeline_fragment_context.h | PipelineFragmentContext::submit() | Fragment → Pipeline 构建 |
be/src/exec/pipeline/operator.h | Operator (Source/Sink/Transform) | 算子接口 |
be/src/exec/pipeline/dependency.h | Dependency | Pipeline 间依赖 |
be/src/exec/pipeline/task_scheduler.h | TaskScheduler | 就绪队列 / 阻塞队列调度 |
be/src/exec/pipeline/task_queue.h | TaskQueue | Task 队列数据结构 |
be/src/exec/operator/olap_scan_operator.h | OlapScanOperator | Scan Source Operator (读存储) |
2. 核心概念详解
Pipeline 定义
一个 Pipeline 是一系列 有序 Operator 组成的链——数据从 Source 进入, 经过 0 到多个 Transform, 最终到达 Sink。
Pipeline {
Source Operator (OlapScanOperator) ← 产生数据
Transform (FilterOperator) ← 过滤行
Transform (ProjectOperator) ← 投影列
Sink Operator (AggSinkOperator) ← 消费数据, 聚合
}
Operator 分类
| 类型 | 接口 | 职责 | 示例 |
|---|---|---|---|
| Source | get_block(Block* block) | 从某处产生数据 (磁盘/网络) | OlapScanOperator, ExchangeSourceOperator |
| Transform | get_block(), sink() | 过滤/投影/转换 | FilterOperator, ProjectOperator |
| Sink | sink(Block* block) | 消费数据 (聚合/输出) | AggSinkOperator, ExchangeSinkOperator, ResultSinkOperator |
核心类定义
Pipeline — 一个 Operator 链, 不可再分割
PipelineTask — 一个 Pipeline 的一次执行实例 (可执行单元)
Dependency — Pipeline 之间的依赖关系 (如 ExchangeSource 等待 ExchangeSink)
TaskScheduler — 调度器: 决定哪个 Task 获得 CPU 时间片
TaskQueue — 就绪队列 + 阻塞队列

Pipeline 依赖图 (DAG)
Fragment 2 (Leaf)
┌──────────────────────────┐
│ Pipeline A │
│ OlapScanOp → FilterOp → │ ExchangeSinkOp ──────┐
│ │ │
└──────────────────────────┘ │
│
Fragment 1 (Intermediate) │
┌──────────────────────────────────────┐ │
│ Pipeline B │ │
│ ExchangeSourceOp ← ─ ─ ─ ─ ─ ─ ─ ─ ┘ │
│ → AggOp → ExchangeSinkOp ──────────────────┐ │
└──────────────────────────────────────┘ │ │
│ │
Fragment 0 (Root) │ │
┌───────────────────────────────────────┐ │ │
│ Pipeline C │ │ │
│ ExchangeSourceOp ← ─ ─ ─ ─ ─ ─ ─ ─ ─ ┘ │ │
│ → ResultSinkOp │ │
└───────────────────────────────────────┘ │ │
Dependency 的工作原理: ExchangeSource 依赖 ExchangeSink — 数据没到时 Source 阻塞, Pipeline Task 被移入 Blocked Queue; 数据到达后 Dependency 被更新, Task 移入 Ready Queue, 获得线程后继续。
3. PipelineX——LocalExchange 支持
PipelineX 是 Pipeline 的扩展, 支持 单 BE 内部的多 Pipeline 并行:
Fragment X
┌────────────────────────────────────┐
│ Pipeline A_0: Scan → Shuffle │ ─┐
│ Pipeline A_1: Scan → Shuffle │ │ LocalExchange (本 BE 内 Shuffle)
│ Pipeline A_2: Scan → Shuffle │ │
│ │ │
│ Pipeline B: Shuffle Receive → Agg │ ◄┘
└────────────────────────────────────┘
与 Distribute Exchange 的区别: Distribute Exchange 跨 BE 通过 Brpc 传输数据; LocalExchange 在本 BE 内通过内存/共享 Queue 传输, 延迟低。
4. 线程模型——Driver/WorkerPool/ScheduleUnit
核心线程模型
┌──────────────────────────────────────┐
│ Thread Pool (WorkerPool) │
│ Worker 0: [执行 Task A 100ms] → YIELD│
│ Worker 0: [执行 Task C 100ms] → YIELD│
│ Worker 1: [执行 Task B 100ms] → YIELD│
│ Worker 1: [执行 Task A 100ms] → YIELD│
│ │
│ 关键: 所有 Worker 只处理 Ready Queue │
│ 中的 Task, 每个 Task 最多执行 │
│ TIME_SLICE (默认 100ms) 后 yield │
│ │
│ ┌─────────────┐ ┌──────────────────┐│
│ │ Ready Queue │◄──│ Blocked Queue ││
│ │(竞争式执行) │ │ (Dependency 满足 ││
│ │ │ │ 后移入 Ready) ││
│ └─────────────┘ └──────────────────┘│
└──────────────────────────────────────┘
TaskScheduler
// task_scheduler.h — 调度核心
class TaskScheduler {
// 核心调度循环
void schedule() {
while (!stopped) {
// 1. 从 Ready Queue 取一个 Task
PipelineTask* task = ready_queue.take();
// 2. 检查 Dependency (如阻塞, 移入 Blocked Queue)
// 3. 执行 Task (最多 100ms)
task->execute(&done);
if (!done) {
// 未完成 → 放回 Ready Queue 尾部
ready_queue.put(task);
}
}
}
};
5. 反压机制实现
为什么需要反压
合并阶段的 Agg 生产速度 < Scan 阶段的数据生产速度 → 中间 Buffer 不断增长 → OOM。
反压流程
1. ExchangeSink 的 Output Buffer 满 (超过 buffer_limit)
2. ExchangeSink 设置 Dependency 的 "BLOCKED" 状态
3. 上一个 Pipeline Task (O) 调用 get_block() → 返回 BLOCKED
4. Task Scheduler 发现 Task BLOCKED → 移入 Blocked Queue
5. Scan 停止生产数据, 释放内存
6. 当 Output Buffer 有空间时
7. Dependency 状态更新 → Task 移入 Ready Queue
8. Scan 恢复生产
关键代码
// ExchangeSinkOperator::sink() 伪代码
Status ExchangeSinkOperator::sink(Block* block) {
// 1. 尝试将 block 写入 Output Buffer
if (output_buffer->is_full()) {
// 2. Buffer 满 → 等待 (block task)
return Status::Yield(); // Task 自动移入 Blocked Queue
}
output_buffer->add_block(block);
return Status::OK();
}
6. 实例追踪——一条查询的完整 Pipeline 执行生命周期
时点 状态
────────────────────────────────────────────────────────────────
T0 FE 发送 TPipelineFragmentParams (Thrift) 到 BE
T1 BE FragmentMgr::exec_plan_fragment(params)
T2 PipelineFragmentContext 创建
T3 Pipeline DAG 构建 (Exchange 边界 → 多个 Pipeline)
T4 PipelineTask 创建 (每个 Pipeline 1-N 个 Task)
T5 Dependency 构建 (ExchangeSource 依赖 ExchangeSink)
T6 PipelineFragmentContext::submit()
→ 所有 Root Pipeline Tasks 推入 Ready Queue
T7 TaskScheduler 分配 Worker 线程
T8 ~100ms PipelineTask::execute() — 第一个 100ms 时间片
T9 ~200ms PipelineTask::execute() — 第二个 100ms 时间片
... ...
T_final 所有 PipelineTask 完成
→ ResultSink 发送结果到 FE
→ PipelineFragmentContext 销毁
7. 与旧执行引擎 be/src/exec/ 的迁移路径
虽然 Pipeline 已成为默认引擎, 但旧的 be/src/exec/ 算子仍然存在且被 Pipeline 重用。迁移策略:
- 已经迁移: 所有 Scan/Agg/Join/Sort 的 Operator 重新实现为 Pipeline Operator (
be/src/exec/operator/) - 仍在使用: 部分旧算子(
exec/目录)仍被 Pipeline Operator 内部调用 - 已废弃: 旧的 ExecNode tree 执行逻辑 (
exec/exec_node.cpp) 不再用于新查询
如果你在阅读代码时看到两套路径, 优先看 be/src/exec/operator/ (新的 Pipeline Operator 实现) 和 be/src/exec/pipeline/ (Pipeline 框架)。
8. 常见问题 / 面试题
Q1: 为什么 TIME_SLICE 是 100ms 而不是 1s 或 5ms? A: 太长 → 其他查询的 Task 得不到 CPU, 短查询的延迟增大。太短 → 上下文切换频繁 (每次 yield 都有 thread context switch 开销)。100ms 在延迟敏感度和吞吐量之间取得平衡。
Q2: Pipeline 模型如何处理数据倾斜? A: 数据倾斜(某些 Key 的数据量远大于其他 Key)在 Pipeline 中体现为某个 PipelineTask 的 YIELD 次数远多于其他 Task。TaskScheduler 的公平调度 (Round-Robin ready queue) 不能解决倾斜, 但也不能加剧本地倾斜。跨 BE 的倾斜需要通过数据重分布 (Repartition/Shuffle) 解决。
Q3: 什么情况下 PipelineTask 会进入 Blocked Queue? A: (1) 等待 Exchange 数据 (ExchangeSource 等待上游 ExchangeSink); (2) Output Buffer 满 (反压); (3) IO 等待 (Scan 等待磁盘); (4) 等待 RuntimeFilter 到达; (5) 等待 Shared Hash Table (Join Build 端未完成)。
Q4: 并行度怎么控制?
A: Pipeline 的数量由 (1) Leaf Fragment 的 Tablet 数量(每个 Tablet 一个 Scan PipelineTask); (2) 中间 Fragment 的 Partition 数 (SHUFFLE 的分桶数) 决定。可以通过查询 profile 中的 InstanceNum 和 PipelineTaskNum 观察。
Q5: 反压机制会导致死锁吗? A: 理论上 Pipeline DAG 是单向的 (无环), 所以不会死锁。但可能"活锁" — 所有 Task 都在等待某些 Resource(如内存) 但都没有释放。Doris 通过"最小保留内存 + 内存限制"来防止活锁: 每个查询有最小保留内存, 低于这个阈值时 Reservation 机制保证查询不被反压彻底挂起。

下一步
- 向量化执行: doc-05-vectorized-execution.md——Pipeline 中每个 Operator 如何进行向量化计算
- 存储扫描: doc-06-storage-tablet-rowset-segment.md——OlapScanOperator 如何读数据