Doris architecture

ANALYTICAL DATABASE / SOURCE READING / LESSON 02

Write lifecycle

Follow Doris write paths through transactions, tablets, rowsets, and segment generation.

Reading
90-120 min
Track
Doris architecture
Source
Chinese source notes

The source notes for this track are currently maintained in Chinese.

预计阅读时间: 90-120 分钟 前置阅读: doc-00(架构概念), doc-01(理解 FE 协调和 BE 执行基底) 下一次阅读: doc-06(Rowset/Segment 存储), doc-07(Compaction), doc-10(导入体系全览)


1. 写入方式总览与对比

Doris 支持 6 种数据写入方式, 每种针对不同的数据源和延迟要求。

方式协议数据源延迟适用规模关键类
Stream LoadHTTP (multipart)本地文件/程序输出秒级单次 <10GBStreamLoadHandler / StreamLoadAction
Broker LoadSQL 提交HDFS / S3 上的大文件分钟级GB~TB 级BrokerLoadJob
Routine LoadSQL 提交Kafka (持续消费)秒级持续流式RoutineLoadManager / RoutineLoadTask
INSERT INTOMySQL SQL客户端 SQL秒级小批量InsertStmtExecutor / GroupCommitManager
Group Commit自动(INSERT 聚合)多客户端 INSERT 合并毫秒级高频小批量GroupCommitManager
Binlog LoadCanal/Binlog 订阅MySQL Binlog (CDC)秒级CDC 持续BinlogManager

架构对比图:

数据写入端到端——从 HTTP 到 Rowset 的全源码级穿越 图 01


2. Stream Load 全流程——作为写入主线

Stream Load 是理解 Doris 写入机制的最佳入口: 它覆盖了从 FE 协调到 BE 写入到 Publish 的完整路径, 且足够简单(单次 HTTP 请求)。以下逐段展开。

2a. FE 协调阶段——TxnManager 和 Tablet 分配

架构定位

客户端发送 HTTP PUT 请求到 FE 的 HTTP Server, 包含 CSV/JSON 数据 body。FE 负责: (1) 鉴权和表解析, (2) 开启事务, (3) 分配 Tablet(决定每个 Tablet 写到哪个 BE), (4) 转发请求到分配的 BE。

源码导航

文件关键类/方法职责
fe/fe-core/.../httpv2/rest/LoadAction.javaLoadAction.execute()HTTP 请求入口, 处理 multipart body
fe/fe-core/.../load/StreamLoadHandler.javaStreamLoadHandler.handleStreamLoad() (via TStreamLoadPutResult)Stream Load 核心协调逻辑
fe/fe-core/.../transaction/GlobalTransactionMgr.javabeginTransaction(), commitTransaction()全局事务管理
fe/fe-core/.../transaction/DatabaseTransactionMgr.javacommitTransaction()单 DB 事务管理
fe/fe-core/.../planner/StreamLoadPlanner.javaStreamLoadPlanner.plan()生成写入计划

调用链

Client: curl -X PUT -T data.csv http://fe:8030/api/db/table/_stream_load

FE HTTP Server (端口 8030)
  → LoadAction.execute(HttpRequest, HttpResponse)
    → StreamLoadHandler.executeLoad(TStreamLoadPutRequest)
      → 1. 鉴权: ConnectContext.get().getCurrentUserIdentity()
      → 2. 解析表: Database.getTable(OlapTable)
      → 3. beginTransaction():
           GlobalTransactionMgr.beginTransaction(dbId, tableId, label)
           → Txn 状态: PREPARE
      → 4. 分配 Tablet (按 Hash/Bucket):
           为每个 Tablet 选择一个 BE (基于负载/heartbeat 指标)
      → 5. 构造 TStreamLoadPutResult → 返回给客户端
           (告诉客户端: 数据应该 POST 到哪个 BE 的哪个端口)

Client: POST 数据到返回的 BE 地址

事务状态机

PREPARE → [beginTxn] → PREPARED
  → [数据写入成功] → COMMITTED
  → [写入失败/Txn超时] → ABORTED
  → [Publish 成功] → VISIBLE

关键代码

// StreamLoadHandler.java — 核心协调逻辑
public void executeLoad(TStreamLoadPutRequest request) {
    // 1. 开启事务
    long txnId = GlobalTransactionMgr.beginTransaction(dbId, tableId, label);
    // 2. 解析写入计划 (决定每个 Tablet 写到哪个 BE)
    StreamLoadPlanner planner = new StreamLoadPlanner(db, olapTable, request);
    TExecPlanFragmentParams plan = planner.plan();
    // 3. 返回给客户端: BE 地址列表 + TxnId
}

2b. BE 写入阶段——MemTable → DeltaWriter → RowsetBuilder

架构定位

BE 端接收客户端的数据流, 通过 HTTP Body → MemTable(内存排序) → Flush(磁盘 Rowset) 的流水线完成写入。这是写入路径最复杂的部分。

源码导航

