预计阅读时间: 90-120 分钟 前置阅读: doc-00(架构概念), doc-01(理解 FE 协调和 BE 执行基底) 下一次阅读: doc-06(Rowset/Segment 存储), doc-07(Compaction), doc-10(导入体系全览)
1. 写入方式总览与对比
Doris 支持 6 种数据写入方式, 每种针对不同的数据源和延迟要求。
| 方式 | 协议 | 数据源 | 延迟 | 适用规模 | 关键类 |
|---|---|---|---|---|---|
| Stream Load | HTTP (multipart) | 本地文件/程序输出 | 秒级 | 单次 <10GB | StreamLoadHandler / StreamLoadAction |
| Broker Load | SQL 提交 | HDFS / S3 上的大文件 | 分钟级 | GB~TB 级 | BrokerLoadJob |
| Routine Load | SQL 提交 | Kafka (持续消费) | 秒级 | 持续流式 | RoutineLoadManager / RoutineLoadTask |
| INSERT INTO | MySQL SQL | 客户端 SQL | 秒级 | 小批量 | InsertStmtExecutor / GroupCommitManager |
| Group Commit | 自动(INSERT 聚合) | 多客户端 INSERT 合并 | 毫秒级 | 高频小批量 | GroupCommitManager |
| Binlog Load | Canal/Binlog 订阅 | MySQL Binlog (CDC) | 秒级 | CDC 持续 | BinlogManager |
架构对比图:

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.java | LoadAction.execute() | HTTP 请求入口, 处理 multipart body |
fe/fe-core/.../load/StreamLoadHandler.java | StreamLoadHandler.handleStreamLoad() (via TStreamLoadPutResult) | Stream Load 核心协调逻辑 |
fe/fe-core/.../transaction/GlobalTransactionMgr.java | beginTransaction(), commitTransaction() | 全局事务管理 |
fe/fe-core/.../transaction/DatabaseTransactionMgr.java | commitTransaction() | 单 DB 事务管理 |
fe/fe-core/.../planner/StreamLoadPlanner.java | StreamLoadPlanner.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.h | LoadStreamAction | BE HTTP 接收数据 |
be/src/load/delta_writer/delta_writer.h | DeltaWriter::write(), close() | 写入协调: MemTable → Flush |
be/src/load/memtable/memtable.h | MemTable::insert() | 内存排序缓冲(SkipList) |
be/src/load/memtable/memtable_writer.cpp | MemTableWriter | MemTable 的写入线程 |
be/src/load/memtable/memtable_flush_executor.cpp | MemTableFlushExecutor | Flush MemTable 到磁盘 |
be/src/storage/rowset_builder.h | BetaRowsetBuilder::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
数据流图
数据流图

标注每个阶段的状态: MemTable 阶段数据在内存, Flush 后数据在磁盘但不可见, Publish 后可见
2c. Publish 阶段——事务可见性的最后一步
架构定位
BE 写入完成后, FE 发起 Publish, 将 Tablet 的 Version 从当前升级到最新, 使新 Rowset 对查询可见。这是两阶段提交的第二阶段。
源码导航
| 文件 | 关键类/方法 | 职责 |
|---|---|---|
fe/fe-core/.../transaction/DatabaseTransactionMgr.java | commitTransaction(), publishVersion() | FE 事务提交 + 触发 Publish |
fe/fe-core/.../transaction/PublishVersionDaemon.java | PublishVersionDaemon | 定期发布版本(后台线程) |
be/src/agent/task_worker_pool.cpp | PublishVersionWorkerPool::publish_version_callback() | BE 处理 Publish 任务 |
be/src/storage/tablet.h | Tablet::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 Load | Broker Load |
|---|---|---|
| 提交方式 | HTTP PUT (同步) | SQL LOAD LABEL ... (异步) |
| 数据路径 | Client → BE | HDFS/S3 → Broker/BE 直接拉取 |
| 数据格式 | CSV/JSON (HTTP body) | Parquet/ORC/CSV/JSON (外部文件) |
| 事务模型 | 单 Txn | 每个 Task 一个 Txn, 整体成功或全部失败 |
| 关键类 | StreamLoadHandler | BrokerLoadJob, 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 Load | Routine Load |
|---|---|---|
| 触发方式 | 手动 HTTP 请求 | 自动消费 Kafka |
| 连接模型 | 一次性 | 长连接(Kafka Consumer) |
| 事务粒度 | 每个 HTTP 请求 1 Txn | 每个 Kafka Partition 定期 Commit |
| 关键类 | StreamLoadHandler | RoutineLoadManager, 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 会导致:
- 查询性能下降: 扫描时需要合并多个 Rowset
- 元数据膨胀: 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

下一步
- 存储格式: doc-06-storage-tablet-rowset-segment.md——Tablet/Rowset/Segment 的物理格式
- Compaction 详解: doc-07-storage-compaction-version.md——版本管理与 Compaction
- 导入体系全景: doc-10-data-ingestion.md——所有导入方式深度对比
常见问题 / 面试题
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)写满。在低负载时靠间隔触发(保证低延迟),高负载时靠大小触发(保证吞吐)。