第 09 关 · ★★★

批量与实时写入

Stream Load、Group Commit、Broker/TVF 与 INSERT:统一导入规范与重试策略。

已点亮 · 最佳 分

Doris 批量导入体系:Stream Load、Group Commit、TVF、INSERT 与 Broker Load

数据怎样抵达 Doris,事务何时形成,异常与重试如何处理

导入坐标系Stream LoadGroup CommitTVF + INSERT质量与幂等
1 / 28

Day 9|Doris 批量导入体系:Stream Load、Group Commit、TVF、INSERT 与 Broker Load

Day 9 学习总览
Day 9 学习总览

Day 8 完成了表结构设计。我们已经知道,一行数据会先进入某个 Partition,再按照分桶规则进入某个 Tablet,最终在多个副本上形成 Rowset 与 Segment。Day 9 把视角放到这条路径的入口:数据怎样抵达 Doris,什么时候形成事务,哪台 BE 负责协调,异常行怎样处理,客户端怎样安全重试,批次过小为什么会拖累 Compaction,远端文件又该通过哪条路径进入目标表。

数据导入看起来像“把文件写进表”,实际涉及一组彼此关联的工程决策:

  • 数据位于客户端、本地文件、对象存储、HDFS,还是 Doris 内部表;
  • 客户端是否需要同步等待结果;
  • 一批数据有多大,多久产生一批,有多少并发写入方;
  • 上游能否稳定生成业务批次号;
  • 异常数据需要整批拒绝,还是允许有限比例的过滤;
  • 写入目标是 Duplicate、Aggregate 还是 Unique Key 表;
  • 数据可见延迟、吞吐、版本数量和重试成本怎样平衡;
  • 新系统是否仍应采用 Broker Load,还是转向 TVF 与 INSERT INTO ... SELECT

这些问题决定一套导入链路能否长期稳定。单次成功只说明数据进入了 Doris。生产可用还要求任务能够重试、对账、观测、限流、回滚,并且不会把 FE 元数据、BE MemTable、磁盘 I/O 和后台 Compaction 持续推向高压状态。

本日学习目标

完成本篇后,你应当能够:

  1. 从数据源、时效、批量大小、并发和容错要求选择导入方式;
  2. 解释 Stream Load 从 HTTP 请求到数据可见的完整链路;
  3. 正确使用 Expect: 100-continue、Label、格式、列映射和过滤参数;
  4. 读懂 Stream Load 返回 JSON,并据此决定成功、重试或人工介入;
  5. 理解 Group Commit 的队列、同步模式、异步模式、WAL 与回退条件;
  6. 说明高频小批次怎样形成版本堆积和 Compaction 压力;
  7. 使用 S3/HDFS TVF 在写入前预览、过滤和转换远端文件;
  8. 使用 INSERT INTO ... SELECT 完成内部 ETL 和外部文件落表;
  9. 理解 Broker Load 的异步执行机制、现有价值与迁移方向;
  10. 建立 Label、数据质量、对账、批次、并发和故障恢复的生产规范。

实验前提

本篇沿用 Day 5 的单 FE、单 BE 学习环境。实验表使用单副本,示例中的 S3、HDFS、对象存储凭据均为占位符。生产环境必须使用独立账号、最小权限和密钥管理服务,严禁把真实 AK/SK 提交到代码仓库。

版本口径

截至 2026-08-16,Apache Doris 官网将 4.1.3 标记为 Latest,将 4.0.8 标记为 Stable。当前 4.x 官方文档已经把 Broker Load 标记为 Deprecated,并说明计划在 5.0 移除;新建的远端文件导入链路应优先评估 S3/HDFS TVF、Catalog 表与 INSERT INTO ... SELECT。现有 Broker Load 作业仍可在对应 4.x 版本中运行,迁移需要先完成语义与性能回归。


一、先建立导入坐标系:数据在哪里,谁主动,结果怎样返回

Doris 导入方式全景
Doris 导入方式全景

培训材料把 Doris 导入方式分成同步与异步两类:Stream Load 和 INSERT INTO 会让客户端等待结果;Broker Load 与 Routine Load 由 Doris 在后台持续或异步执行。这个框架很适合建立第一层认知。进入 4.x 之后,还需要增加两个维度:数据是由客户端推送,还是由 Doris 主动拉取;任务是一次性批次,还是持续数据流。

1. 客户端推送

客户端已经拿到数据内容,通过网络把字节流发给 Doris。典型方式是 Stream Load:

本地 CSV / JSON / Parquet / ORC
             ↓ HTTP PUT
            Doris

这条路径适合应用程序、脚本、Connector 或数据采集服务。客户端负责组织批次、生成 Label、处理返回值和实施重试。

2. Doris 主动拉取

数据已经位于 S3、OSS、COS、OBS、HDFS 等远端存储。Doris 读取文件并写入目标表。当前可选路径主要有:

  • S3/HDFS TVF + INSERT INTO ... SELECT
  • Catalog 表 + INSERT INTO ... SELECT
  • 4.x 中仍可使用的 Broker Load / S3 Load / HDFS Load;
  • 面向持续新增文件的 Streaming Job。

主动拉取能减少客户端搬运大文件的网络压力,也方便 Doris 根据文件和集群资源组织并行读取。

