简介:这是一份由Doris官方提供的Kettle-Spoon数据抽取插件doris-stream-loader,面向使用Kettle进行ETL开发、需要将数据实时高效写入Doris分析型数据库的大数据工程师。插件内含3个jar文件与1个xml文件,压缩包仅502KB,jar包涵盖插件核心实现与界面模块,xml则用于插件版本及元信息配置,结构简洁,无需额外依赖即可部署。使用前需确认Kettle版本为9.4.0.0-343,解压后放入data-integration\plugins目录并重启Spoon,即可在“转换”的“批量加载”中找到该插件。配置时重点填写Fenodes节点地址(格式为ip:http_port,默认8030)以及数据库连接信息,并注意表字段大小写需与流字段保持一致,以避免数据映射异常。目前已有1506人学习下载。借助该插件,用户可跳过繁琐的自定义开发,简化Kettle与Doris之间的数据通道,显著提升大批量数据导入效率,适合需要稳定、高吞吐数据同步的实时分析场景。 做了几年数据仓库,每天打交道最多的就是两件事:从各种库里把数据抽出来,再灌进目标库。Doris 的查询性能我是很服的,但“怎么快速把数据喂给 Doris”这个问题,早期真没少让我头疼。后来换成了 doris 官方提供的 kettle-spoon 插件 doris-stream-loader,数据抽取效率一下子提上来了,这篇文章就把我的踩坑过程和完整配置思路写出来,给同样用 Kettle Spoon 做抽取同步的朋友做个参考。
Kettle Spoon 在 ETL 工具里算是团队必装了,图形化拖拽、免费开源、插件生态全,但默认没有面向 Doris 的高效写入组件。如果图省事直接配一个 JDBC 输出,把 Kettle 的每个 RowSet 都转成 insert 语句执行,数据量一上来就是灾难:一条事务一条 commit,网络来回还被 JDBC 驱动限制得死死的。而 Doris 官方推荐的数据导入方式,恰恰是走 HTTP 协议的 Stream Load。doris-stream-loader 插件做的事情,就是把这个 Stream Load 的能力封装成 Kettle 的一个输出步骤,让普通 ETL 工程师不用写 Java、不用调 REST API,在 Spoon 界面里配几个参数就能享受到批量化流式写入的吞吐优势。
这个方案适合谁?适合已经在用 Kettle 做离线数仓同步、又想把目标表切到 Doris 的团队。也适合想统一团队成员技术栈、不想每个任务都去单独写 DataX 脚本的项目组。下面的内容我会按照“方案思路、核心机制、实操配置、常见问题”这条线展开,尽量把文档里不会写明白的细节也一并交代清楚。
1. 整体设计与思路拆解
1.1 为什么数据抽取链路里需要 Doris 官方插件
先从业务背景聊起。Doris 是个 MPP 架构的分析型数据库,它的强项是海量数据下的高并发查询,所以越来越多报表平台、用户行为分析、日志分析系统把 Doris 作为明细层和汇总层的存储。可数据进来之前,所有团队都要面对“抽取”这个老问题:源端可能是 MySQL、Oracle、SQL Server,也可能是 CSV、Kafka 里的日志。传统做法是先用 Kettle 把数据从源端读出来,做清洗、关联、去重、字段映射,最后写入目标端。
如果目标端是 MySQL,Kettle 本身有“表输出”步骤,用 JDBC 批量插入性能尚可。但目标端换成 Doris 以后,情况就变了。Doris 的 JDBC 驱动在很多版本里还承担着“给外部查询和少量写入”的任务,拿来大批量灌数并不理想。更合理的做法是走 Doris 的导入通道。Doris 支持多种导入方式,包括 Broker Load、Spark Load、Stream Load、Routine Load。Stream Load 是最直接的一种,它让客户端通过 HTTP 发送一批数据,Doris 的 FE 节点接收请求、解析元数据,然后把数据按照分桶规则分发给 BE 节点执行写入。doris-stream-loader 插件把这个过程原封不动地封装成了 Kettle 的 Output Step,等于是把 Doris 最擅长的导入路径,直接接到了 Kettle 最通用的图形化流程里。
1.2 与其它数据导入方案的取舍对比
很多人会问:既然 Stream Load 这么好,为什么不用 DataX,或者自己写个脚本 curl 发请求,偏要装 Kettle 插件?我的结论是,取决于团队的使用习惯和任务编排现状。
| 方案 | 使用方式 | 吞吐性能 | 维护成本 | 适用场景 |
|---|---|---|---|---|
| JDBC 批量写入 | Kettle 自带表输出步骤 | 低中 | 低 | 测试、几百行小表 |
| DataX + doriswriter | Python/JSON 配置任务 | 高 | 中,需要额外部署 DataX 服务 | 离线批量同步,独立调度 |
| Spark Load | 通过 Spark 任务导入 | 很高 | 高,需要维护 Spark 集群 | 超大规模离线任务 |
| doris-stream-loader | Kettle 插件步骤 | 高 | 低,嵌入Kettle即可 | Kettle ETL 流程里直接同步到 Doris |
从这个表能看出来,doris-stream-loader 最核心的竞争力是“不改变团队现有的 ETL 架构”。你不需要再起一台 DataX 服务,也不用让分析师去学一套新的配置语法,已经在 Kettle 上跑得好好的清洗逻辑、转换逻辑可以原样保留,只需要把最后的输出步骤换成“Doris Stream Loader”,数据抽取效率就能立刻获得数量级的提升。
这里也想提一下它解决的一个真实痛点:以前我们团队既要维护 Kettle 作业,又要维护 DataX 脚本,两套任务各自跑各自的,字段映射对不上、血缘关系理不清的问题经常发生。统一到 Kettle 插件以后,整条链路都在同一个工具里,运维成本确实降了不少。
2. 核心细节解析与实操要点
2.1 doris-stream-loader 的工作机制
任何工具用之前,最好先把原理看明白,否则遇到问题就只能瞎猜。doris-stream-loader 的实际工作流程是这样的:Kettle 的数据流会按行进入这个输出步骤,插件先把行数据在内存里攒成一个小批次,然后按照你配置的分隔符组装成 CSV 格式的文本,再把这些文本作为 HTTP 请求的 body,发送给 Doris 的 FE 节点的 Stream Load 接口。
这里最重要的一个设计是“流式上传”。既然发送的是 HTTP body,那么理论上可以不等待整个批次全部生成完毕再发,而是边生成边通过 chunked 编码把数据推过去。这样一来,Kettle 这端的内存占用是可控的,不会因为几千万行数据积压在内存里就 OOM;Doris 那端接收到 HTTP 后,会把它当成一个完整的导入 Label 任务,由 BE 节点接收并写入。由于 BE 节点在写入时会复用自身的 MemTable 结构,批量提交的数据可以直接进入列式存储的合并流程,比一条条 insert 再走事务提交要快得多。
为啥它比 JDBC 快?一句话总结就是:JDBC 是“每次来一条写一条”,Stream Load 是“攒一批后整包寄出去”。网络往返次数少了,导入引擎的批量优化也发挥出来了。实际测试中,百万行级别数据量下,JDBC 可能要跑十几分钟,Stream Load 通常几十秒就能完成。
2.2 插件安装与核心参数说明
安装这个插件不算复杂,但版本问题很容易踩坑。doris-stream-loader 目前在 GitHub 上有独立仓库,也提供了编译好的 jar 包。拿到 jar 包后,放进 Kettle 安装目录下的plugins/steps子目录里,比如>CREATE TABLE dwd.user_amount ( id INT, user_name VARCHAR(50), city VARCHAR(20), amount DECIMAL(12,2), create_time DATETIME ) DUPLICATE KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 10 PROPERTIES ("replication_num" = "1");
注意如果是测试环境,replication_num可以设为 1 节省资源;生产环境建议至少 3 保证副本安全。这里用 DUPLICATE KEY 模型是因为很多明细数据不需要更新,只是追加,Doris 对这类数据的导入效率是最高的。
3.2 在 Spoon 里配置一个最简单的同步转换
打开 Kettle,新建一个转换,添加一个“表输入”步骤作为数据源。假设源表是 MySQL 的一张业务表,查询语句可以这样写:
SELECT id, user_name, city, amount, create_time FROM mysql_source.biz_user_amount WHERE create_time >= '2025-01-01 00:00:00'然后添加“Doris Stream Loader”输出步骤,将两个步骤连接起来。双击输出步骤,按上面参数表配置:FE 节点、库名、表名、用户名密码、目标列、分隔符。比如我习惯配成:
- FE nodes:
192.168.1.10:8030 - Database:
dwd - Table:
user_amount - Columns:
id,user_name,city,amount,create_time - Separator:
\t - Label Prefix:
kettle_user_amount_ - Stream Load Properties:
{"format":"csv","column_separator":"\\t","max_filter_ratio":"0.1","timeout":"300"}
这里有个细节,如果插件界面上有独立的“Separator”字段,又在“Stream Load Properties”里写了column_separator,两者可能会互相覆盖。建议只在一个地方配置,避免冲突。
接下来保存转换,点击运行。如果一切顺利,日志里会出现类似Stream load success. NumberTotalRows: 5000000, NumberLoadedRows: 5000000, LoadBytes: 102400000的信息,说明这批数据已经成功导入。
3.3 验证数据并对比抽取性能
导入完成后,去 Doris 里执行一条查询确认数据是否正确:
SELECT city, COUNT(*), SUM(amount) FROM dwd.user_amount GROUP BY city;顺便说一下我做过的对比测试。同样是从 MySQL 抽 500 万行明细数据到 Doris,用 Kettle 自带的表输出步骤(JDBC 方式)跑,耗时大概 12 分钟;换成 doris-stream-loader 之后,全程用了 1 分 40 秒。注意这里的对比不是黑 JDBC,而是说明在 Doris 的场景里,JDBC 那条路确实不适合批量导入。如果你在处理几百 GB 级别的表,这个差异还会更明显。
还有一个实用做法,就是让 Kettle 的输入 SQL 带上分片条件,比如:
SELECT * FROM biz_user_amount WHERE id > ? AND id <= ?然后在 Kettle 里用“复制发送到结果”或自定义循环实现多个分片并行读取,让下游的 Doris Stream Load 步骤也能跟着并行运行,整体吞吐直接成倍提升。不过并发不要开太高,否则会给 FE 节点造成过大压力,导入任务可能会因为 HTTP 连接数超限而报错。
4. 常见问题与排查技巧实录
4.1 我遇到过的典型报错和处理方式
只有踩过坑,才能把配置记得牢。我把自己实际遇到过的几个高频问题整理成了一张速查表,每一条都给出了排查方向和解决建议。
| 现象 | 可能原因 | 排查与解决 |
|---|---|---|
日志提示Connection refused | FE HTTP 端口写错,或者网络不通 | 确认 FE 的http_port默认是 8030,用curl http://fe_ip:8030/api/health测试连通性 |
导入后提示Label Already Exists | Label 重复,Kettle 重跑时没有生成新 Label | 检查 Label Prefix 是否带了时间戳/随机数,或者去 Doris 执行SHOW LOAD WHERE LABEL LIKE '前缀%'清理旧任务 |
提示errCode = 2, [217] The, response is not [OK] | 导入数据与目标表 schema 不匹配 | 检查列数、分隔符、Columns 顺序,重点看是否有空串、时间格式异常 |
大量行被 filter,NumberLoadedRows远小于NumberTotalRows | max_filter_ratio实际为 0,数据质量问题 | 在 Stream Load Properties 里设置"max_filter_ratio":"0.1",先允许 10% 脏数据,再针对性查日志看 filter 原因 |
报错Time Out | 批次数据量太大,超过 Stream Load 超时时间 | 调大timeout字段(单位秒),或者减小 Kettle 每次发送的行数 |
注意:Doris 的 Stream Load 对时间格式要求比较严格。如果你源库里的日期是
2025/01/01 10:00:00,而 Doris 目标列是 DATETIME,最好在 Kettle 里先通过“字段选择”或“字符串操作”把它统一转成yyyy-MM-dd HH:mm:ss,否则很容易被 filter。
4.2 效率调优与避坑心得
插件本身效率高,但用不好也容易被周边环节拖慢。第一个要调的是 Kettle 的“每一批提交的行数”,有些版本翻译叫“Commit size”。如果设置得太小,比如默认 1000 行就发一次 HTTP,那么 500 万行数据要发 5000 个请求,再快也经不住这种网络开销。我一般会调到 50000 到 200000 之间。行数增加后,内存会有一定上涨,你可以通过调整 Kettle 的-Xmx参数来给足 JVM 空间。
第二个是并发度。doris-stream-loader 作为 Kettle 的步骤,是支持“复制步骤”并行执行的。你可以在输出步骤上右键选择“改变开始复制的数量”,把并发提到 4 或 8。但 Doris 侧同一时间也不适合有太多 Stream Load 任务并行,尤其是导入大批量任务时,BE 节点的 CPU、磁盘 IO 都可能被打满。稳妥的做法是先单独跑一次单并发任务,记录耗时,再慢慢往上加并发,观察 Doris 主机的负载,找到一个不触发瓶颈的平衡点。
第三个避坑点是,不要在同一个转换里同时连多个 Doris 输出步骤去写同一张表。你可能会想用两个输入流分别清洗后合并写入,但 Stream Load 的 Label 机制要求导入任务在 Doris 侧用唯一标记,两个输出步骤如果使用相同的前缀又没带随机后缀,很容易互相覆盖或报重复 Label。可以用“复制分发到多个输出步骤”这种方式,确保每个输出步骤的 Label 前缀不同。
还有一个我特别想提醒的细节:如果你用 doris-stream-loader 去写一个带有自动生成列或默认值的表,一定要在 Columns 里把源数据要写入的所有列写全,不要依赖 Doris 表字段默认值。因为 Stream Load 的 CSV 格式默认按列顺序解析,如果 Columns 漏了列,原始数据会对不上目标表,最后要么报错要么写入错位。
最后再分享一个我自己的使用习惯。每个导入作业的 Label Prefix 里,除了表名,我还会拼上这次调度的时间戳,比如kettle_user_amount_20250218_,这样即使任务因故障重跑,也不会碰到旧 Label。虽然 Doris 本身有 label 去重机制,能防重复导入,但设计一个好的命名规则,能让你在SHOW LOAD里排查问题时省下不少时间。
做数据抽取这件事,Doris 已经给出了很多通道,而 doris-stream-loader 最打动我的地方,是它把 Doris 的高性能导入能力和 Kettle 的易用性结合得很自然。只要别在版本兼容和参数细节上偷懒,按照上面这套方法和排查路径走,你的数据抽取效率大概率也能出现肉眼可见的提升。
本文还有配套的精品资源,点击获取