一、架构原理深度解析

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 的所

更多推荐