3. Doris 内部计算后落表

数据源可以是 Doris 内部表、Catalog 外表或 TVF。SQL 完成过滤、Join、聚合和字段转换,OlapTableSink 把结果写入目标表:

INSERT INTO dws_sales_daily
SELECT
    order_date,
    shop_id,
    SUM(pay_amount) AS pay_amount
FROM dwd_orders
WHERE order_date = '2026-08-16'
GROUP BY order_date, shop_id;

这类任务适合数仓分层、历史回填、口径修复和轻量 ELT。数据量较大时,可以通过 Job 机制异步执行,并单独设置资源组与运行窗口。

4. 持续数据流

Kafka、数据库 Binlog 和持续新增的对象存储文件需要长期运行的任务。Routine Load、Kafka Connector、Flink Doris Connector、Streaming Job 与 CDC Stream TVF 会在 Day 10 系统展开。Day 9 只保留一条边界:持续流的核心难题包括 Offset、Checkpoint、乱序、重启恢复和 Schema 演进,单次文件导入无法覆盖这些职责。

5. 选型时先问五个问题

在讨论参数之前,先完成下面五项判断:

  1. 数据源位于哪里?
  2. 单批数据量和文件数量是多少?
  3. 数据多久必须可见?
  4. 客户端是否能够稳定合批和生成业务批次号?
  5. 任务失败后,谁负责重试、从哪里继续、怎样证明没有重复?

导入方式的名字只是入口。生产设计需要围绕这五个问题形成明确答案。


二、所有导入方式共享一条核心写入链路

Doris 通用导入链路
Doris 通用导入链路

Stream Load、Broker Load、Routine Load 和 INSERT INTO 的提交方式不同,进入 Doris 后会经历一条相近的核心路径。

1. 提交任务与参数校验

客户端通过 HTTP、MySQL 协议或持久化 Job 提交任务。FE 或协调 BE 会校验:

  • 用户是否拥有目标表的 INSERT 权限;
  • 数据库、表、分区和字段是否存在;
  • 文件格式、分隔符、JSON 路径和列映射是否合法;
  • Label 是否冲突;
  • 目标表的数据模型和更新参数是否匹配;
  • 当前用户是否具备目标 Compute Group 或资源标签的访问权限。

2. 生成导入计划

FE 根据表的分区、分桶、副本和字段定义生成 Load Plan,并选择协调节点。远端文件导入还会结合文件数量、文件大小和 BE 数量拆分读取任务。

3. 读取、解析与轻量转换

协调 BE 接收 HTTP 数据流,或从 S3/HDFS 主动读取文件。解析器完成 CSV 切分、JSON 解析、Parquet/ORC 解码,并执行列映射、表达式转换和过滤。

这里需要区分四个概念:

  • 列映射:源字段与目标字段建立对应关系;
  • 列转换:使用表达式生成目标值;
  • 前置过滤:在转换前过滤原始行,Broker Load 等方式支持;
  • 后置过滤:在转换后按照目标语义过滤,Stream Load 可通过 where 使用。

4. 按分区与分桶路由

解析后的每一行根据分区规则找到目标 Partition,再根据 Hash 或 Random 分桶进入 Tablet。数据可能从协调 BE Shuffle 到持有目标 Tablet 副本的其他 BE。

5. MemTable、Rowset 与 Segment

每个活跃 Tablet 在内存中维护 MemTable。数据按照 Key 排序,Aggregate 或 Unique Key 表还会执行聚合、去重、主键索引查找或 Delete Bitmap 维护。MemTable 达到阈值或任务结束后 Flush,形成列式 Segment,并归入本次事务对应的 Rowset。

一批数据同时覆盖大量分区和 Tablet 时,会激活大量 MemTable。内存压力会触发提前 Flush,生成更多小 Segment。导入吞吐、表的分桶数量、单批覆盖的分区数和 Compaction 健康度需要一起观察。

6. 事务提交与版本发布

常规导入任务会开启事务。各 BE 写入完成后,协调节点向 FE 或 Meta Service 提交事务。多数副本写入成功后,系统发布新版本,事务进入 VISIBLE,查询才能读到新数据。失败任务会回滚临时数据,避免出现半批可见。

Group Commit 对事务边界进行了合并:多个小请求进入同一个服务端批次,由一个合并事务提交。它的可见性和 Label 行为需要单独理解,后文会详细说明。


三、Stream Load:最通用的 HTTP 批量写入接口

Stream Load 请求与写入路径
Stream Load 请求与写入路径

Stream Load 通过一次 HTTP PUT 把文件或字节流写入一张 Doris 表。客户端在同一连接上拿到 JSON 结果,因此很适合脚本、微服务和 Connector。

1. 请求怎样流转

标准路径可以概括为:

Client
  → FE HTTP 端口
  → HTTP 307 重定向
  → Coordinator BE
  → 解析与路由
  → Tablet 所在 BE
  → Commit / Publish
  → JSON Response

使用 curl --location-trusted 时,客户端会跟随 FE 的 307 重定向,并在重定向后继续携带认证信息。也可以直接向健康 BE 发起请求,减少一次跳转;生产环境需要配合服务发现、负载均衡和故障摘除机制。

