10 分钟打通 Doris 数据导入:Stream Load 接口与多语言 SDK 实战指南
【免费下载链接】dorisApache Doris is a real-time analytics and hybrid search database for AI agents.项目地址: https://gitcode.com/GitHub_Trending/doris/doris
上游有 Kafka 消费组、有 Flink 作业、还有定时任务在产出 CSV,数据散在五六套系统里。你要做的事很朴素:给它们一个统一入口,把数据灌进 Doris 做实时分析。Stream Load 就是这个入口——它是 Doris 的 RESTful 数据加载接口,客户端发一个 HTTP 请求,数据就变成一次可追踪的导入事务。
Stream Load 在数据链路中的位置
一句话讲清心智模型:你只跟 FE 的 HTTP 端口说话,数据实际落到 BE 上。
流程是:客户端PUT到 FE 的/api/{db}/{table}/_stream_load(默认 8030 端口,见 conf/fe.conf 的http_port = 8030);FE 校验后把请求转发给某个 BE,BE 解析、写入并发布版本,最后 FE 把 BE 的导入结果 JSON 原样回给你。仓库里这条路由注册在 be/src/service/http_service.cpp,只注册了HttpMethod::PUT。
为什么用 PUT 而不是 POST?因为这次请求的语义是"把这份数据写到这张表",PUT 表达的是对资源的完整写入,而不是提交一个表单动作。对你唯一的实际影响是:客户端必须显式指定 PUT 方法,浏览器地址栏是试不出来的。
最小可用请求
先跑通,再谈参数。下面这段 Python 就是完整可运行的最小示例,向db0.t_user写两行 CSV:
import requests from requests.auth import HTTPBasicAuth url = 'http://127.0.0.1:8030/api/db0/t_user/_stream_load' headers = { 'Content-Type': 'text/plain; charset=UTF-8', 'format': 'csv', 'column_separator': ',', 'Expect': '100-continue', } resp = requests.put(url, headers=headers, data='1,Tom\n2,Jelly', auth=HTTPBasicAuth('root', '')) print(resp.status_code, resp.text)仓库里的完整版本在 samples/stream_load/python/DorisStreamLoad.py,注意它把should_strip_auth关了——因为 FE 可能 307 跳转到 BE,跳走后默认库会丢掉 Authorization 头。
必需 Header 只有三个:
| Header | 作用 |
|---|---|
Content-Type | 声明请求体是文本流(数据本体,不是表单字段) |
format | 告诉 BE 按什么格式解析,csv或json |
Expect: 100-continue | 客户端先问"你敢收吗",BE 确认后才传数据体,避免白传 |
成功时响应里Status为Success,并带NumberLoadedRows、TxnId等字段。注意:HTTP 200 只说明链路通了,不代表导入成功,必须看 JSON 里的Status与Message,Java 示例 samples/stream_load/java/DorisStreamLoad.java 的注释里就专门强调了这一点。
分场景加参数
最小请求之外,参数按场景补即可,不用背全表。
上游列顺序和表不一致——加columns指定映射,例如columns: id,name。JSON 场景还能在columns里写表达式(bucket=floor(member_id/1000000)),这是 Go 示例 samples/stream_load/go/doris_stream_load.go 的用法:
req.Header.Add("columns", fmt.Sprintf("member_id,bucket=floor(member_id/1000000),members=to_bitmap(member_id)")) req.Header.Add("format", "json") req.Header.Add("Label", fmt.Sprintf("crowd_%d_%d_%d", crowd, logID, num))JSON 是对象数组——加strip_outer_array: true剥掉最外层[ ];嵌套字段用jsonpaths提取。
失败容忍度——脏数据行默认会整单回滚。可以接受少量过滤时,配合max_filter_ratio设过滤比例上限(BE 侧配置项),超了照样失败。
幂等控制——label是导入任务的唯一标识。同一个 label 重复提交会直接返回Label Already Exists且不会二次写入(Java 示例注释里有这个响应的完整样例)。这是后面"断点续传"的地基。
多语言 SDK 选型对照
Doris 官方没有独立的语言 SDK 包,社区在 samples/stream_load 目录维护了四种语言的完整参考实现,全部是裸 HTTP,依赖很轻:
| 语言 | 参考实现 | 一句话点评 |
|---|---|---|
| Python | DorisStreamLoad.py | requests几行搞定,写胶水脚本和数据验证首选 |
| Java | DorisStreamLoad.java | 基于 Apache HttpClient,要自己处理 307 跳转时保留认证头,Flink/Spark 集成常用 |
| Go | doris_stream_load.go | net/http标准库实现,天然带连接池,高吞吐采集器适合 |
| Rust | doris_stream_load.rs | reqwest+ tokio 异步发请求,内存安全,适合高性能数据管道 |
let response = client .put(url) .headers(headers) .body("1,Tom\n2,Jelly") .send().await?;选型逻辑很简单:跟着数据管道的既有语言走。采集组件是 Go 就用 Go,Java 生态里加个类就行,不需要跨语言起服务。
生产避坑清单
现象:偶发 307/401,数据时有时无原因:FE 把请求 307 转发到 BE 时,部分 HTTP 库默认会丢弃跨主机的 Authorization 头。 处置:让客户端在重定向后保留认证头。Python 里设session.should_strip_auth = lambda *a: False;Java 里用自定义DefaultRedirectStrategy允许 PUT 重定向,参考 samples/stream_load/java/DorisStreamLoad.java 第 119 行的写法。
现象:返回 HTTP 200,但表里没数据原因:200 只代表 BE 服务可达,导入成败在响应体 JSON 里。 处置:判断Status == "Success"且NumberFilteredRows是否符合预期,把Message打进日志。
现象:同一批数据被写了两次原因:超时后盲目重试,重试请求换了 label,等于新开了一次导入。 处置:重试时复用原 label。label 已被使用且状态FINISHED时,说明数据其实已写入,无需再发。
现象:CSV 导入报字段数不匹配原因:上游列顺序与表结构不一致,或行尾多了空字符。 处置:用columns显式指定列映射,导入前先在客户端做一次字段数校验。
可落地的进阶做法
- 断点续传:给每批数据生成确定性 label(如
bizid + 批次号 + 时间片),把"已成功的 label 集合"落到本地或 Redis。重跑时先查 label 状态,已完成的批次直接跳过。幂等性由 Doris 的 label 机制兜底,客户端不需要自己做去重存储逻辑。 - 批量任务跟踪:响应 JSON 里有
TxnId,导入结束后也可在 FE 用SHOW STREAM LOAD查任务状态。把(label, TxnId, 批次号, 提交时间)记成一张流水表,就能回答"哪批失败了、失败在哪",配合NumberFilteredRows还能发现"成功但数据被大量过滤"的隐性事故。 - 客户端预处理:大文件先分片再逐片提交,每片一个 label;提交前在客户端完成格式归一(时间戳格式、编码转换),把
max_filter_ratio当最后防线而不是日常手段。过滤率长期大于 0 的管道,问题一定在上游。
下一步:先把你最上游的那一路数据源(Kafka 消费端或定时任务)接到最小可用请求上,用固定 label 连跑三天,观察Label Already Exists的命中次数,再决定要不要上分片与流水表。
【免费下载链接】dorisApache Doris is a real-time analytics and hybrid search database for AI agents.项目地址: https://gitcode.com/GitHub_Trending/doris/doris
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考