Doris 实战与架构

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

SQL 查询端到端——从 MySQL 协议到磁盘读取的全源码级穿越

追踪一条 SQL 如何从 MySQL 协议进入 FE,再被优化、调度并在 BE 端读盘返回。

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

预计阅读时间: 90-120 分钟 前置阅读: doc-00-architecture-overview.md——架构概览与代码地图 下一次阅读: doc-02(写入链路), doc-03(优化器), doc-04(Pipeline)


1. MySQL Wire Protocol——SQL 如何进入系统

架构定位

MySQL Wire Protocol 是所有客户端请求的唯一入口(无论是 mysql CLI、JDBC 还是 HTTP 转发)。Doris 实现了完整的 MySQL 协议, 包括鉴权、SSL、Prepared Statement 和 Multi-ResultSet。

源码导航

文件关键类/方法职责
fe/fe-core/.../mysql/MysqlServer.javaMysqlServer, accept()启动 MySQL 服务, 绑定端口(默认 9030)
fe/fe-core/.../mysql/nio/NMysqlServer.javaNMysqlServerNIO Accept 循环 (Java NIO Selector)
fe/fe-core/.../mysql/MysqlChannel.javaMysqlChannel单连接读写, MySQL Packet 编解码
fe/fe-core/.../qe/ConnectProcessor.javaprocessOnce()读一个 packet, 分发到 handleQuery/handleStmtPrepare 等
fe/fe-core/.../qe/ConnectContext.javaConnectContext连接上下文: 当前用户、数据库、Session 变量

调用链

MysqlServer.accept()
  → NMysqlServer (NIO accept loop)
    → MysqlChannel.read()            // 读取 MySQL packet
      → ConnectProcessor.processOnce() // 单 packet 处理循环
        → dispatch()                   // 按 MySQL 命令类型分发
          → handleQuery(stmt)          // COM_QUERY: 普通 SQL
          → handleStmtPrepare(stmt)    // COM_STMT_PREPARE
          → handleStmtExecute(param)   // COM_STMT_EXECUTE

关键代码片段

// ConnectProcessor.java:822 — 核心循环入口
// 每次处理一个 MySQL packet, 完成后等待下一个
public void processOnce() throws IOException, NotImplementedException {
    // 1. 鉴权(packet 0 → handshake)
    // 2. 循环: read packet → dispatch → handleQuery/DDL/...
    dispatch();
}

// ConnectProcessor.java:217 — handleQuery 入口
protected void handleQuery(String originStmt) throws ConnectionException {
    // 1. 解析 OriginStatement
    // 2. new StmtExecutor(ctx, parsedStmt)
    // 3. executor.execute() — 进入 SQL 执行管线
}

时序图

SQL 查询端到端——从 MySQL 协议到磁盘读取的全源码级穿越 图 01

11. 常见问题 / 面试题(汇总)

Q: Doris 为什么选 MySQL 协议而不是 PostgreSQL? A: MySQL 是中文互联网最流行的数据库协议(JDBC/ODBC/Python 生态成熟), 用户不需要额外安装驱动。PostgreSQL 协议虽然功能更强(如完整类型系统), 但现有 MySQL 工具链可直接连接 Doris。

Q: processOnce 是单线程处理一个连接吗? A: 是。每个 MySQL 连接对应一个 ConnectProcessor 实例, 由 ConnectScheduler 线程池调度。一个连接同一时刻只有一个线程在处理, 但多个连接可以并发处理。这是标准的 MySQL 线程模型。

Q: 查询超时怎么处理? A: ConnectContext.getExecTimeout() 设置最大执行时间, 超时后 Coordinator 向所有涉及的 BE 发 Cancel RPC, 停止 Pipeline 执行。


2. Nereids 语法解析——SQL String → AST

架构定位

SQL 字符串首先由 ANTLR 语法解析器转换为抽象语法树(AST)。Doris 使用自己的 ANTLR4 语法定义 (gensrc/antlr4/), 生成 Lexer 和 Parser。

源码导航