2. 一个可运行的 CSV 示例

curl --location-trusted \
  -u "${DORIS_USER}:${DORIS_PASSWORD}" \
  -H "Expect:100-continue" \
  -H "label:day09_orders_20260816_0001" \
  -H "format:csv" \
  -H "column_separator:," \
  -H "columns:order_id,user_id,order_time,amount,status" \
  -T data/orders.csv \
  -X PUT \
  "http://${DORIS_FE_HOST}:8030/api/day09_lab/orders_detail/_stream_load"

Expect:100-continue 让服务端先完成认证和请求头校验,再接收文件正文。大文件在权限或参数错误时可以尽早终止,减少无效上传。

3. Stream Load 支持哪些格式

当前 4.x 文档覆盖 CSV、JSON、Parquet、ORC 和 CSV 带表头变体。格式选择会影响解析成本、文件大小和并行能力:

  • CSV 简单、通用,分隔符和转义规则必须明确;
  • JSON 适合半结构化事件,需确定逐行 JSON、外层数组和 JSONPath;
  • Parquet、ORC 自带类型和列式编码,适合文件已经由上游批处理系统生成的场景;
  • 压缩文件可以节省网络,解压过程会消耗 CPU,单个不可切分的大压缩文件还会限制并行度。

4. 常用请求头

Stream Load 请求合同与返回值
Stream Load 请求合同与返回值

参数 作用 生产关注点
label 标识本批事务 使用业务批次号,保证可追踪和安全重试
format CSV、JSON、Parquet、ORC 与真实文件格式一致
column_separator CSV 列分隔符 注意制表符、转义和多字节分隔符
line_delimiter CSV 行分隔符 跨平台文件需核验换行符
columns 列映射与转换 显式写出源字段和计算字段
where 转换后的过滤条件 被过滤行会进入统计
strict_mode 类型转换失败时的处理 当前默认 false,生产任务应显式设置
max_filter_ratio 允许过滤的最大比例 默认值为 0,质量门禁应当谨慎放宽
timeout 导入超时 与文件大小、网络和集群负载匹配
partial_columns Unique Key 部分列更新 会改变补列和更新路径
two_phase_commit Stream Load 2PC 用于跨系统 Checkpoint 协调
group_commit Group Commit 模式 off_modesync_modeasync_mode

列映射可以完成简单清洗:

-H 'columns:order_id,user_id,src_time,src_amount,status,
    order_time=str_to_date(src_time,"%Y-%m-%d %H:%i:%s"),
    amount=cast(src_amount as decimal(18,2))'
-H 'where:amount >= 0'

复杂逻辑、跨表关联和可复用业务规则更适合放在上游处理层,或使用 TVF / Catalog 配合 INSERT INTO ... SELECT。导入表达式越复杂,失败定位和版本升级回归成本越高。

5. 如何读懂返回 JSON

典型成功结果包括:

{
  "TxnId": 30021,
  "Label": "day09_orders_20260816_0001",
  "Status": "Success",
  "NumberTotalRows": 100000,
  "NumberLoadedRows": 99998,
  "NumberFilteredRows": 2,
  "LoadBytes": 14738291,
  "LoadTimeMs": 1280
}

客户端至少需要检查:

  • Status 是否属于成功状态;
  • NumberLoadedRows + NumberFilteredRows 是否与总行数相符;
  • 过滤比例是否低于业务门限;
  • ErrorURL 是否存在;
  • LoadTimeMs 是否异常抬升;
  • Label 和 TxnId 是否已记录到任务审计表。

HTTP 状态码成功并不能代替业务结果检查。Doris 可能返回 HTTP 200,同时在 JSON 中给出数据质量失败、Label 冲突或其他任务状态。调用方需要解析响应体。

6. Label 与安全重试

Label 在一个数据库内唯一,用于识别一次导入事务。官方文档允许系统在缺省时生成 UUID;生产任务仍应使用确定性的业务 Label,例如:

orders_20260816_10m_0006

网络超时后,客户端无法确认服务端是否已经提交。使用同一个 Label 重试时:

  • 前一次成功:返回 Label Already Exists,不会重复写入;
  • 前一次失败:相同 Label 可以再次执行;
  • 前一次仍在运行:需要查询任务状态,避免盲目并发重试。

Label 有清理周期和数量上限。它提供一定窗口内的幂等保护,长期批次去重仍需要上游任务表、数据对账和业务主键语义共同保证。

7. JSON 与半结构化批次的额外检查

JSON 导入需要先确认输入形态。逐行 JSON 应开启对应的逐行读取参数;一个文件只有一个外层数组时,需要按数组模式解析;嵌套字段可通过 JSONPath 或列映射提取。生产前应使用几十行真实样本覆盖以下情况:字段缺失、字段为 JSON Null、字符串 "null"、空字符串、数字被引号包裹、时间精度变化、数组为空、同一字段在不同记录中出现不同类型。

当 JSON 字段长期变化且查询路径不稳定时,可以把稳定公共字段拆为实体列,把长尾扩展字段写入 VARIANT。导入阶段仍需保留原始事件标识、事件时间、来源系统和 Schema 版本,方便后续追溯。直接把所有字段映射为字符串会降低类型校验能力,也会把质量问题推迟到查询阶段。

