Doris architecture

ANALYTICAL DATABASE / SOURCE READING / LESSON 04

Pipeline execution engine

Break down fragments, pipelines, operators, and BE tasks in Doris' parallel execution engine.

Reading
45 min
Track
Doris architecture
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 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.hPipelineOperator 链定义
be/src/exec/pipeline/pipeline_task.hPipelineTask::execute()Task 执行循环
be/src/exec/pipeline/pipeline_fragment_context.hPipelineFragmentContext::submit()Fragment → Pipeline 构建
be/src/exec/pipeline/operator.hOperator (Source/Sink/Transform)算子接口
be/src/exec/pipeline/dependency.hDependencyPipeline 间依赖
be/src/exec/pipeline/task_scheduler.hTaskScheduler就绪队列 / 阻塞队列调度
be/src/exec/pipeline/task_queue.hTaskQueueTask 队列数据结构
be/src/exec/operator/olap_scan_operator.hOlapScanOperatorScan Source Operator (读存储)

2. 核心概念详解

Pipeline 定义

一个 Pipeline 是一系列 有序 Operator 组成的链——数据从 Source 进入, 经过 0 到多个 Transform, 最终到达 Sink。

Pipeline {
    Source Operator  (OlapScanOperator)  ← 产生数据
    Transform        (FilterOperator)     ← 过滤行
    Transform        (ProjectOperator)    ← 投影列
    Sink Operator    (AggSinkOperator)    ← 消费数据, 聚合
}

Operator 分类

类型接口职责示例
Sourceget_block(Block* block)从某处产生数据 (磁盘/网络)OlapScanOperator, ExchangeSourceOperator
Transformget_block(), sink()过滤/投影/转换FilterOperator, ProjectOperator
Sinksink(Block* block)消费数据 (聚合/输出)AggSinkOperator, ExchangeSinkOperator, ResultSinkOperator

核心类定义

Pipeline                  — 一个 Operator 链, 不可再分割
PipelineTask              — 一个 Pipeline 的一次执行实例 (可执行单元)
Dependency                — Pipeline 之间的依赖关系 (如 ExchangeSource 等待 ExchangeSink)
TaskScheduler             — 调度器: 决定哪个 Task 获得 CPU 时间片
TaskQueue                 — 就绪队列 + 阻塞队列

Pipeline 执行引擎——查询如何在 BE 端并行执行 图 01

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 重用。迁移策略:

  1. 已经迁移: 所有 Scan/Agg/Join/Sort 的 Operator 重新实现为 Pipeline Operator (be/src/exec/operator/)
  2. 仍在使用: 部分旧算子(exec/ 目录)仍被 Pipeline Operator 内部调用
  3. 已废弃: 旧的 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 中的 InstanceNumPipelineTaskNum 观察。

Q5: 反压机制会导致死锁吗? A: 理论上 Pipeline DAG 是单向的 (无环), 所以不会死锁。但可能"活锁" — 所有 Task 都在等待某些 Resource(如内存) 但都没有释放。Doris 通过"最小保留内存 + 内存限制"来防止活锁: 每个查询有最小保留内存, 低于这个阈值时 Reservation 机制保证查询不被反压彻底挂起。


Pipeline 执行 DAG: Fragment → Pipeline → Task

下一步