文件关键类/方法职责
gensrc/antlr4/org/apache/doris/nereids/DorisLexer.g4ANTLR Lexer 规则Token 定义
gensrc/antlr4/org/apache/doris/nereids/DorisParser.g4ANTLR Parser 规则SQL 语法规则
fe/fe-core/.../nereids/parser/NereidsParser.javaparseSingle(), parseMultiple()调用 ANTLR 解析 SQL
fe/fe-core/.../nereids/trees/plans/commands/各种 StatementBase 子类AST 节点(Query, Insert, DDL...)
fe/fe-core/.../nereids/trees/expressions/Expression 子类表达式 AST

调用链

NereidsParser.parseSingle(sql)
  → DorisLexer (ANTLR)           // SQL → Token 流
    → DorisParser (ANTLR)        // Token 流 → ParseTree
      → AstBuilder.visit()       // ANTLR Visitor 模式
        → StatementBase          // 返回 AST 根节点

AST 节点体系

StatementBase
├── QueryStatement          (SELECT ...)
├── InsertIntoTableCommand  (INSERT INTO ...)
├── CreateTableCommand      (CREATE TABLE ...)
├── AlterTableCommand       (ALTER TABLE ...)
└── ExplainCommand          (EXPLAIN ...)

SELECT COUNT(*) FROM t WHERE a > 10 GROUP BY b 为例的 AST:

QueryStatement
└── LogicalResultSink          ← 最外层
    └── LogicalAggregate       ← GROUP BY b, COUNT(*)
        └── LogicalFilter      ← WHERE a > 10
            └── LogicalOlapScan ← FROM t

关键代码

// NereidsParser.java — 解析入口
public List<StatementBase> parseMultiple(String sql) {
    // parseSingle 调用 ANTLR, 返回单个 StatementBase
    return parseSingle(sql).stream().collect(...)
}

11. 常见问题 / 面试题(汇总)

Q: 为什么不直接用 Calcite 的 SQL Parser? A: Nereids 的目标是全自研优化器, 不依赖 Calcite。自研 Parser 可以精确支持 Doris 特有的 SQL 语法(Doris 的 DISTRIBUTED BY、PROPERTIES 等), 对 ANTLR 的 Visitor 生成有完全的语法控制。

Q: SQL hint 怎么解析和生效? A: Hint 在 Parser 阶段作为特殊的 CommentOrHint 节点附着在 StatementBase 上, 后续在 Optimizer 阶段被规则消费(如 leading hint 影响 Join 顺序)。


3. Nereids 语义分析——AST → Logical Plan

架构定位

Analyzer 阶段将语法上合法的 AST 转换为语义上正确的 Logical Plan。核心工作: Bind(名字绑定到 Catalog 实体)、Resolve(类型推断/表达式类型检查)、Rewrite(表达式标准化)。

源码导航

文件关键类/方法职责
fe/fe-core/.../nereids/analyzer/Analyzer.javaAnalyzer分析入口, 驱动 Scope 和 Binder
fe/fe-core/.../nereids/analyzer/Scope.javaScope作用域: 列名 → 表的映射
fe/fe-core/.../nereids/analyzer/CatalogBinder.javaCatalogBinder将表名绑定到 Catalog 中的 OlapTable
fe/fe-core/.../nereids/analyzer/RelationManager.javaRelationManager管理分析过程的表/子查询关系
fe/fe-core/.../nereids/rules/analysis/Bind* Rule 集合各种 Bind 规则

调用链

NereidsPlanner.analyze()
  → Analyzer.analyze(plan)
    → 1. Bind: 将 names 绑定到 Catalog 中的实体
         CatalogBinder.bind("t") → Catalog.getDb("db").getTable("t")
    → 2. Resolve: 确定每列的数据类型
         a → INT, b → VARCHAR(100)
    → 3. Rewrite: 表达式标准化
         a > 10 → GreaterThan(SlotRef(a), Literal(10))
    → 返回 BoundPlan

关键代码

// NereidsPlanner.java:288-292 — analyze 流程
analyze();
// 如果 analyzedPlan 成功, 包含完整的类型和 Catalog 引用
analyzedPlan = cascadesContext.getRewritePlan();

数据流转图