8. 客户端 SDK 应统一封装哪些行为

企业内部不宜让每个团队独立拼接 curl。可以提供统一 Stream Load SDK 或脚本模板,至少封装:

  • FE/BE 地址发现、连接超时、读取超时和有限重试;
  • Label 生成规则与冲突查询;
  • CSV、JSON、Parquet、ORC 的标准 Header 模板;
  • 返回 JSON 的状态机,区分成功、幂等命中、结果不确定和明确失败;
  • ErrorURL 下载、脱敏和异常样本归档;
  • 源行数、加载行数、过滤行数、字节数和耗时指标上报;
  • 密码与令牌从密钥系统读取,日志中自动脱敏;
  • 请求体校验和、源文件清单、批次审计和告警事件;
  • 限流、熔断和最大并发保护,防止故障期间形成重试风暴。

SDK 的返回对象应包含业务批次号、Doris Label、TxnId、状态、行数、过滤比例、ErrorURL、耗时和建议动作。上层调度器根据建议动作选择结束、等待、查询状态、有限重试或转人工处理。


四、Group Commit:把高频小写入合并成健康事务

Group Commit 的队列、模式与事务
Group Commit 的队列、模式与事务

高频小批次会让每个请求都创建事务、写 FE 元数据、Flush MemTable,并在 Tablet 内新增 Rowset 和 Segment。业务端看见的是低延迟,存储层承受的是版本数量、文件数量和 Compaction 写放大。

1. Group Commit 做了什么

Group Commit 在 BE 端为同一张表维护共享写入队列。多个并发请求进入同一个 LoadBlockQueue,达到时间阈值或数据量阈值后,由一个事务统一提交。当前官方文档给出的默认刷新条件是 10 秒或 128 MB,可通过表属性调整:

ALTER TABLE events SET (
  "group_commit_interval_ms" = "1000",
  "group_commit_data_bytes" = "67108864"
);

阈值需要围绕可见延迟和吞吐目标验证。1 秒刷新适合秒级可见,10 秒刷新能形成更大批次。数据量阈值应结合单行大小、并发数和 BE 内存评估。

2. 三种模式

off_mode   常规导入路径
sync_mode  合批后等待事务提交和数据可见
async_mode 写入本地 WAL 后立即返回 PREPARE,稍后统一提交

INSERT 可以通过会话变量开启:

SET group_commit = async_mode;

INSERT INTO events VALUES
(1, '2026-08-16 13:00:00.000', 'open'),
(2, '2026-08-16 13:00:00.100', 'click');

Stream Load 可以通过请求头开启:

-H "group_commit:async_mode"

async_mode 返回 PREPARE 时,数据已经写入接收 BE 的本地 WAL,尚未进入查询可见状态。调用方若要求“返回后立刻能查到”,应选择 sync_mode,或在业务侧等待可见性确认。

3. Group Commit 的回退条件

Group Commit 会在部分请求形态下回到常规导入路径,常见条件包括:

  • 显式事务;
  • 用户自行指定 Label;
  • Stream Load 2PC;
  • partial_columns:true
  • VALUES 中包含特定表达式;
  • 异步模式可用 WAL 空间不足。

调用方只看到请求成功,未必知道该请求是否真正合入共享事务。性能验证应记录返回的 Label 和 TxnId;多个请求共享同一组值,通常说明合批生效。

4. 顺序语义与 Unique Key

Group Commit 会合并不同客户端的写入,不能把网络到达顺序直接当成业务版本顺序。Unique Key 表若存在乱序更新,需要配置 Sequence Column 或 Sequence Mapping,让 Doris依据业务时间、Binlog 位点或单调版本号裁决新旧。

5. Group Commit 的适用范围

它适合:

  • 大量微服务进行几十 TPS 以上的小批量 JDBC 写入;
  • 日志、IoT、埋点采集端难以在客户端形成大批次;
  • 小 Stream Load 已经造成版本堆积;
  • 业务允许秒级可见,且希望降低事务与 Compaction 压力。

单批已经达到数百 MB 时,常规 Stream Load 更直接。需要自定义 Label、2PC 或部分列更新时,应使用对应的常规写入能力。


五、Broker Load:理解现有能力,也要看清迁移方向

Broker Load 与 TVF + INSERT 的演进
Broker Load 与 TVF + INSERT 的演进

培训材料把 Broker Load 作为 HDFS/S3 大规模历史导入的主要方式:FE 生成执行计划,多个 BE 并行拉取远端文件,任务异步运行,并通过 SHOW LOAD 查看状态。这个机制在很多现有 Doris 集群中仍然承担历史迁移和日批任务。

当前 4.x 官方文档已经增加重要提示:Broker Load 被标记为 Deprecated,并计划在 5.0 移除,新建链路应使用 TVF 或 Catalog 表配合 INSERT INTO ... SELECT

1. 当前 Broker Load 的工作方式