文件关键类/方法职责
be/src/service/http/action/load_stream_action.hLoadStreamActionBE HTTP 接收数据
be/src/load/delta_writer/delta_writer.hDeltaWriter::write(), close()写入协调: MemTable → Flush
be/src/load/memtable/memtable.hMemTable::insert()内存排序缓冲(SkipList)
be/src/load/memtable/memtable_writer.cppMemTableWriterMemTable 的写入线程
be/src/load/memtable/memtable_flush_executor.cppMemTableFlushExecutorFlush MemTable 到磁盘
be/src/storage/rowset_builder.hBetaRowsetBuilder::build_rowset()从 Flush 结果构建 Rowset 文件

调用链

BE HTTP: LoadStreamAction
  → DeltaWriter::init()                     // 初始化写入上下文
  → 循环: DeltaWriter::write(Block*)        // 接收数据块
      → MemTable::insert(row)              // SkipList 内存排序
        → 每个 Tablet 对应一个 MemTable
  → DeltaWriter::close()
      → MemTableFlushExecutor::flush()
          → MemTable::flush_to_segment()    // 内存排序数据写磁盘
            → SegmentWriter::init()
              → SegmentWriter::append_block()
                → SegmentWriter::finalize() // 生成 Segment 文件
          → BetaRowsetBuilder::build_rowset()  // 从 Segment 构建 Rowset
            → RowsetMeta + Segment 文件清单

关键代码

// delta_writer.h:64 — write 接口
virtual Status write(const Block* block, const DorisVector<uint32_t>& row_idxs) = 0;

// delta_writer.h:138 — 触发 MemTable flush 成为 Rowset
Status commit_txn(const PSlaveTabletNodes& slave_tablet_nodes);

// MemTable 内部用 SkipList (MemTableRowSet) 做内存排序
// 目的: 写入数据按 (SortKey, RowId) 排序, 加速后续 Compaction

数据流图

数据流图

数据写入端到端——从 HTTP 到 Rowset 的全源码级穿越 图 02

标注每个阶段的状态: MemTable 阶段数据在内存, Flush 后数据在磁盘但不可见, Publish 后可见

2c. Publish 阶段——事务可见性的最后一步

架构定位

BE 写入完成后, FE 发起 Publish, 将 Tablet 的 Version 从当前升级到最新, 使新 Rowset 对查询可见。这是两阶段提交的第二阶段。

源码导航

文件关键类/方法职责
fe/fe-core/.../transaction/DatabaseTransactionMgr.javacommitTransaction(), publishVersion()FE 事务提交 + 触发 Publish
fe/fe-core/.../transaction/PublishVersionDaemon.javaPublishVersionDaemon定期发布版本(后台线程)
be/src/agent/task_worker_pool.cppPublishVersionWorkerPool::publish_version_callback()BE 处理 Publish 任务
be/src/storage/tablet.hTablet::publish_version()Tablet 版本升级

调用链

FE: DatabaseTransactionMgr.commitTransaction()
  → 1. 检查所有 BE 的写入状态
  → 2. 构造 PublishVersionTask (含新 Version 号)
  → 3. BE Agent RPC: PublishVersion
       → EnginePublishVersionTask::execute()
         → Tablet::publish_version()
           → 1. 将 Rowset 加入 Tablet 的 Rowset 列表
           → 2. 更新 Tablet 的 max_version
           → 3. Rowset 变为 VISIBLE (可被查询扫描)
  → 4. 标记事务为 VISIBLE

后台守护: PublishVersionDaemon 定期检查超时未 Publish 的事务并重试

Version 升级示意图

Publish 前:  Tablet version = 10
              Rowset: [0-9] (Base Compaction 历史) + [10] (刚刚写入的 Delta Rowset)
                                            ↑ 状态: COMMITTED (不可见)

Publish 后:  Tablet version = 10
              Rowset: [0-9] + [10]           ← 现在可见
              visible_version = 10

关键代码

// task_worker_pool.cpp:2109 — BE Publish 回调
void PublishVersionWorkerPool::publish_version_callback(const TAgentTaskRequest& req) {
    const auto& publish_version_req = req.publish_version_req;
    EnginePublishVersionTask engine_task(_engine, publish_version_req, &error_tablet_ids, ...);
    engine_task.execute();  // 遍历所有涉及的 Tablet, 执行 publish_version
}

3. Broker Load——外部存储批量导入

架构定位

Broker Load 从 HDFS/S3 等外部存储拉取大量数据(GB~TB级)。与 Stream Load 的核心区别: (1) 数据不经过客户端, 直接由 Broker 进程或 BE 从外部存储读取; (2) 是异步作业, 由 FE 提交 + 调度; (3) 支持多表、多分区批量导入。

与 Stream Load 的差异点

维度Stream LoadBroker Load
提交方式HTTP PUT (同步)SQL LOAD LABEL ... (异步)
数据路径Client → BEHDFS/S3 → Broker/BE 直接拉取
数据格式CSV/JSON (HTTP body)Parquet/ORC/CSV/JSON (外部文件)
事务模型单 Txn每个 Task 一个 Txn, 整体成功或全部失败
关键类StreamLoadHandlerBrokerLoadJob, BrokerFileGroup