SQL 查询端到端——从 MySQL 协议到磁盘读取的全源码级穿越 图 02

11. 常见问题 / 面试题(汇总)

Q: 如果表不存在, 在哪一步报错? A: 在 Bind 阶段——CatalogBinder.bind("non_existent_table") 时抛 AnalysisException("Unknown table")

Q: 子查询怎么分析? A: 子查询创建一个子 Scope(嵌套作用域), 内部可以引用外部 Scope 的列(correlated subquery), 外部不能引用内部的列。


4. Nereids 优化器——Logical Plan → Optimized Plan

架构定位

Nereids 使用 Cascades 框架的优化器, 核心资源是 Memo(存储所有等价计划组)和 Rule(定义等价变换规则)。这是整个查询系统最复杂的组件(200+ Java 文件)。

源码导航

文件关键类/方法职责
fe/fe-core/.../nereids/jobs/executor/Optimizer.javaOptimizer优化流程编排: explore → optimize
fe/fe-core/.../nereids/memo/Memo.javaMemo存储所有等价计划 Group
fe/fe-core/.../nereids/memo/Group.javaGroup一个等价组(语义等价的多个 Plan)
fe/fe-core/.../nereids/memo/GroupExpression.javaGroupExpressionGroup 中的一个 Plan, 即一个算子树节点
fe/fe-core/.../nereids/rules/Rule.javaRule规则接口 → 子类包括 RBO 和 CBO
fe/fe-core/.../nereids/rules/RuleType.javaRuleType规则类型枚举(100+ 项)
fe/fe-core/.../nereids/cost/CostCalculator.javaCostCalculator计算计划的 CPU/IO/Network/内存成本
fe/fe-core/.../nereids/stats/Statistics.javaStatistics列统计信息: 行数/NDV/null 比例/直方图

Cascades 框架核心流程

Optimizer.optimize()
  → Phase 1: EXPLORATION (探索等价逻辑计划)
      ┌─────────┐  Rule Match   ┌──────────────┐
      │  Memo   │ ◄─────────── │  Rule 触发     │
      │ (计划组) │ ────────────► │  新 Plan → Memo │
      ├─────────┤  Plan Insert  └──────────────┘
      │ Group 0 │               FilterPushDown
      │ Group 1 │               ColumnPruning
      │ Group 2 │               JoinReorder
      │  ...    │               ...
      └─────────┘

  → Phase 2: IMPLEMENTATION (逻辑到物理的转换)
      LogicalXxx → PhysicalXxx
      LogicalOlapScan → PhysicalOlapScan + DistributionSpec

  → Phase 3: OPTIMIZATION (选择最优物理计划)
      Cascades Optimize: 比较 Group 中的所有 Expression,
      选择 Cost 最小的

核心 Rule 清单(部分重点)

规则类型触发条件效果
FilterPushDownRBOFilter 在 Join 之上将 Filter 推到 Join 之下或叶子
ColumnPruningRBOProject 中不用所有列移除不需要的列引用
PartitionPrunerRBOFilter 中包含分区列裁剪不需要的 Partition
JoinReorderCBO多表 Join (n≥3)基于统计信息重排 Join 顺序
PushDownAggThroughJoinRBOAgg 在 Join 之上把部分聚合下推到 Join 之下
PushDownLimitThroughJoinRBOLimit N 在 Join 之上将 Limit 推到每个表扫描
EliminateOuterJoinRBO可证明 Outer 等价于 Inner消除不必要的 Outer Join
MaterializedViewRewriteCBO存在匹配的物化视图用物化视图扫描替代原始表扫描
SimplifyAggGroupBySetsRBOGROUP BY 可简化合并 GROUPING SETS
EliminateLimitZeroRBOLIMIT 0简化无结果查询

关键代码

// NereidsPlanner.java:207-210 — 完整优化流程
public void plan(StatementBase queryStmt) {
    // 1. analyze: Bind + Resolve + Rewrite
    // 2. rewrite: RBO 规则应用
    // 3. optimize: CBO + Implementation
    // 4. physicalPlan: 生成最终物理计划
}