LOAD LABEL day09_lab.orders_s3_20260816
(
    DATA INFILE("s3://example-bucket/orders/dt=2026-08-16/*.parquet")
    INTO TABLE orders_detail
    FORMAT AS "parquet"
    (order_id, user_id, order_time, amount, status)
)
WITH S3
(
    "provider" = "S3",
    "AWS_ENDPOINT" = "https://s3.us-west-2.amazonaws.com",
    "AWS_REGION" = "us-west-2",
    "AWS_ACCESS_KEY" = "<from-secret-manager>",
    "AWS_SECRET_KEY" = "<from-secret-manager>"
)
PROPERTIES
(
    "timeout" = "3600"
);

任务提交后立即返回,状态通过下面的命令查看:

SHOW LOAD FROM day09_lab
WHERE LABEL = 'orders_s3_20260816';

未进入 FINISHEDCANCELLED 的任务可以取消:

CANCEL LOAD FROM day09_lab
WHERE LABEL = 'orders_s3_20260816';

2. 是否仍需要独立 Broker 进程

Doris 已经内置 S3 和 HDFS 访问能力,这两类存储通常不需要额外部署 Broker 进程。自定义协议或特定生态接入仍可能依赖 Broker。旧课件中“Broker 负责所有 HDFS/S3 访问”的理解需要按集群版本和存储协议核验。

3. 现有 Broker Load 作业怎样迁移

迁移可以分成四步:

  1. 用 S3/HDFS TVF 对同一批文件执行只读查询,核验 Schema 和行数;
  2. 把 Broker Load 中的字段列表、SET、过滤条件改写为 SQL SELECT
  3. 使用 INSERT INTO ... SELECT FROM S3/HDFS(...) 写入影子表;
  4. 对比行数、聚合值、异常行、耗时、资源消耗和失败恢复,再逐批切换调度。

Broker Load 的异步、服务端重试和批任务状态是重要能力。TVF + INSERT 默认同步执行;需要后台运行时,应配合 Job 机制,不能只替换 SQL 文本后忽略调度与恢复语义。


六、S3/HDFS TVF + INSERT:先把远端文件当成表查询

TVF 预览、转换与落表
TVF 预览、转换与落表

TVF 的价值在于“可预览”。远端文件被暴露成一张临时关系表,可以先查看字段、抽样、过滤、转换、聚合,再决定是否写入目标表。

1. 先预览

SELECT
    order_id,
    user_id,
    order_time,
    amount,
    status
FROM S3
(
    "uri" = "s3://example-bucket/orders/dt=2026-08-16/*.parquet",
    "format" = "parquet",
    "provider" = "S3",
    "s3.endpoint" = "https://s3.us-west-2.amazonaws.com",
    "s3.region" = "us-west-2",
    "s3.access_key" = "<from-secret-manager>",
    "s3.secret_key" = "<from-secret-manager>"
)
LIMIT 20;

预览阶段应确认:

  • 文件能否访问;
  • 自动推断类型是否符合预期;
  • 日期、金额和空值语义是否正确;
  • 文件路径是否包含预期分区;
  • 多文件 Schema 是否一致;
  • 行数和文件数是否符合上游清单。

2. 再转换与落表

INSERT INTO orders_detail
WITH LABEL orders_tvf_20260816
SELECT
    CAST(order_id AS BIGINT),
    CAST(user_id AS BIGINT),
    CAST(order_time AS DATETIME(3)),
    CAST(amount AS DECIMAL(18, 2)),
    UPPER(status)
FROM S3
(
    "uri" = "s3://example-bucket/orders/dt=2026-08-16/*.parquet",
    "format" = "parquet",
    "provider" = "S3",
    "s3.endpoint" = "https://s3.us-west-2.amazonaws.com",
    "s3.region" = "us-west-2",
    "s3.access_key" = "<from-secret-manager>",
    "s3.secret_key" = "<from-secret-manager>"
)
WHERE amount >= 0;

SELECT 可以完成列裁剪、类型转换、过滤、Join 和聚合。复杂类型文件、跨文件清洗和临时历史回填通常会从这条路径获得更强的可解释性。

3. 严格模式与容错

INSERT INTO 使用 enable_insert_strict 控制异常数据处理。设置为 false 后,S3/HDFS/LOCAL TVF 写入还可以使用 insert_max_filter_ratio 控制最大过滤比例。生产任务应当把参数、过滤行数和异常样本写入验收记录。

4. 同步与异步

普通 INSERT INTO ... SELECT 会同步返回结果。大任务可以提交为 Job,由 Doris 后台执行。异步化以后仍需处理:任务唯一标识、运行窗口、资源组、失败重试、历史状态和补数流程。


七、INSERT INTO:内部 ELT、测试写入和小批量入口

INSERT INTO 通过 MySQL 协议提交,前半段执行查询,末端使用 OlapTableSink 写入目标表。

1. INSERT INTO ... VALUES

它适合:

  • 功能测试;
  • 少量维表维护;
  • 人工修复少量记录;
  • 配合 Group Commit 的高频小写入;
  • JDBC PreparedStatement 的成批参数写入。

批量初始化和大文件导入应使用 Stream Load、TVF + INSERT 或其他批量路径。逐行执行 INSERT VALUES 会产生大量解析、计划、事务和版本开销。

2. INSERT INTO ... SELECT

