Doris 实战与架构

分析型数据库 / 源码阅读 / LESSON 04

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

拆解 Pipeline 执行引擎,理解 Fragment、Pipeline、Operator 和 Task 如何在 BE 端并行运行。

阅读时间
45 分钟
学习路径
Doris 实战与架构
内容来源
Doris 深度笔记

预计阅读时间: 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

下一步