预计阅读时间: 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.java | MysqlServer, accept() | 启动 MySQL 服务, 绑定端口(默认 9030) |
fe/fe-core/.../mysql/nio/NMysqlServer.java | NMysqlServer | NIO Accept 循环 (Java NIO Selector) |
fe/fe-core/.../mysql/MysqlChannel.java | MysqlChannel | 单连接读写, MySQL Packet 编解码 |
fe/fe-core/.../qe/ConnectProcessor.java | processOnce() | 读一个 packet, 分发到 handleQuery/handleStmtPrepare 等 |
fe/fe-core/.../qe/ConnectContext.java | ConnectContext | 连接上下文: 当前用户、数据库、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 执行管线
}
时序图

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.g4 | ANTLR Lexer 规则 | Token 定义 |
gensrc/antlr4/org/apache/doris/nereids/DorisParser.g4 | ANTLR Parser 规则 | SQL 语法规则 |
fe/fe-core/.../nereids/parser/NereidsParser.java | parseSingle(), 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.java | Analyzer | 分析入口, 驱动 Scope 和 Binder |
fe/fe-core/.../nereids/analyzer/Scope.java | Scope | 作用域: 列名 → 表的映射 |
fe/fe-core/.../nereids/analyzer/CatalogBinder.java | CatalogBinder | 将表名绑定到 Catalog 中的 OlapTable |
fe/fe-core/.../nereids/analyzer/RelationManager.java | RelationManager | 管理分析过程的表/子查询关系 |
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();
数据流转图

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.java | Optimizer | 优化流程编排: explore → optimize |
fe/fe-core/.../nereids/memo/Memo.java | Memo | 存储所有等价计划 Group |
fe/fe-core/.../nereids/memo/Group.java | Group | 一个等价组(语义等价的多个 Plan) |
fe/fe-core/.../nereids/memo/GroupExpression.java | GroupExpression | Group 中的一个 Plan, 即一个算子树节点 |
fe/fe-core/.../nereids/rules/Rule.java | Rule | 规则接口 → 子类包括 RBO 和 CBO |
fe/fe-core/.../nereids/rules/RuleType.java | RuleType | 规则类型枚举(100+ 项) |
fe/fe-core/.../nereids/cost/CostCalculator.java | CostCalculator | 计算计划的 CPU/IO/Network/内存成本 |
fe/fe-core/.../nereids/stats/Statistics.java | Statistics | 列统计信息: 行数/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 清单(部分重点)
| 规则 | 类型 | 触发条件 | 效果 |
|---|---|---|---|
FilterPushDown | RBO | Filter 在 Join 之上 | 将 Filter 推到 Join 之下或叶子 |
ColumnPruning | RBO | Project 中不用所有列 | 移除不需要的列引用 |
PartitionPruner | RBO | Filter 中包含分区列 | 裁剪不需要的 Partition |
JoinReorder | CBO | 多表 Join (n≥3) | 基于统计信息重排 Join 顺序 |
PushDownAggThroughJoin | RBO | Agg 在 Join 之上 | 把部分聚合下推到 Join 之下 |
PushDownLimitThroughJoin | RBO | Limit N 在 Join 之上 | 将 Limit 推到每个表扫描 |
EliminateOuterJoin | RBO | 可证明 Outer 等价于 Inner | 消除不必要的 Outer Join |
MaterializedViewRewrite | CBO | 存在匹配的物化视图 | 用物化视图扫描替代原始表扫描 |
SimplifyAggGroupBySets | RBO | GROUP BY 可简化 | 合并 GROUPING SETS |
EliminateLimitZero | RBO | LIMIT 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.java | DistributedPlanner | 生成分布式计划(Fragment 切分) |
物理算子对应表
| Logical Plan | Physical Plan | 分布式属性 |
|---|---|---|
LogicalOlapScan | PhysicalOlapScan | 从 Tablet 分布决定 shuffle |
LogicalFilter | PhysicalFilter | 保留父算子分布 |
LogicalHashJoin | PhysicalHashJoin | BROADCAST(右表) 或 SHUFFLE(两表) |
LogicalAggregate | PhysicalHashAggregate | 两阶段: LocalAgg(Shuffle) + GlobalAgg |
LogicalSort | PhysicalTopN 或 PhysicalQuickSort | LocalTopN → 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.java | exec(), deliverExecFragmentRequests() | 协调查询执行 |
fe/fe-core/.../planner/DistributedPlanner.java | plan() | 切分 Fragment |
fe/fe-core/.../planner/PlanFragment.java | PlanFragment, PlanFragmentId | Fragment 数据单元 |
调用链
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 分配算法
- Leaf Fragment (Scan): 每个 Tablet 对其副本所在 BE 分配一个 Instance。如果 Tablet 有 3 个副本, 选择当前负载最低的 BE(根据 Heatbeat 上报的指标)
- Intermediate Fragment: 按数据量或 hash 分布到多个 BE
- 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.cpp | FragmentMgr::exec_plan_fragment() | BE 端 Fragment 入口 |
be/src/exec/pipeline/pipeline_fragment_context.h | PipelineFragmentContext::submit() | 构造 Pipeline 并提交执行 |
be/src/exec/pipeline/pipeline.h | Pipeline | 一个 Pipeline(多个 Operator 的顺序链) |
be/src/exec/pipeline/pipeline_task.h | PipelineTask::execute() | 单个 Task 的执行循环 |
be/src/exec/pipeline/operator.h | Operator 接口 (实际在 pipeline 目录下) | 算子基类: Source/Sink/Transform |
be/src/exec/pipeline/dependency.h | Dependency | Pipeline 间的依赖关系 |
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.cpp | OlapScanOperator::open(), get_block() | Scan 算子主逻辑 |
be/src/exec/scan/olap_scanner.cpp | OlapScanner::get_block() | 实际的 Tablet 扫描实现 |
be/src/storage/tablet_reader.h | TabletReader | 从 Tablet 读取 Rowset → 构造 RowBlock |
be/src/storage/rowset/rowset_reader.h | RowsetReader | 从 Rowset 中读取 |
be/src/storage/rowset/segment_reader.h | SegmentReader | 读取单个 Segment 文件的列数据 |
be/src/storage/rowset/segment_v2/segment_iterator.cpp | SegmentIterator | 按条件迭代 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.h | Block | 向量化数据块: 多列的集合 |
be/src/vec/columns/column.h | IColumn | 列向量接口 |
be/src/vec/columns/column_vector.h | ColumnVector<T> | 具体类型列向量 (int/float/string) |
be/src/vec/data_types/data_type.h | IDataType | 数据类型接口 |
be/src/vec/exprs/vexpr.h | VExpr, VExprContext | 向量化表达式 |
be/src/vec/functions/function.h | IFunction | 函数定义接口 |
be/src/vec/functions/function_simple_unary.h | FunctionSimpleUnary | 一元函数模板类 |
be/src/vec/aggregate_functions/aggregate_function.h | IAggregateFunction | 聚合函数接口 |
核心类型系统
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.h | ExchangeSinkOperator | BE 侧数据输出到 Exchange |
be/src/exec/operator/exchange_source_operator.h | ExchangeSourceOperator | BE 侧从 Exchange 接收数据 |
be/src/runtime/data_stream_sender.h | DataStreamSender | 通过 Brpc 发送数据 |
fe/fe-core/.../qe/ResultReceiver.java | ResultReceiver | FE 侧接收 BE 结果 |
fe/fe-core/.../qe/StmtExecutor.java | sendResultSet() | 将结果转为 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。

下一步
- 理解写入: doc-02-write-lifecycle.md——Stream Load 写入端到端
- 深入优化器: doc-03-nereids-optimizer.md——Cascades 框架与 Rule 深度
- 深入执行引擎: doc-04-pipeline-execution.md——Pipeline 详解