它适合:

  • Doris 内部表到表转换;
  • Hive、Iceberg、JDBC Catalog 等外部表落地;
  • S3/HDFS/LOCAL TVF 文件导入;
  • 分层数仓的 DWD、DWS、ADS 构建;
  • 历史数据重算和分区回填。

3. 显式 Label

当前 4.x 支持为 INSERT 指定 Label:

INSERT INTO dws_sales_daily
WITH LABEL dws_sales_daily_20260816
SELECT ...;

Label 能让调度系统根据业务批次查询状态,并在不确定结果时使用同一标识重试。调度平台还应维护任务实例表,记录源数据范围、目标分区、Label、TxnId、行数、校验值和重试次数。


八、批次设计:在可见延迟、吞吐和存储健康之间求平衡

批次大小、延迟与吞吐的权衡
批次大小、延迟与吞吐的权衡

导入性能往往受批次设计影响。每个常规导入都是事务,会产生元数据提交、MemTable Flush、Rowset 和 Segment。批次越小,数据越快可见,单位数据承担的固定成本越高。

1. 客户端可控时,优先合批

官方最佳实践建议把客户端可控的写入积累到数百 MB 至数 GB,再提交一次导入。具体数值需要结合 SLA:

  • 秒级看板可以采用 1–10 秒微批;
  • 分钟级报表可以形成更大批次;
  • 日批历史导入按分区和文件组拆分;
  • 单个远端大批任务建议控制在 100 GB 以内,降低失败重试成本。

培训材料给出了 100 MB–1 GB 或 10–60 秒一批的课堂建议,可以作为 POC 起点。真实生产值应通过压测确定。

2. 客户端无法合批时,使用 Group Commit

大量独立微服务、IoT 设备或采集进程很难共享客户端缓冲区。Group Commit 让 BE 在服务端合并请求,减少事务和版本数量。它降低了服务端压力,也会引入一个可配置的可见等待窗口。

3. 控制单批覆盖的分区数

每个被写入的 Tablet 都可能激活 MemTable。一批数据横跨几百个分区时,MemTable 会高度分散,提前 Flush 和小文件数量都会增加。历史数据应按天或按月顺序回灌,每批集中到少量分区。

4. 文件切分决定并行度

Parquet、ORC 和压缩文件的并行度通常受文件数量影响。一个超大文件只能被有限任务读取,多个大小均衡的文件更容易让多个 BE 并行工作。未压缩 CSV、JSON 在部分路径中可以内部切分,仍需通过 Profile 和任务详情确认实际并行度。

5. 表设计会反向影响导入

  • Duplicate Key 路径最短,适合原始事实和日志;
  • Aggregate Key 需要按 Key 聚合;
  • Unique Key MoW 需要主键索引、历史行查找和 Delete Bitmap;
  • 倒排索引、Bitmap 索引等在写入时同步维护,会占用 CPU 与存储;
  • 桶数过多会让一批数据分散到大量 Tablet;
  • 随机分桶场景可评估 load_to_single_tablet=true,降低分发和小文件开销。

6. 并发需要边界

当前官方最佳实践建议,Stream Load 单 BE 并发优先控制在 128 以内,并明确提示不能超过相关线程池硬限制。并发提升到一定程度后,WebServer 线程、解析 CPU、网络、MemTable 内存和磁盘 Flush 会先后成为瓶颈。压测应逐级增加并发,观察吞吐是否继续增长,以及 P99、失败率和 Compaction Score 是否恶化。

7. 建立导入运行看板

导入看板至少覆盖四组指标。任务层关注成功率、失败率、运行时长、排队时间、可见延迟和重试次数;数据层关注 Total、Loaded、Filtered、Bytes、每秒行数和每秒字节数;资源层关注 FE 事务与元数据压力、BE CPU、内存、磁盘、网络、活跃 MemTable 和远端存储 QPS;存储健康层关注 Tablet Version Count、Segment 数、Compaction Score、磁盘水位和副本健康。

指标需要按照数据库、表、数据源、任务类型和业务批次聚合。平均值容易掩盖少数异常任务,运行时长和可见延迟应同时展示 P50、P95、P99。上线门禁还要设置趋势告警:吞吐未下降时,Version Count 和 Compaction Score 持续上升,通常意味着存储层正在积累技术债务。


九、数据质量与幂等:把“写成功”定义清楚

数据质量、Label 与对账闭环
数据质量、Label 与对账闭环

生产导入需要一套统一成功标准。

1. 技术成功

  • 请求或异步任务进入成功状态;
  • 事务已提交并达到 VISIBLE
  • 目标副本写入达到系统要求;
  • 没有未处理的 ErrorURL 或 ErrorMsg。

2. 数据成功

  • 源行数、加载行数、过滤行数可解释;
  • 目标分区行数符合预期;
  • 金额、数量、最大最小时间等关键聚合一致;
  • Unique Key 表的最终状态和 Sequence 规则正确;
  • Aggregate Key 表的聚合粒度与指标口径正确;
  • 字符编码、时区、NULL 和 Decimal 精度没有漂移。

3. 业务成功

  • 数据在承诺时间内可查;
  • 下游报表或任务已完成验证;
  • 批次状态、Label、源范围和目标范围已经登记;
  • 重试不会产生重复;
  • 失败批次有明确的补数和回滚路径。

