Apache Doris 底层原理深度学习笔记(持续更新)
·
一、架构原理深度解析
1.1 整体架构设计
┌─────────────────────────────────────────────────────────┐
│ MySQL Client │
└─────────────────────────────────────────────────────────┘
│
MySQL Protocol
│
┌─────────────────────────────────────────────────────────┐
│ Frontend (FE) Layer │
│ ┌──────────────────────────────────────────────────┐ │
│ │ Master FE (BDBJE - Metadata Storage) │ │
│ │ - Catalog Manager (Database/Table Metadata) │ │
│ │ - Load Manager (Import Job Management) │ │
│ │ - Routine Load Manager │ │
│ │ - Alter Job Manager │ │
│ └──────────────────────────────────────────────────┘ │
│ ┌──────────────────────────────────────────────────┐ │
│ │ Query Engine │ │
│ │ - SQL Parser (JavaCC) │ │
│ │ │ └─> AST (Abstract Syntax Tree) │ │
│ │ - Analyzer (Semantic Analysis) │ │
│ │ │ └─> Analyzed Statement │ │
│ │ - Planner │ │
│ │ │ ├─> SingleNodePlanner │ │
│ │ │ └─> DistributedPlanner │ │
│ │ │ └─> PlanFragment (DataSink/DataStream) │ │
│ │ - Coordinator (Query Execution Coordination) │ │
│ │ └─> Fragment Instances Distribution │ │
│ └──────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────┘
│
Thrift RPC (PlanFragmentExec)
│
┌─────────────────────────────────────────────────────────┐
│ Backend (BE) Layer │
│ ┌──────────────────────────────────────────────────┐ │
│ │ Fragment Executor │ │
│ │ - ExecNode Tree (Pipeline Execution) │ │
│ │ ├─> OlapScanNode │ │
│ │ ├─> AggregationNode │ │
│ │ ├─> HashJoinNode │ │
│ │ └─> ExchangeNode │ │
│ └──────────────────────────────────────────────────┘ │
│ ┌──────────────────────────────────────────────────┐ │
│ │ Storage Engine │ │
│ │ - Tablet Management │ │
│ │ - Rowset (Data Version Management) │ │
│ │ - Segment (Column Data Storage) │ │
│ │ ├─> Column Data (Encoded & Compressed) │ │
│ │ ├─> Index (ShortKey/ZoneMap/BloomFilter) │ │
│ │ └─> Footer (Metadata) │ │
│ │ - MemTable (Write Buffer) │ │
│ │ - Compaction Engine │ │
│ └──────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────┘
│
Local File System / HDFS
1.2 FE 元数据管理原理
BDBJE (Berkeley DB Java Edition) 实现高可用
// FE 元数据存储结构
// 基于 BDBJE 实现的分布式状态机复制
// 元数据组织结构
MetaContext {
Database database;
Map<Long, Table> idToTable;
Map<String, Database> nameToDatabase;
Map<Long, Backend> idToBackend;
Map<String, Resource> nameToResource;
}
// Journal 机制:所有元数据变更都通过 Journal 持久化
// Journal 类型示例
enum OperationType {
OP_CREATE_DB, // 创建数据库
OP_DROP_DB, // 删除数据库
OP_CREATE_TABLE, // 创建表
OP_ADD_PARTITION, // 添加分区
OP_FINISH_CONSISTENCY_CHECK, // 一致性检查完成
OP_BACKEND_STATE_CHANGE, // BE 状态变更
...
}
// 元数据写入流程
1. Master FE 接收 DDL 请求
2. 修改内存中的元数据
3. 生成 Journal 日志
4. 通过 BDBJE 复制到 Follower FE
5. 大多数节点写入成功后返回
6. Follower/Observer 异步回放 Journal
元数据检查点 (Image)
# FE 会定期生成元数据快照
# 目录结构: doris-meta/image/
image.123456 # 快照文件,数字是 Journal ID
image.version # 版本信息
# 恢复流程
1. 加载最新的 image 文件
2. 回放 image 之后的所有 Journal
3. 重建完整的元数据状态
1.3 BE 存储引擎深度解析
存储层次结构
Database
└─> Table
└─> Partition (Range/List 分区)
└─> MaterializedIndex (Base Table + Rollup)
└─> Tablet (分桶单位)
└─> Replica (多副本)
└─> Rowset (版本管理单位)
└─> Segment (物理存储文件)
Tablet 结构详解
// Tablet 是 Doris 数据管理的基本单位
class Tablet {
TabletId tablet_id;
SchemaHash schema_hash;
TabletState state;
// 版本管理
vector<Rowset> rowsets; // 所有历史版本
Version committed_version; // 已提交的最大版本
// Rowset 结构
class Rowset {
Version version; // [start_version, end_version]
int num_segments;
vector<Segment> segments;
// 每个 Segment 对应一个 .dat 文件
class Segment {
// 列式存储,每列独立编码压缩
vector<ColumnData> columns;
IndexData index_data;
SegmentFooter footer;
}
}
}
数据文件格式深度剖析
Segment File (.dat) 物理结构:
┌───────────────────────────────────────────────────────┐
│ Magic Number (8 bytes) │
├───────────────────────────────────────────────────────┤
│ Column 0 Data Block │
│ ┌─────────────────────────────────────────────────┐ │
│ │ Data Page 0 (Compressed) │ │
│ │ - Encoding: Dictionary/RLE/BitPacking │ │
│ │ - Compression: LZ4/Zstd/Snappy │ │
│ ├─────────────────────────────────────────────────┤ │
│ │ Data Page 1 (Compressed) │ │
│ ├─────────────────────────────────────────────────┤ │
│ │ ... │ │
│ └─────────────────────────────────────────────────┘ │
├───────────────────────────────────────────────────────┤
│ Column 1 Data Block │
├───────────────────────────────────────────────────────┤
│ ... │
├───────────────────────────────────────────────────────┤
│ ShortKey Index │
│ - 每 1024 行建立一个索引项 │
│ - 存储: (Key, RowId, FileOffset) │
├───────────────────────────────────────────────────────┤
│ ZoneMap Index (Min/Max Statistics) │
│ - 每个 Data Page 的统计信息 │
│ - 用于快速过滤不满足条件的 Page │
├───────────────────────────────────────────────────────┤
│ Bloom Filter Index │
│ - 配置的列生成 Bloom Filter │
│ - 快速判断值是否存在 │
├───────────────────────────────────────────────────────┤
│ Bitmap Index (可选) │
│ - 低基数列的位图索引 │
├───────────────────────────────────────────────────────┤
│ Footer (Metadata) │
│ - Schema 信息 │
│ - 各 Column 的 Offset 和 Length │
│ - 索引的 Offset │
│ - 统计信息 (行数、大小等) │
├───────────────────────────────────────────────────────┤
│ Footer Length (4 bytes) │
├───────────────────────────────────────────────────────┤
│ Checksum (4 bytes) │
├───────────────────────────────────────────────────────┤
│ Magic Number (8 bytes) │
└───────────────────────────────────────────────────────┘
编码方式详解
// 1. Dictionary Encoding (字典编码)
// 适用场景: 重复值多的字符串列
// 原理: 构建字典,用整数 ID 替代实际值
Column: ["Beijing", "Shanghai", "Beijing", "Shenzhen", "Beijing"]
Dictionary: {0: "Beijing", 1: "Shanghai", 2: "Shenzhen"}
Encoded: [0, 1, 0, 2, 0]
Compression Ratio: ~70%
// 2. Run-Length Encoding (RLE)
// 适用场景: 连续重复值多的列
// 原理: 记录值和重复次数
Original: [1, 1, 1, 1, 2, 2, 3, 3, 3]
Encoded: [(1, 4), (2, 2), (3, 3)]
// 3. Bit-Packing
// 适用场景: 小整数列
// 原理: 使用实际需要的最小位数存储
Values: [5, 7, 3, 6] // 最大值 7,需要 3 bits
Packed: 101 111 011 110 (12 bits instead of 32 bits)
// 4. Delta Encoding (差值编码)
// 适用场景: 递增序列
Original: [100, 102, 105, 108, 110]
Delta: [100, +2, +3, +3, +2]
// 5. Frame-of-Reference (FOR)
// 适用场景: 数值范围集中的列
Original: [1000, 1002, 1005, 1008]
Base: 1000
Encoded: [0, 2, 5, 8]
二、查询执行原理
2.1 查询编译流程
SQL 执行全流程
SQL Query
↓
[1] Parser (JavaCC)
↓
Abstract Syntax Tree (AST)
↓
[2] Analyzer
- 语义分析
- 元数据绑定 (Table/Column Resolution)
- 类型检查
- 权限验证
↓
Analyzed Statement
↓
[3] Rewriter
- 子查询展开
- 表达式简化
- 常量折叠
↓
[4] Planner
- 单机逻辑计划
- 分布式逻辑计划
- 物理计划
↓
PlanFragment Tree
↓
[5] Coordinator
- Fragment 分发到各个 BE
- 执行进度跟踪
- 结果收集
↓
[6] BE Execution
- Pipeline 执行
- 向量化计算
↓
Query Result
示例: 复杂查询的执行计划
-- 原始查询
SELECT
o.city,
COUNT(DISTINCT o.user_id) as users,
SUM(o.amount) as total_amount
FROM orders o
JOIN users u ON o.user_id = u.user_id
WHERE o.order_date >= '2024-01-01'
AND u.age >= 18
GROUP BY o.city
HAVING SUM(o.amount) > 10000
ORDER BY total_amount DESC
LIMIT 10;
-- 执行计划分析
EXPLAIN VERBOSE
SELECT ...;
生成的执行计划结构
Fragment 0 (Coordinator)
├─> RESULT SINK
└─> TOP-N (LIMIT 10)
└─> EXCHANGE (GATHER from Fragment 1)
Fragment 1 (BE Nodes)
├─> DATA SINK (Stream to Fragment 0)
└─> TOP-N (Local LIMIT 10)
└─> HAVING FILTER (SUM(amount) > 10000)
└─> AGGREGATION (GROUP BY city)
├─> Aggregate Functions: COUNT(DISTINCT user_id), SUM(amount)
└─> HASH JOIN (INNER)
├─> LEFT CHILD: EXCHANGE (SHUFFLE by user_id from Fragment 2)
└─> RIGHT CHILD: EXCHANGE (SHUFFLE by user_id from Fragment 3)
Fragment 2 (BE Nodes - Orders Table Scan)
├─> DATA SINK (Shuffle by user_id to Fragment 1)
└─> OLAP SCAN (orders)
├─> Predicates: order_date >= '2024-01-01'
├─> Partition Pruning: Applied
├─> ShortKey Index: Applied
├─> ZoneMap Index: Applied
└─> Columns: user_id, city, amount
Fragment 3 (BE Nodes - Users Table Scan)
├─> DATA SINK (Shuffle by user_id to Fragment 1)
└─> OLAP SCAN (users)
├─> Predicates: age >= 18
└─> Columns: user_id
2.2 数据扫描优化
OlapScanNode 执行细节
// OlapScanNode 是 BE 读取数据的核心组件
class OlapScanNode {
// 1. Tablet 选择和分配
vector<Tablet> tablets; // 根据分桶和分区计算
// 2. 谓词下推 (Predicate Pushdown)
vector<Predicate> pushed_predicates;
// 3. 列裁剪 (Column Pruning)
vector<int> required_column_ids;
// 4. 索引选择
bool use_short_key_index = true;
bool use_zone_map_index = true;
bool use_bloom_filter = true;
// 5. 读取优化
int chunk_size = 4096; // 向量化批处理大小
bool use_vectorized_reader = true;
}
// Tablet 扫描流程
1. 根据分区裁剪确定需要扫描的 Partition
2. 根据分桶计算确定需要扫描的 Tablet
3. 选择合适的 Rowset Version (MVCC)
4. 在每个 Segment 中:
a. 应用 ShortKey Index 定位数据范围
b. 应用 ZoneMap Index 过滤 Page
c. 应用 Bloom Filter 快速否定查询
d. 读取满足条件的 Data Page
e. 解压缩和解码
f. 应用剩余谓词过滤
5. 返回符合条件的数据批次
索引使用示例
-- 创建表时的索引配置
CREATE TABLE sales_data
(
order_id BIGINT,
user_id BIGINT,
product_id BIGINT,
city VARCHAR(50),
order_date DATE,
amount DECIMAL(10,2)
)
DUPLICATE KEY(order_id, user_id) -- ShortKey Index 基于这些列
DISTRIBUTED BY HASH(order_id) BUCKETS 32
PROPERTIES (
"bloom_filter_columns" = "product_id,city", -- Bloom Filter
"storage_medium" = "SSD"
);
-- 创建 Bitmap 索引
CREATE INDEX idx_city ON sales_data(city) USING BITMAP;
-- 查询时索引使用分析
SELECT * FROM sales_data
WHERE order_id = 12345 -- ShortKey Index: 精确定位
AND city = 'Beijing' -- Bloom Filter: 快速过滤
AND order_date >= '2024-01-01' -- ZoneMap: 跳过不符合的 Page
AND product_id IN (1,2,3); -- Bloom Filter: 快速判断存在性
-- 查看索引使用情况
EXPLAIN SELECT * FROM sales_data WHERE city = 'Beijing';
-- 输出会显示:
-- OlapScanNode
-- Predicates: city = 'Beijing'
-- Indexes: SHORT_KEY_INDEX, BLOOM_FILTER, ZONE_MAP
2.3 向量化执行引擎
传统 Row-based vs 向量化 Column-based
// 传统行式执行 (Tuple-at-a-time)
for (each row) {
result = filter(row);
if (result) {
aggregate(row);
}
}
// 问题:
// - 大量函数调用开销
// - CPU 分支预测失败
// - 缓存命中率低
// 向量化执行 (Batch-at-a-time)
const int BATCH_SIZE = 4096;
while (has_data) {
Column* col_batch = read_column(BATCH_SIZE);
// SIMD 指令并行处理
bool* filter_result = filter_vectorized(col_batch);
// 批量聚合
aggregate_vectorized(col_batch, filter_result);
}
// 优势:
// - 减少函数调用
// - 利用 CPU SIMD 指令
// - 提高缓存命中率
// - 提升 10-100 倍性能
向量化算子示例
// 向量化过滤算子
class VectorizedFilterOperator {
// 使用 AVX2 指令集进行 SIMD 加速
void filter(
const int64_t* column_data,
const int64_t threshold,
bool* result,
int batch_size
) {
__m256i threshold_vec = _mm256_set1_epi64x(threshold);
for (int i = 0; i < batch_size; i += 4) {
// 一次加载 4 个 int64 值
__m256i data_vec = _mm256_loadu_si256(
(__m256i*)&column_data[i]
);
// SIMD 比较: 4 个值并行比较
__m256i cmp_result = _mm256_cmpgt_epi64(
data_vec, threshold_vec
);
// 存储结果
_mm256_storeu_si256(
(__m256i*)&result[i], cmp_result
);
}
}
}
// 向量化聚合算子
class VectorizedAggregationOperator {
// Hash Aggregation 向量化
void aggregate_sum(
const int64_t* keys,
const int64_t* values,
HashMap& hash_table,
int batch_size
) {
// 批量哈希计算
uint64_t* hash_values = prefetch_hash(keys, batch_size);
// 批量更新哈希表
for (int i = 0; i < batch_size; i += 8) {
// Prefetch 下一批数据
_mm_prefetch(&keys[i + 8], _MM_HINT_T0);
// 更新聚合值
for (int j = 0; j < 8; j++) {
hash_table[keys[i+j]] += values[i+j];
}
}
}
}
三、数据导入原理
3.1 MemTable 和 Flush 机制
写入流程详解
// BE 端数据写入的核心数据结构
class MemTable {
// 内存中的写缓冲区,按 Key 排序
SkipList<Row> data; // 跳表实现,O(log n) 插入
size_t memory_usage;
size_t max_memory_limit = 100MB; // 默认 100MB
// 插入数据
void insert(const Row& row) {
data.insert(row);
memory_usage += row.size();
// 达到阈值触发 Flush
if (memory_usage >= max_memory_limit) {
flush_to_disk();
}
}
// Flush 到磁盘
void flush_to_disk() {
// 1. 对 MemTable 数据排序(已通过 SkipList 排序)
// 2. 按列组织数据
// 3. 编码和压缩
// 4. 生成 Segment 文件
// 5. 写入 Rowset 元数据
Segment* seg = new Segment();
// 分列编码
for (each column) {
ColumnWriter writer(column_schema);
// 选择最优编码方式
Encoding encoding = choose_encoding(column_data);
// 写入数据页
for (each page in column_data) {
encode_and_compress(page, encoding);
writer.write(page);
}
// 写入索引
writer.write_index();
}
// 写入 Footer
seg->write_footer();
// 生成 Rowset
Rowset* rowset = create_rowset(seg);
// 更新 Tablet 元数据
tablet->add_rowset(rowset);
}
}
Stream Load 完整流程
Client
│ HTTP PUT /api/db/table/_stream_load
↓
FE Coordinator
│ 1. 解析 HTTP 请求
│ 2. 创建 Load Job (分配 Label 和 TxnId)
│ 3. 选择 BE 节点作为 Coordinator
↓
BE Coordinator
│ 4. 解析数据格式 (CSV/JSON)
│ 5. 数据分发到各个 BE
│ - 根据 Distribution Key Hash 计算目标 BE
│ - 通过 Thrift RPC 发送数据
↓
BE Executors (多个)
│ 6. 接收数据写入 MemTable
│ 7. MemTable 满时 Flush 到 Segment
│ 8. 向 Coordinator 报告写入状态
↓
BE Coordinator
│ 9. 收集所有 BE 的写入结果
│ 10. 向 FE 报告 Load 完成
↓
FE Master
│ 11. 提交事务 (Publish Version)
│ 12. 更新元数据
│ 13. 返回结果给 Client
↓
Client receives response
Stream Load 实战示例
# 创建测试表
mysql -h 127.0.0.1 -P 9030 -u root << EOF
CREATE DATABASE IF NOT EXISTS test_db;
USE test_db;
CREATE TABLE user_events (
event_id BIGINT,
user_id BIGINT,
event_type VARCHAR(50),
event_time DATETIME,
properties JSON
)
DUPLICATE KEY(event_id)
DISTRIBUTED BY HASH(event_id) BUCKETS 32;
EOF
# 准备大量数据 (100万行)
python3 << 'PYTHON'
import json
import random
from datetime import datetime, timedelta
with open('events.json', 'w') as f:
base_time = datetime(2024, 1, 1)
for i in range(1000000):
event = {
"event_id": i,
"user_id": random.randint(1, 100000),
"event_type": random.choice(["click", "view", "purchase"]),
"event_time": (base_time + timedelta(seconds=i)).strftime("%Y-%m-%d %H:%M:%S"),
"properties": json.dumps({"page": f"page_{random.randint(1,100)}"})
}
f.write(json.dumps(event) + '\n')
PYTHON
# Stream Load 导入
curl --location-trusted -u root: \
-H "label:load_events_$(date +%s)" \
-H "format:json" \
-H "max_filter_ratio:0.1" \
-H "timeout:300" \
-H "strict_mode:false" \
-T events.json \
http://localhost:8030/api/test_db/user_events/_stream_load
# 查看导入结果
mysql -h 127.0.0.1 -P 9030 -u root -e "
USE test_db;
SHOW LOAD ORDER BY CreateTime DESC LIMIT 1\G
SELECT COUNT(*) FROM user_events;
"
3.2 事务和版本管理 (MVCC)
Doris 的 MVCC 实现
// Tablet 中的版本管理
class Tablet {
// 每个 Tablet 维护多个版本的数据
map<Version, Rowset*> version_map;
// Version 定义: [start, end]
struct Version {
int64_t start;
int64_t end;
bool contains(int64_t v) {
return v >= start && v <= end;
}
}
// 示例: Tablet 的版本演进
// Initial: [0-0] (空)
// Load 1: [0-0], [1-1]
// Load 2: [0-0], [1-1], [2-2]
// Compact: [0-0], [1-2] (合并 [1-1] 和 [2-2])
// Load 3: [0-0], [1-2], [3-3]
}
// 查询时读取版本
class OlapScanner {
vector<Rowset*> select_rowsets(int64_t read_version) {
vector<Rowset*> result;
// 选择所有覆盖 read_version 的 Rowset
for (auto& [ver, rowset] : tablet->version_map) {
if (ver.contains(read_version)) {
result.push_back(rowset);
}
}
return result;
}
// 读取时需要对多个 Rowset 做归并排序
Row read_next() {
// Merge 多个 Rowset 的数据
// 根据 Key 排序,保证读取顺序
return merge_rowsets(rowsets);
}
}
事务提交流程
1. Begin Transaction
FE: txn_id = begin_txn(label)
2. Load Data (多个 BE 并行)
BE1: write_data(txn_id, tablet_1, data) -> version [10-10]
BE2: write_data(txn_id, tablet_2, data) -> version [10-10]
BE3: write_data(txn_id, tablet_3, data) -> version [10-10]
此时数据已写入,但版本未发布 (不可见)
3. Commit Transaction
Coordinator -> FE: commit_txn(txn_id)
FE:
- 检查所有 Tablet 写入成功
- 生成全局唯一的 version: 10
- 标记事务为 COMMITTED
4. Publish Version (异步)
FE -> All BE: publish_version(txn_id, version=10)
BE:
- 更新 Tablet 元数据
- 将 version [10-10] 标记为可见
- 返回 publish 成功
5. Transaction Complete
FE:
- 标记事务为 VISIBLE
- 清理事务状态
此时查询可以读取 version <= 10 的所
更多推荐
所有评论(0)