// Memo.java — 核心抽象
// Memo 存储所有等价计划组, 每个 Group 包含多个 GroupExpression
// GroupExpression 的 children 是 Group ID (而非具体 Plan)
// 这允许组合爆炸式的等价计划探索

11. 常见问题 / 面试题(汇总)

Q: 为什么不用 Apache Calcite 作为优化器? A: (1) 性能: Nereids 为 Doris 量身优化, 比通用 Calcite 更轻更快; (2) 可控: 完全的代码控制权, 不依赖 Calcite 的升级节奏; (3) 分布式: Calcite 的 CBO 缺乏对 MPP 分布式计划的深层建模, Nereids 可以直接处理 Fragment/Exchange。

Q: RBO 和 CBO 有什么区别? A: RBO (Rule-Based Optimization) 是基于启发式规则的优化, 如"Filter 一律下推", 不考虑数据分布, 速度快但可能不最优。CBO (Cost-Based Optimization) 是基于统计信息的优化, 如"根据表的行数和 Join 列的 NDV 选择 Hash Join 还是 Broadcast Join", 需要维护统计信息但决策更准。Nereids 中的 CBO 通过 Cascades 框架和 Cost Model 实现。

Q: 优化器规则太多会不会反而影响查询性能(optimization overhead)? A: 会。Nereids 通过以下机制控制: (1) RuleSet 按查询类型启用不同规则集; (2) Rule 的 pattern() 方法快速过滤不匹配的 Plan; (3) Cost Improvement 阈值 (如新计划的 cost 必须比当前最优至少低 1% 才替换); (4) Explore 轮次限制。


5. 物理计划生成——Optimized Plan → Physical Plan

架构定位

逻辑计划(LogicalXxx)转换为物理计划(PhysicalXxx)。物理计划增加了分布式属性: 每个算子的数据分布方式(HASH/SHUFFLE/BROADCAST/NONE)和 Fragment 归属。

源码导航

文件关键类/方法职责
fe/fe-core/.../nereids/plans/physical/Physical* 类集合物理算子
fe/fe-core/.../nereids/properties/PhysicalProperties数据分布属性
fe/fe-core/.../planner/DistributedPlanner.javaDistributedPlanner生成分布式计划(Fragment 切分)

物理算子对应表

Logical PlanPhysical Plan分布式属性
LogicalOlapScanPhysicalOlapScan从 Tablet 分布决定 shuffle
LogicalFilterPhysicalFilter保留父算子分布
LogicalHashJoinPhysicalHashJoinBROADCAST(右表) 或 SHUFFLE(两表)
LogicalAggregatePhysicalHashAggregate两阶段: LocalAgg(Shuffle) + GlobalAgg
LogicalSortPhysicalTopNPhysicalQuickSortLocalTopN → MergeSort

11. 常见问题 / 面试题(汇总)

Q: 为什么需要 Partitioner 属性? A: 分布式计划的正确性要求满足数据分布约束。例如 HashJoin 要求两表在 Join Key 上的分布一致(或有 BROADCAST 特性)。物理计划中的 PhysicalProperties 描述了每个算子对数据的分布要求, 如果上游不满足则插入 Exchange(数据重分布)。


6. Plan Fragment 切分与调度

架构定位

Coordinator 将 Physical Plan 按 Exchange 边界切分为 Fragment(每个 Fragment 是一个独立的 Pipeline 执行单元)。Fragment Tree 结构: Root Fragment(聚合/排序) + Leaf Fragment(扫描/过滤)。

源码导航

文件关键类/方法职责
fe/fe-core/.../qe/Coordinator.javaexec(), deliverExecFragmentRequests()协调查询执行
fe/fe-core/.../planner/DistributedPlanner.javaplan()切分 Fragment
fe/fe-core/.../planner/PlanFragment.javaPlanFragment, PlanFragmentIdFragment 数据单元

调用链

Coordinator.exec()
  → 1. 调用 DistributedPlanner.plan() — 切分 Fragment
  → 2. 确定每个 Fragment 的 BE 分配
       (根据 Tablet 分布选择存储所在 BE 作为 Scan Fragment 宿主)
  → 3. 构造 TPipelineFragmentParams (Thrift 序列化)
  → 4. RPC: deliverExecFragmentRequests() →
       每个 BE 的 BackendService.exec_plan_fragment()
  → 5. 启动 Root Fragment (最后一个 Fragment 启动)
    → 6. 等待 ResultReceiver 收集结果