4. Strict Mode 与过滤比例

当前 Stream Load 的 strict_mode 默认值为 falsemax_filter_ratio 默认值为 0。严格模式开启后,非空源值在类型转换后变成 NULL 的记录会被过滤;严格模式关闭时,转换失败可能写成 NULL,目标列不允许 NULL 时该行仍会被过滤。max_filter_ratio 控制过滤行占“过滤行 + 成功加载行”的最大比例,前后置条件产生的 Unselected Rows 不计入该比例。放宽比例前,需要回答:

  • 哪些错误可以丢弃;
  • 丢弃比例如何告警;
  • 异常行保存在哪里;
  • 何时触发人工阻断;
  • 修复后怎样补回;
  • 对账口径是否排除了过滤行。

直接把比例调大容易掩盖上游质量问题。ErrorURL 中的错误样本应进入质量平台或隔离表,形成可追踪闭环。

5. Exactly-Once 的边界

Label 能保证同一数据库内、Label 保留窗口中的同批事务只成功一次。上游若提供 At-Least-Once,配合稳定 Label 可以实现批次级去重。跨系统 Checkpoint 的 Exactly-Once 通常需要 Stream Load 2PC 或 Connector 的事务协调。

Group Commit 会使用系统生成的合并 Label,并绕过用户自定义 Label 与 2PC 场景。选择 Group Commit 时,需要接受它的事务边界和可见性语义。


十、故障排查:先定位阶段,再决定重试方式

Doris 导入故障定位树
Doris 导入故障定位树

1. Label Already Exists

这通常表示同一 Label 已经成功或正在执行。先查询状态:

  • 已成功:把本次重试视为幂等命中;
  • 仍在运行:等待或检查超时;
  • 已失败:确认该 Label 是否允许重用,再执行重试。

2. Too many filtered rows

常见原因包括字段类型不匹配、目标列非空、字符串超长、分区不匹配、日期格式错误和列顺序错误。优先打开 ErrorURL,查看具体错误行和原因,再调整 Schema、列映射或源数据。

3. -235 Too many segments 或版本堆积

高频小批、过多分区、桶数过多和 Compaction 跟不上都可能触发。处理顺序通常是:

  1. 增大客户端批次;
  2. 启用 Group Commit;
  3. 减少单批覆盖的分区;
  4. 检查桶数与 Tablet 数;
  5. 观察 Compaction Score、磁盘和 CPU;
  6. 在明确瓶颈后再调整 Compaction 参数。

4. 导入超时或长时间不可见

区分读取慢、解析慢、Shuffle 慢、Flush 慢和 Publish 慢。检查:

  • 远端存储带宽和 QPS;
  • 文件是否过大且缺少并行;
  • BE CPU、内存、磁盘和网络;
  • 副本是否健康;
  • 目标分区是否过多;
  • 是否存在大规模 Compaction 或 Schema Change;
  • SHOW LOAD、返回 JSON 和 BE 日志中的阶段耗时。

5. S3/HDFS 访问失败

核验 Endpoint、Region、URI、AK/SK 权限、临时凭据有效期、网络出口、DNS、Kerberos 和服务端时间。对象存储权限应限制到必要 Bucket 和 Prefix,任务日志不得输出明文密钥。

6. BE OOM 或频繁提前 Flush

重点观察一批数据覆盖的 Tablet 数、并发任务数、单行大小、宽表列数、索引构建开销和 MemTable 内存。历史回灌按分区串行或有限并行,通常比全表混合导入更稳定。


十一、导入方式选择矩阵

场景 推荐入口 关键理由 主要注意点
本地文件、应用批次、需要同步结果 Stream Load HTTP 简单、Label 幂等、批次原子提交 控制单文件与并发,解析返回 JSON
高频小批 JDBC / Stream Load Group Commit 服务端合批,降低事务与版本压力 模式、WAL、可见延迟和回退条件
S3/HDFS 文件,需要预览与转换 TVF + INSERT INTO SELECT SQL 可读、可过滤、可转换 默认同步,大任务配合 Job
Doris 内部或 Catalog 表 ELT INSERT INTO SELECT 与 SQL 计算链路统一 资源隔离、Label、分区回填
现有 4.x 大规模远端异步作业 Broker Load 服务端异步和任务状态成熟 已弃用,新建链路规划迁移
Kafka 持续消费 Routine Load / Kafka Connector 管理 Offset、重启与长期任务 Day 10 详细展开
MySQL/PG 全量与增量同步 Flink CDC / Streaming Job Checkpoint、Schema 演进、CDC 语义 Day 10 详细展开

选择顺序可以概括为:先识别数据源和持续性,再确定同步或异步,随后确定批次和事务语义,最后处理格式、转换、质量和运维。


十二、本日实战:完成四条导入路径和一次失败演练

Day 9 实验与验收闭环
Day 9 实验与验收闭环

完整包提供实验 SQL、数据生成器和 Shell 脚本。建议按下面的顺序执行。

任务 1:创建实验表

创建:

  • orders_detail:Duplicate Key 订单明细;
  • events_gc:用于 Group Commit 的事件表;
  • orders_stage:用于 INSERT INTO SELECT 的中间表。

任务 2:生成样例数据

运行:

python3 scripts/generate_orders.py --rows 100000

脚本使用固定随机种子生成 CSV,并额外生成一份包含非法金额和非法日期的错误样本。

任务 3:执行正常 Stream Load

cp .env.example .env
# 修改连接信息后
bash scripts/stream-load-orders.sh

验收:

  • Status = Success
  • 加载行数为 100,000;
  • 过滤行数为 0;
  • 目标表行数一致;
  • Label 和 TxnId 已保存。

任务 4:使用同一 Label 重试

再次执行同一脚本,确认没有新增重复数据,并记录返回状态。随后更换 Label,验证 Duplicate Key 表会新增一批明细。

任务 5:执行错误数据导入

分别测试:

  • max_filter_ratio = 0
  • max_filter_ratio = 0.05
  • strict_mode = true
  • strict_mode = false

记录任务状态、过滤行数、ErrorURL 和最终目标行数,写出你所在团队允许的质量门限。

任务 6:验证 Group Commit

开启 async_mode,连续执行多次小 INSERT,观察返回的 Label、TxnId 和 PREPARE 状态;等待刷新后查询数据。再切换 sync_mode,比较响应时间与立即可见性。

任务 7:完成 INSERT INTO SELECT

orders_stage 清洗金额和状态后写入 orders_detail,使用显式 Label,并通过 SHOW LOAD 查看历史任务。

任务 8:阅读远端导入模板

根据你的环境补全 S3/HDFS TVF 和 Broker Load 模板。没有远端存储时,只进行 SQL 评审,避免使用无效凭据反复请求公共 Endpoint。

本日交付物

输出一份《Doris 数据导入方案与验收单》,至少包含:

  1. 数据源、格式、规模和增长速度;
  2. 可见延迟和吞吐目标;
  3. 导入方式及选择依据;
  4. 批次大小、时间窗口和并发;
  5. Label 生成规则;
  6. 严格模式和过滤比例;
  7. 异常数据留存与补数方式;
  8. 任务监控、告警和重试策略;
  9. 行数、聚合值和业务结果对账;
  10. 版本堆积、Compaction 和容量监控;
  11. 故障恢复与回滚;
  12. Broker Load 存量作业的迁移计划。

十三、Knowledge Check

1. Stream Load 为什么需要 Label?

用于唯一标识批次,在网络超时和客户端重试时提供事务去重依据,并支持按业务批次查询任务状态。

2. HTTP 200 是否等于导入成功?

不能直接等同。客户端还需要解析返回 JSON 中的 Status、加载行数、过滤行数、ErrorURL 和事务状态。

3. 高频小批次为什么会伤害集群?

每批都会产生事务、元数据记录、MemTable Flush、Rowset 和 Segment。版本数量和小文件持续增长后,Compaction、磁盘 I/O 和 FE 元数据压力都会上升。

4. sync_modeasync_mode 的主要差别是什么?

sync_mode 等待合并事务提交并可见后返回;async_mode 在写入接收 BE 的 WAL 后返回 PREPARE,数据稍后可见。

5. 哪些请求会绕过 Group Commit?

显式事务、自定义 Label、Stream Load 2PC、部分列更新和部分不兼容的 VALUES 表达式等场景会走常规路径。

6. 当前新建 S3/HDFS 批量导入链路优先采用什么方式?

优先评估 S3/HDFS TVF 或 Catalog 表配合 INSERT INTO ... SELECT。Broker Load 在 4.x 中仍可用于存量兼容,官方已标记弃用。

7. TVF 的最大工程价值是什么?

远端文件可以先作为关系表查询,完成 Schema 验证、抽样、过滤和转换,再写入目标表,导入逻辑更容易审查和复用。

8. 为什么历史回灌要按分区分批?

可以减少同时激活的 Tablet 和 MemTable,降低提前 Flush、小文件、内存和 Compaction 压力,也能缩小失败重试范围。

9. ErrorURL 应怎样处理?

下载错误样本,解析失败原因,进入质量隔离与告警流程;修复源数据或列映射后补数,并保留批次审计记录。

10. 一次导入验收至少要核对哪些数字?

源行数、总行数、加载行数、过滤行数、目标行数、关键金额或数量汇总、最小最大时间,以及 Unique/Aggregate 模型对应的最终语义。


十四、本日小结

今天建立了一套完整的 Doris 批量导入方法:

  • Stream Load 提供同步 HTTP 批次写入、Label 幂等和原子提交;
  • Group Commit 在服务端合并高频小写入,降低事务、版本和 Compaction 压力;
  • S3/HDFS TVF 让远端文件先进入 SQL 视野,再通过 INSERT INTO SELECT 完成清洗与落表;
  • INSERT INTO 统一内部 ELT、Catalog 数据落地和少量写入;
  • Broker Load 仍服务于部分 4.x 存量异步任务,新建链路需要规划向 TVF、Catalog 与 Job 迁移;
  • 批次、Label、严格模式、异常行、对账、监控和重试共同构成生产导入合同。

Day 10 将进入实时同步体系,系统讲解 Routine Load、Kafka、Flink Doris Connector、Flink CDC、Streaming Job、Offset、Checkpoint、Sequence Column 与 Schema 演进。


官方资料索引