调用链概要

FE: BrokerLoadJob
  → 1. 拉取文件清单 (从 HDFS NameNode 或 S3 List)
  → 2. 按 Tablet 拆分文件范围 → Task
  → 3. 分配给 BE → BE 从 HDFS/S3 读取数据 → 正常写入路径 (DeltaWriter)
  → 4. 所有 Task 完成 → commit + Publish

4. Routine Load——Kafka 持续导入

架构定位

Routine Load 维护一个或多个 Kafka Consumer, 持续消费消息并写入 Doris, 支持 Exactly-Once 语义(基于 Kafka offset + Doris Label 去重)。

关键差异

维度Stream LoadRoutine Load
触发方式手动 HTTP 请求自动消费 Kafka
连接模型一次性长连接(Kafka Consumer)
事务粒度每个 HTTP 请求 1 Txn每个 Kafka Partition 定期 Commit
关键类StreamLoadHandlerRoutineLoadManager, RoutineLoadTask

Exactly-Once 语义的实现

Kafka Topic [Partition 0] → RoutineLoadTask
  → 1. 从 Kafka 拉取一批消息 (offset 101-200)
  → 2. 写入 Doris (Label = "routine_load_job_20260708_0_101")
  → 3. 写入成功 → 提交 Kafka offset
  → 4. 写入失败 → 不提交 offset → 重试 (Label 相同 → Doris 去重保证幂等)

Label 去重: Doris 的 Txn Label 全局唯一, 相同 Label 的事务只能成功一次。如果 RoutineLoadTask 重试同一个批次, Doris 发现 Label 已存在则直接返回成功, 保证端到端的 Exactly-Once。


5. Compaction——写入的后台收尾

为什么需要 Compaction

每次写入生成一个 Delta Rowset。多个小 Rowset 会导致:

  1. 查询性能下降: 扫描时需要合并多个 Rowset
  2. 元数据膨胀: Tablet 中的 Rowset 列表过长

Compaction 定期将多个小 Rowset 合并为一个大 Rowset, 回收重复/删除的数据。

Compaction 类型

Base Compaction: 合并所有 Rowset → 单个大 Rowset
                 (每日低峰执行, 合并后删除旧 Rowset)

Cumulative Compaction: 合并最近 N 个 Delta Rowset → 一个 Merged Rowset
                       (持续后台执行, 控制 Rowset 数量)

生命周期示意

写入: Data → MemTable → Flush → Rowset [Version 11]
       Data → MemTable → Flush → Rowset [Version 12]
       Data → MemTable → Flush → Rowset [Version 13]

查询前: Rowset[0-9] + Rowset[10] + Rowset[11] + Rowset[12] + Rowset[13]

Cumulative Compaction:
  Rowset[11] + Rowset[12] + Rowset[13] → Merged Rowset [11-13]

更新后的查询: Rowset[0-9] + Rowset[10] + Rowset[11-13]  ← 少一个 Rowset

详细的 Compaction 机制 → doc-07


六种导入方式对比

下一步


常见问题 / 面试题

Q1: Stream Load 的 Label 去重如何保证 Exactly-Once? A: Label 在 FE 端全局唯一。客户端重试同一 Label 时,FE 检查上次事务状态:VISIBLE → 幂等返回成功;PREPARE → 返回原 txnId 继续使用;ABORTED → 创建新事务。Label 有有效期限制(label_keep_max_second,默认3天),超期清理后重试会退化到 At-Least-Once。

Q2: MemTable 为什么要用 SkipList 做内存排序? A: 写入数据按 (SortKey, RowId) 排序后写入 Segment,Compaction 时多个 Rowset 可以按序 Merge,不需要全量排序。SkipList 的插入复杂度 O(log N),内存友好(节点大小可控)。

Q3: Publish 阶段有什么失败场景? 如何恢复? A: 常见的:BE 在 Publish 期间宕机 → PublishVersionTask 失败 → PublishVersionDaemon 后台重试(间隔递增)。如果所有副本都失败 → 事务最终 Abort → 用户重试整个 Load。BE 恢复后通过 Clone 从健康副本补数据。

Q4: Broker Load 和 Routine Load 的 Exactly-Once 语义有何不同? A: Broker Load 是"整体成功或全部失败"(所有 Task 完成才 Commit),没有中间状态。Routine Load 是"逐批提交"(每个 Kafka offset 范围一个 Txn),通过 Kafka offset + Doris Label 组合保证每个批次的 Exactly-Once。

Q5: Group Commit 的合并策略如何决定何时 Flush? A: 两个条件任一满足即 Flush:(1) group_commit_interval_ms(默认 10ms)超时;(2) group_commit_max_data_size(默认 128MB)写满。在低负载时靠间隔触发(保证低延迟),高负载时靠大小触发(保证吞吐)。