切分示例

SELECT COUNT(*) FROM orders
WHERE dt='2026-07'
GROUP BY status

Physical Plan:
    PhysicalHashAgg (Global)         ← Fragment 0 (Root, 1 BE 执行)
    └── Exchange (SHUFFLE by status)
    PhysicalHashAgg (Local)          ← Fragment 1 (中间, 按 hash 分布到多 BE)
    └── Exchange (GATHER)
    PhysicalOlapScan (orders)        ← Fragment 2 (Leaf, 每个 Tablet 所在 BE 执行)

Fragment 分配算法

  1. Leaf Fragment (Scan): 每个 Tablet 对其副本所在 BE 分配一个 Instance。如果 Tablet 有 3 个副本, 选择当前负载最低的 BE(根据 Heatbeat 上报的指标)
  2. Intermediate Fragment: 按数据量或 hash 分布到多个 BE
  3. Root Fragment: 协调者 FE 自身运行(简单的聚合)或分配给 1 个 BE(复杂聚合/排序)

11. 常见问题 / 面试题(汇总)

Q: 如果所有副本所在 BE 都挂了怎么办? A: Coordinator 会收到 RPC 失败, 记录到 query profile 中, 查询返回错误。不会自动重试(因为 Doris 查询无状态)。用户需要重新执行 SQL。


7. BE Pipeline 执行——Fragment Instance → Pipeline

架构定位

BE 端接收 Fragment 参数 → 构造 Pipeline(由 Operator 组成的 DAG) → 提交 PipelineTask 执行 → 线程池调度 Driver 的执行循环。

源码导航

文件关键类/方法职责
be/src/runtime/fragment_mgr.cppFragmentMgr::exec_plan_fragment()BE 端 Fragment 入口
be/src/exec/pipeline/pipeline_fragment_context.hPipelineFragmentContext::submit()构造 Pipeline 并提交执行
be/src/exec/pipeline/pipeline.hPipeline一个 Pipeline(多个 Operator 的顺序链)
be/src/exec/pipeline/pipeline_task.hPipelineTask::execute()单个 Task 的执行循环
be/src/exec/pipeline/operator.hOperator 接口 (实际在 pipeline 目录下)算子基类: Source/Sink/Transform
be/src/exec/pipeline/dependency.hDependencyPipeline 间的依赖关系

Pipeline 执行模型

Fragment Instance
   └── Pipeline DAG
        ├── Pipeline 0: OlapScanOp(Source) → FilterOp → AggSinkOp
        │     ↑ 依赖 (通过 Exchange)
        └── Pipeline 1: ExchangeSourceOp → AggOp → ResultSinkOp

核心执行模型

PipelineFragmentContext::submit()
  → 1. 从 Thrift 参数构建 Operator Tree
  → 2. 构建 Pipeline 依赖图 (Dependency)
  → 3. 构造 PipelineTask (每个 Pipeline 1个或多个)
  → 4. 推入线程池: Driver::execute()
       → while (!done) {
            task->execute(&done);
            // 一轮执行后 yield, 避免一个 Task 占用线程太久
            // yield 阈值: time_slice (默认 100ms)
         }

线程模型

┌──────────────────────────────────────────────┐
│          Thread Pool (多线程)                  │
│  ┌────┐  ┌────┐  ┌────┐  ┌────┐  ┌────┐   │
│  │ T0 │  │ T1 │  │ T2 │  │ T3 │  │ T4 │  │
│  └────┘  └────┘  └────┘  └────┘  └────┘   │
│    │       │       │       │       │        │
│    ▼       ▼       ▼       ▼       ▼        │
│ ┌──────────────────────────────────────┐    │
│ │ Task Queue (Ready Queue)              │    │
│ │ [Task A][Task B][Task C]...[Task Z]   │    │
│ └──────────────────────────────────────┘    │
│                                              │
│  ┌────────────────────────┐                 │
│ │ Blocked Queue            │                 │
│ │ [Task B: waiting dep]    │                 │
│ │  → 依赖满足后移入 Ready  │                 │
│ └────────────────────────┘                 │
└──────────────────────────────────────────────┘

