Doris architecture

ANALYTICAL DATABASE / SOURCE READING / LESSON 10

Data ingestion system

Compare Doris ingestion methods and their transaction lifecycles across Stream Load, Routine Load, Broker Load, and related paths.

Reading
50 分钟 | **前置阅读**: [doc-02](doc-02-write-lifecycle.md) (写入主线), [doc-07](doc-07-storage-compaction-version.md) (Compaction)
Track
Doris architecture
Source
Chinese source notes

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

预计阅读时间: 50 分钟 | 前置阅读: doc-02 (写入主线), doc-07 (Compaction) 下一次阅读: doc-12 (BE副本), doc-11 (FE HA)


1. 导入方式全景对比

方式协议数据源最大延迟事务保证最佳规模
Stream LoadHTTP PUT本地 CSV/JSON秒级Exactly-Once<10GB/次
INSERT INTOMySQL SQL客户端 VALUES毫秒-秒Exactly-Once<1MB/次
Group CommitINSERT 自动多客户端小批量合并毫秒-秒Exactly-Once高频小批量
Broker LoadSQL 提交HDFS/S3 大文件分钟At-Least-OnceGB-TB/次
Routine LoadSQL 提交Kafka 持续消费秒级Exactly-Once持续流
Binlog LoadCanal 订阅MySQL Binlog (CDC)秒级Exactly-OnceCDC

源码导航

文件关键类/方法职责
fe/fe-core/.../load/StreamLoadHandler.javaexecuteLoad()Stream Load 协调
fe/fe-core/.../load/BrokerLoadJob.javaBrokerLoadJobBroker Load 作业
fe/fe-core/.../load/RoutineLoadManager.javaRoutineLoadManagerRoutine Load 管理
fe/fe-core/.../transaction/GlobalTransactionMgr.javabeginTransaction(), commitTransaction()全局事务管理
fe/fe-core/.../load/GroupCommitManager.javaGroupCommitManagerGroup Commit 合并
be/src/load/delta_writer/delta_writer.hDeltaWriterBE 写入协调
be/src/load/memtable/memtable.hMemTable内存排序缓冲

2. 统一事务模型

数据导入体系全览 图 01

TxnManager——所有导入方式共享

所有导入方式
    └── GlobalTransactionMgr  (fe/fe-core/.../transaction/)
        ├── beginTransaction(dbId, tableId, label)
        │   → 创建 TxnState {label, status=PREPARE, ...}
        │   → Label 去重检查 (再次提交相同 label → 返回已有 txnId)
        │
        ├── commitTransaction(txnId, tabletCommitInfos)
        │   → 检查所有 BE 写入完成
        │   → 标记 Txn status=COMMITTED
        │   → 触发 PublishVersionDaemon
        │
        └── abortTransaction(txnId)
            → 清理已写数据, 释放资源

Label 去重——Exactly-Once 的基石

Label = "my_load_20260708_001"

第一次提交:  beginTxn("my_load_20260708_001")
  → txnId=1001, status=PREPARE → 写入数据 → commit → status=VISIBLE

第二次提交 (重试): beginTxn("my_load_20260708_001")
  → Label 已存在, 检查上一次的状态:
    - status=VISIBLE → 返回 "Duplicate label, already committed"
    - status=PREPARE → 返回 txnId=1001 (继续使用同一个事务)
    - status=ABORTED → 创建新 txnId

3. Group Commit——高频小批量写入

原理

多个客户端的 INSERT 请求在 BE 的 MemTable 层面合并, 然后一起 Flush, 大幅减少 Compaction 压力:

Client A: INSERT INTO t VALUES (1, 'a')      → MemTable
Client B: INSERT INTO t VALUES (2, 'b')      → MemTable
Client C: INSERT INTO t VALUES (3, 'c')      → MemTable
                    │
            ┌───────▼───────┐
            │  Group Commit  │  ← 等待 10ms 或达到 N 行
            │  一次性 Flush   │
            └───────┬───────┘
                    ▼
           1 个 Rowset (含3行)

参数

参数默认说明
group_commit_interval_ms10ms最长等待时间
group_commit_max_data_size128MB单次合并最大数据量

4. Broker Load——大批量外部文件导入

流程

FE: SUBMIT LOAD LABEL ... DATA INFILE "hdfs://path/*" INTO TABLE t
  → BrokerLoadJob
    → 1. 列出 HDFS 文件列表 (NameNode Client)
    → 2. 为每个文件范围划分 Task
    → 3. 分配 Task 到 BE
    → 4. BE 从 HDFS 拉取数据 → 正常写入 (DeltaWriter)
    → 5. 所有 Task 完成 → commit + Publish

关键: 数据从 HDFS → BE 是直接流式读取, 不经过 FE 中转。FE 只协调 Task 元数据。


5. Routine Load + Kafka——Exactly-Once 实现

Kafka Consumer 模型

Kafka Topic [MyTopic]
├── Partition 0 → RoutineLoadTask[0]  (1个 BE 上的 Consumer)
├── Partition 1 → RoutineLoadTask[1]
└── Partition 2 → RoutineLoadTask[2]

每个 Task:
  while (true):
    1. KafkaConsumer.poll() → 一批消息
    2. Label = "jobName_partition_offset" (唯一)
    3. FE: beginTxn(label) → txnId
    4. BE: DeltaWriter.write(messages) → MemTable
    5. FE: commitTxn(txnId)
    6. KafkaConsumer.commitSync(offset)  ← 现在才提交 offset

故障恢复: Task 崩溃 → Kafka offset 未提交 → 重启后从上次 committed offset 开始 → Label 相同 → Doris 去重保证幂等 → Exactly-Once


6. 性能调优参数

参数说明影响
write_batch_sizeMemTable 每批行数太小→频繁 Flush; 太大→内存涨
send_batch_parallelismStream Load 并行度单节点瓶颈→调大
load_parallel_instance_numBroker Load 并行度BE 数 × 2-4
routine_load_task_consume_secondKafka 每批消费间隔秒级
max_tolerable_backend_down_numBroker Load 容忍故障 BE数据管道鲁棒性

7. 常见问题 / 面试题

Q1: Exactly-Once 的边界在哪? Label 有效期多久? A: Label 有效期由 label_keep_max_num + label_keep_max_second 控制(默认 1000 个/3天)。超出范围的旧 Label 被清理后, 重试就会产生重复数据 → 退化到 At-Least-Once。

Q2: Group Commit vs Stream Load 选哪个? A: Group Commit 适合高频小批量 (如每秒数百个 INSERT), 自动合并, 无需客户端改造。Stream Load 适合周期性大文件导入 (如每小时一次 CSV import)。Group Commit 本质上是 Stream Load 的"微批优化"——内部也转换为 Stream Load。

Q3: Routine Load 的 Partition 数与 Doris Tablet 数如何匹配? A: 建议 Kafka Partition 数 = Doris Bucket 数 (或倍数)。Kafka Partition 太少 → 写入并行度不足; Kafka Partition 太多 → 单 Tablet 写入频繁, MemTable 过于碎片化, 导致 Compaction 压力大。


导入事务生命周期 (Txn Lifecycle)

下一步