关键代码

// pipeline_task.cpp:446 — PipelineTask::execute 核心循环
Status PipelineTask::execute(bool* done) {
    // 循环处理当前 Pipeline 中的所有算子
    while (!_is_finalized) {
        // 1. Source: operator->get_block(block)
        // 2. Transform: 下一个 operator 处理 block
        // 3. Sink: 最后一个 operator 输出 (写入 Exchange 或 Result)
        _root_operator->get_block(&block);
    }
}

11. 常见问题 / 面试题(汇总)

Q: 为什么 Pipeline 引擎比旧 exec_model 好? A: Pipeline 将查询拆分为更小的执行单元(PipelineTask), 多个查询的 Task 可以交织执行, 提高 CPU/IO 资源利用率。旧模型一个查询独占线程直到完成, 导致线程池利用率低和内存暴涨。


8. Scan 算子与存储读取——Pipeline → Disk

架构定位

OlapScanOperator 是连接 Pipeline 执行引擎和存储引擎的桥梁。它从 Tablet 读取数据, 应用索引过滤、应用投影(只读需要的列)、应用谓词下推, 返回向量化的 Block。

源码导航

文件关键类/方法职责
be/src/exec/operator/olap_scan_operator.cppOlapScanOperator::open(), get_block()Scan 算子主逻辑
be/src/exec/scan/olap_scanner.cppOlapScanner::get_block()实际的 Tablet 扫描实现
be/src/storage/tablet_reader.hTabletReader从 Tablet 读取 Rowset → 构造 RowBlock
be/src/storage/rowset/rowset_reader.hRowsetReader从 Rowset 中读取
be/src/storage/rowset/segment_reader.hSegmentReader读取单个 Segment 文件的列数据
be/src/storage/rowset/segment_v2/segment_iterator.cppSegmentIterator按条件迭代 Segment 中的行

调用链

OlapScanOperator::get_block(Block* block)
  → OlapScanner::get_block()
    → TabletReader::init()
      → RowsetReader::init()           // 选择要读的 Rowset (按 Version)
        → SegmentReader::init()        // 打开 Segment 文件
          → SegmentIterator::init()     // 构建迭代器
    → TabletReader::next_block()
      → SegmentIterator::next_batch()   // 读取一批行 (4096 默认)
        → [ZoneMap 过滤] → [Ordinal 定位] → [读取 Column Pages]
          → [BloomFilter 过滤] → [Bitmap 过滤]
            → [投影: 只读需要的列]
              → 返回 Block (列向量)

// 核心: 索引过滤发生在存储层, 减少从磁盘读取和向上传递的数据量

索引过滤层次

查询: SELECT status, SUM(amount) FROM t WHERE dt='2026-07' AND city='BJ'

1. PartitionPruner (FE Nereids)
   → 裁剪 Partition: 只保留 dt='2026-07' 的分区

2. KeyRange (FE → BE Thrift)
   → Tablet 范围: 只扫描 key 在 [BJ] 区间内的 Tablet

3. ZoneMap (BE Segment 级)
   → 每列 Min/Max 值: [dt_min, dt_max] 包含 '2026-07'?
   → 如果不包含, 跳过整个 Segment (Row Group)

4. Ordinal Index (BE Page 级)
   → 页级索引: Page 的 ordinal → 文件 offset 映射
   → 定位到包含数据的 DataPage

5. BloomFilter (BE, 可选)
   → city='BJ': BloomFilter 中 'BJ' 存在?
   → 不存在 → 跳过该 Rowset

6. Bitmap Index (BE, 可选)
   → 低基数列(city)的 Bitmap 快速定位行号

11. 常见问题 / 面试题(汇总)

Q: "列存"在这里的体现是什么? A: 物理上, 每一列的 Page 独立存储在一个 Segment 中。查询 SELECT status, amount 只需读这两个 Page, 不需要读取表中所有列(列投影)。这在宽表(100+ 列)上节省巨大的 I/O 开销。


9. 向量化表达式计算——Column → Expr → Compute

架构定位

BE 的向量化执行层把列数据从存储层拿到后, 以列向量为单位进行表达式计算(Filter, Agg, Join 等), 利用 SIMD 加速。所有计算函数都实现为对 IColumn 的批量操作。

源码导航

文件关键类/方法职责
be/src/vec/core/block.hBlock向量化数据块: 多列的集合
be/src/vec/columns/column.hIColumn列向量接口
be/src/vec/columns/column_vector.hColumnVector<T>具体类型列向量 (int/float/string)
be/src/vec/data_types/data_type.hIDataType数据类型接口
be/src/vec/exprs/vexpr.hVExpr, VExprContext向量化表达式
be/src/vec/functions/function.hIFunction函数定义接口
be/src/vec/functions/function_simple_unary.hFunctionSimpleUnary一元函数模板类
be/src/vec/aggregate_functions/aggregate_function.hIAggregateFunction聚合函数接口

核心类型系统

Block (数据块, 类比 Pandas DataFrame)
├── IColumn `dt`       (ColumnVector<Date>)   4096 rows
├── IColumn `status`   (ColumnString)          4096 rows
├── IColumn `amount`   (ColumnVector<Float64>) 4096 rows
└── IColumn `city`     (ColumnString)          4096 rows

所有列操作都是对整个 Column 的批量处理:
```cpp
// 向量化 Filter: amount > 100
// 传统做法(行式): for i in 0..4095: if amount[i] > 100: keep()
// 向量化做法:
ColumnVector<Float64>::apply_filter(
    ColumnVector<Float64>::compare(GREATER_THAN, Literal(100.0))
);
// 上面的 compare() 可以使用 SIMD 指令 (AVX2): _mm256_cmp_pd

新增函数的完整示例

abs(double) → double 为例:

// 1. func_tinyint.cpp — 注册函数
class FunctionAbs : public IFunction {
    // 实现 execute_impl()
    // 输入: Block (一列 double)
    // 输出: ColumnVector<Float64> (每行的 abs 结果)
};

11. 常见问题 / 面试题(汇总)

Q: SIMD 具体在哪些操作中生效? A: 所有数值类型的比较、算术、逻辑操作, 以及部分聚合(如 SUM 的 accumulation)。在 ColumnVector<T>::compare() 等函数中, 编译器或手写的 SIMD intrinsic 批量处理 4/8/16 个元素。判断方法: grep _mm256/_mm512 等在 be/src/vec/ 下。


10. 结果聚合与返回——BE → FE → Client

架构定位

最后一步: BE 将计算结果通过 Brpc 发送回 Coordinator; Coordinator 合并结果; 通过 MySQL 协议返回给客户端。

源码导航

文件关键类/方法职责
be/src/exec/operator/exchange_sink_operator.hExchangeSinkOperatorBE 侧数据输出到 Exchange
be/src/exec/operator/exchange_source_operator.hExchangeSourceOperatorBE 侧从 Exchange 接收数据
be/src/runtime/data_stream_sender.hDataStreamSender通过 Brpc 发送数据
fe/fe-core/.../qe/ResultReceiver.javaResultReceiverFE 侧接收 BE 结果
fe/fe-core/.../qe/StmtExecutor.javasendResultSet()将结果转为 MySQL Row

调用链

BE: ExchangeSinkOperator
  → DataStreamSender::send()
    → Brpc Channel → FE Coordinator
      → ResultReceiver.getNext()
        → RowBatch (Thrift 反序列化)
          → StmtExecutor.sendResultSet()
            → MysqlChannel.writeRow() ← MySQL 协议返回客户端

11. 常见问题 / 面试题(汇总)

Q: 为什么 Ring Buffer 不是纯内存的? A: DataStreamSender 有流控机制, 当接收端(FE 或其他 BE)消费慢时, 发送缓冲区(TBuf)满会触发 backpressure, 导致上游 Pipeline Task 暂停(yield), 直到缓冲区有空间。这防止了 OOM。


SQL 查询端到端时序图

下一步