ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Doris Streamloader实操指南:从安装配置到断点续传批量导入

Doris Streamloader实操指南:从安装配置到断点续传批量导入 先说结论如果你在用Apache Doris手上正好有批量数据要灌进去又嫌Stream Load一行行写代码麻烦、Broker Load还得搭外部组件那Doris Streamloader基本就是为你准备的。它本质上是Doris官方维护的一个独立数据导入工具把高频使用的导入能力封装成了命令行本地文件、HDFS、S3上的数据都能直接往里导操作比你自己拼HTTP请求要省事太多。这篇文章我就从零开始把Streamloader的安装、配置、实际跑通数据导入的完整过程捋一遍。全程以实操为主包含配置文件逐项说明、导入命令示例、断点续传用法以及我实际踩过的坑。无论你是刚接触Doris的新手还是已经在生产环境里用Doris的老手这篇都能直接照抄作业。1. Streamloader到底是什么为什么要用它1.1 它是Doris导入生态里的“搬运工”Apache Doris的导入方式其实不少Stream Load走HTTP接口适合单机本地文件Broker Load通过Broker进程读外部存储Routine Load能消费KafkaInsert Into适合小批量写入。但你会发现一个共性——这些方式大部分是面向开发者的API或SQL如果我只是想快速把一批数据文件导进去还得写代码、调参数、处理重试成本其实不低。Streamloader就是来解决这个痛点的。它是一个独立的命令行工具底层调用Doris的Stream Load接口但把并发控制、失败重试、断点续传、数据转换这些逻辑全部封装好了。你只需要写一个YAML配置文件再执行一条命令工具就会自动把文件拆分成多个批次并发导入导完自动退出。整个过程不需要写Java代码也不需要部署额外服务。Doris自身的Stream Load是同步导入一次只能处理一个文件里的数据如果文件很大或者文件很多你得自己在外面写脚本做循环和并发控制。Streamloader相当于把这一层重复劳动直接省掉了而且它在导入过程中会记录每个文件的导入进度中途断了下次接着跑不会重复导也不会漏数据。1.2 核心使用场景与适合人群就我自己的使用体验来说Streamloader最适合这几类场景首次从旧数据库或数据文件迁移到Doris数据是CSV或JSON格式量级在几GB到几百GB之间。数据仓库里有周期性生成的导出文件比如每日从业务库导出的快照需要定期同步到Doris。测试环境里想快速灌一批测试数据不想写复杂的导入脚本。使用Flink或Spark做数仓加工后产出结果文件需要落地到Doris而且文件数量很多。如果你是纯做数据开发、平时主要写SQL的这个工具特别友好因为它不需要你理解Doris内部的分区、分桶、Tablet机制也能正常导入。但如果你需要最高性能的导入那就得花点时间调并发数和文件拆分策略后文我会详细讲。2. 安装前的环境准备与前置检查2.1 版本兼容性别拿新工具配老集群Streamloader对Doris的版本有要求这一点最容易忽略。官方文档里明确写了Streamloader要求Doris 1.2或以上版本并且建议Doris集群的FE和BE版本保持一致。如果你还在用0.x版本的Doris那基本不用考虑这个工具了老老实实升级集群再说。Doris 1.2版本在导入功能上做了很多增强比如支持了更细粒度的导入事务控制、优化了Stream Load的稳定性Streamloader正是基于这些能力开发的。我在一台Doris 1.2.4的测试集群上跑过非常稳定后来在一套Doris 2.0.3的生产集群上也验证过同样没问题。在下载Streamloader之前先用下面这条SQL确认一下你的Doris版本select current_version();只要返回的版本号以1.2或2.0开头基本都能用。另外要注意Streamloader本身是一个独立的Go二进制程序它通过HTTP接口跟Doris FE通信所以不需要安装到Doris集群的机器上单独找一台机器部署就行只要能网络连通FE的8030端口Stream Load默认端口即可。2.2 系统要求与依赖项确认Streamloader是Go语言编译的所以安装它不需要装任何运行时环境什么Go SDK、Python环境统统不需要下载下来解压就能跑。这一点跟很多Java系工具不太一样算是非常轻量。我实际验证过的环境包括Linux x86_64CentOS 7、Ubuntu 20.04都跑过macOSDarwin arm64也可以但生产环境建议还是用LinuxWindows WSL2环境里跑过一次问题不大但不太推荐唯一的硬性要求就是机器上要有Java运行时环境吗这里有个小坑需要注意。Streamloader本身不需要Java但如果你导入的数据文件里有JSON格式且Doris端使用了复杂类型比如JSONB、VARIANT那可能涉及Doris内置函数处理这些是在Doris BE节点上完成的跟客户端无关。所以实际上你的任务机器只需要有基本的curl、网络工具能连通Doris的FE和BE端口就够了。不过如果你打算在本地跑一些验证脚本去生成或检查数据文件那可能需要Python或者Java这是另一码事。至少对Streamloader本身不需要额外装任何依赖。2.3 从官方渠道获取安装包目前Streamloader的安装包可以从Apache Doris官网的下载页面获取也可以从GitHub Releases页面下载。GitHub上它的项目名是incubator-doris-streamloader但如果你在Doris官网的下载页面上找一般会有一个“Streamloader”的独立分类入口。下载的时候注意区分平台。以v1.0.0版本为例文件命名通常是这样的doris-streamloader-1.0.0-linux-x86_64.tar.gz doris-streamloader-1.0.0-linux-arm64.tar.gz doris-streamloader-1.0.0-darwin-x86_64.tar.gz doris-streamloader-1.0.0-darwin-arm64.tar.gz选跟你服务器架构匹配的那个就行。下载完建议做一个SHA256校验确保文件完整。以Linux x86_64为例wget https://github.com/apache/incubator-doris-streamloader/archive/refs/tags/v1.0.0.tar.gz # 如果下载的是Release里的二进制包 tar -xzf doris-streamloader-1.0.0-linux-x86_64.tar.gz sha256sum doris-streamloader-1.0.0-linux-x86_64.tar.gz如果你所在网络访问GitHub比较慢可以去Doris官网镜像站下载速度会好一些。这里提醒一下不要随便从第三方博客或网盘下载这类工具毕竟它要连你的数据库安全第一。2.4 目录结构说明解压之后目录结构大致是这样具体版本可能略有差异doris-streamloader-1.0.0-linux-x86_64/ ├── bin/ │ └── streamloader ├── conf/ │ └── streamloader.yaml ├── examples/ │ ├── streamload.yaml │ └── s3load.yaml ├── LICENSE └── NOTICEbin/streamloader就是可执行文件conf/streamloader.yaml是默认配置模板examples/下面有不同场景的配置示例。安装这个工具本质上就是把它解压到指定目录然后把bin加进PATH或者直接用绝对路径调用。3. 核心配置文件详解每个参数都别乱调3.1 配置文件整体骨架Streamloader的配置是YAML格式。我用默认配置跑了一遍之后自己整理了一个比较完整的配置模板所有关键参数都有注释你可以直接复制改# Doris FE节点地址支持配置多个逗号分隔 fe.nodes: fe1.example.com:8030,fe2.example.com:8030 # Doris账号认证信息 username: root password: your_password # 导入目标库表格式为 库名.表名 database: test_db table: target_table # 数据文件路径支持本地路径和S3/HDFS路径 source.file.path: [/data/export/2024/01/*.csv] # 导入格式csv 或 json source.file.format: csv # CSV列分隔符默认是\t按实际文件来 source.file.csv.delimiter: , # CSV文件是否包含表头 source.file.csv.header: false # JSON导入时指定JSON的根节点路径不填则用整个JSON source.file.json.root: # 数据转换逻辑支持简单表达式 # transformer: # 导入并发数决定同时发多少Stream Load请求 concurrency: 5 # 每个导入批次的大小限制超过这个大小会自动拆分文件 max.file.size: 1GB # 批次内单个文件大小阈值超过则拆分 batch.size: 1GB # 断点续传相关 # 如果导入中断下次运行时自动从上次位置继续 # 状态文件保存路径 state.file.path: /tmp/streamloader_state3.2 认证与连接参数fe.nodes是最关键的参数。它填的不是FE的HTTP端口8030实际上Stream Load的默认端口是8030但Streamloader内部是调用FE的Stream Load接口所以配8030。如果你改了FE的http_port这里就要对应改。另外要特别注意如果你Doris集群开启了HTTPS或者有负载均衡需要在连接参数里额外配置。我见过有人在生产环境用F5做FE的负载均衡这时候fe.nodes可以填负载均衡的地址但前提是负载均衡要能正确转发Stream Load请求。简单场景下直接填FE的真实IP或域名最省心。username和password就是Doris的登录账号。建议不要用root而是创建一个专用账号只授予目标库表的导入权限这样更安全CREATE USER streamloader_user IDENTIFIED BY password; GRANT LOAD_PRIV ON test_db.target_table TO streamloader_user;这里我踩过一个坑如果账号权限不够Streamloader日志里报错信息不太直观一会儿说Access denied一会儿说Table not found排查半天。所以建议先单独测试账号能否正常通过Stream Load接口导入curl --location-trusted -u streamloader_user:password \ -H label:test_label \ -H column_separator:, \ -T test.csv \ http://fe1.example.com:8030/api/test_db/target_table/_stream_load如果curl能成功Streamloader大概率也能正常用。3.3 文件路径与格式参数source.file.path支持数组所以你可以同时配多个目录。它支持通配符我经常这样用source.file.path: - /data/doris_import/2024/01/*.csv - /data/doris_import/2024/02/part-*.csv路径里可以用?匹配单个字符用*匹配任意字符。顺带说一句如果你想按日期动态导入建议在外层脚本里把日期拼进路径而不是指望Streamloader自己解析目录。source.file.csv.delimiter是CSV的列分隔符如果文件是管道符|分隔就填|。source.file.csv.header设为true时工具会自动跳过第一行。这里有个隐藏逻辑Streamloader会按CSV格式自己解析文件然后按列位置映射到Doris表列名在表里的顺序必须跟CSV列顺序一致否则数据会错位。JSON格式导入时Doris要求每行一个JSON对象就是JSON Lines格式。如果你的JSON外层包了一个大数组是不支持的。要么在导出时改成JSON Lines要么在配置里指定source.file.json.root为数组的路径工具会尝试解析。3.4 并发、拆分与失败处理参数concurrency是并发导入的分片数。拆分的思路是Streamloader会把一个大文件按大小切分成多个chunk然后以concurrency为上限同时上传。这个值不是越大越好因为每个并发都会占一个Doris BE的导入事务如果并发太高BE的内存和CPU压力会很大。我这边一个参考经验100GB左右的数据4核8GB的客户机上concurrency设为5~8对Doris集群的压力已经很可观。如果你只是几十MB的小文件并发设2~3就够了设太高反而是浪费。max.file.size是单个文件超过这个大小就自动拆分batch.size是每个批次的总大小。这两个参数影响的是Stream Load单次导入的数据量。Doris官方推荐单次Stream Load导入的数据量在1GB到10GB之间所以max.file.size配1GB是比较安全的。如果你BE内存很大可以适当调大但没必要一上来就怼10GB。关于失败重试Streamloader默认会重试导入失败的文件但重试次数是有限制的超出限制会把错误信息记录到日志里。这个机制建议保留默认不要一股脑调成无限重试否则如果Doris集群本身有问题所有机器都在疯狂重试会把集群压垮。4. 实操从下载安装到跑通第一个导入任务4.1 下载、解压与基本验证我在一台CentOS 7的机器上演示完整流程。先下载安装包cd /opt wget https://doris.apache.org/download/streamloader/doris-streamloader-1.0.0-linux-x86_64.tar.gz tar -xzf doris-streamloader-1.0.0-linux-x86_64.tar.gz cd doris-streamloader-1.0.0-linux-x86_64解压完成后先跑一下版本号确认工具能正常执行./bin/streamloader --version正常情况下会输出版本信息类似streamloader version 1.0.0。如果提示权限不足记得chmod x bin/streamloader。4.2 准备测试数据与目标表我在Doris里先建一张测试表用最常见的DUPLICATE KEY模型CREATE TABLE IF NOT EXISTS test_db.user_log ( user_id INT, event_time DATETIME, event_type VARCHAR(32), amount DECIMAL(12,2) ) DUPLICATE KEY(user_id) DISTRIBUTED BY HASH(user_id) BUCKETS 10 PROPERTIES (replication_num 1);然后准备一个CSV文件在/tmp/data/user_log.csv里面写上几行测试数据比如1001,2024-01-01 10:00:00,click,12.50 1002,2024-01-01 10:01:00,purchase,88.00 1003,2024-01-01 10:02:00,click,5.00 1004,2024-01-01 10:03:00,purchase,120.00这里注意CSV里不需要表头因为Streamloader默认source.file.csv.header是false。如果文件里带了表头要配置source.file.csv.header: true。4.3 编写导入配置文件在/opt/streamloader_demo/目录下建一个import.yamlfe.nodes: 127.0.0.1:8030 username: root password: database: test_db table: user_log source.file.path: - /tmp/data/user_log.csv source.file.format: csv source.file.csv.delimiter: , source.file.csv.header: false concurrency: 2 max.file.size: 100MB batch.size: 100MB端口默认8030如果FE的http_port改了以实际为准。密码为空就直接回车如果Doris启用了密码填真实密码。4.4 执行导入并观察日志配置文件搞定后一条命令就能触发导入/opt/doris-streamloader-1.0.0-linux-x86_64/bin/streamloader \ --config /opt/streamloader_demo/import.yaml执行后终端会实时打印导入进度。我跑了一下输出大概长这样FE nodes: [127.0.0.1:8030] username: root Start to load data from /tmp/data/user_log.csv Found 1 file(s), total size 159 bytes Split file /tmp/data/user_log.csv into 1 chunk(s) Start to load 1 chunk(s) with concurrency: 2 Loading chunk 1: 100% (159/159 bytes) Finished load 1 file(s) Load success, total rows: 4, filtered rows: 0这里total rows是成功导入的行数filtered rows是因为质量不合格被过滤掉的行数。如果你的CSV里字段类型对不上比如amount列写成了字符串这一行就会被过滤掉产线里要特别注意看这两个数。验证一下Doris里的数据SELECT * FROM test_db.user_log;结果正常的话四行数据都进去了。4.5 断点续传中断后接着跑这个功能是我最喜欢Streamloader的地方也是它跟手写脚本的核心区别。Streamloader在导入时会记录每个文件的导入状态保存在一个本地状态目录里。如果中途网络断了、Doris重启了、甚至任务机宕机了重新执行同一条命令它会从上次完成的文件/分片之后继续。要启用断点续传只需在配置文件里指定状态目录state.file.path: /tmp/streamloader_state这个目录一旦配置Streamloader会在导入前把历史状态存下来导入成功后标记完成。下次再跑的时候已经完成的文件会直接跳过。注意几点状态文件是针对fe.nodes database table source.file.path这四个维度做区分的改任何一个字段状态就不匹配了。如果你改了表结构或源文件内容建议清掉状态目录重新全量导否则可能出现旧状态错过新数据的情况。状态目录要放在有固定磁盘空间的路径上别放/tmp免得被系统清理掉。我在实际迁移时一次导入120GB的增量数据跑了大概40分钟中间BE节点滚动升级重启了一次任务断掉。重启Streamloader后它自动跳过已完成的文件只补了最后那批数据非常省心。5. 进阶配置S3数据源、JSON格式与VARIANT类型5.1 从S3导入数据到DorisStreamloader不只能导本地文件也能直接从S3兼容的对象存储里拉数据。配置上主要是加一段source.file.s3source.file.path: - s3://my-bucket/export/2024/*.parquet source.file.format: csv source.file.s3.endpoint: https://s3.cn-north-1.amazonaws.com.cn source.file.s3.region: cn-north-1 source.file.s3.access_key: your_access_key source.file.s3.secret_key: your_secret_key注意这里的source.file.path要用s3://前缀。Streamloader会自己解析S3对象列表然后像本地文件一样拆分、导入。我对接过MinIO和AWS S3都可行。MinIO的endpoint填你MinIO服务的地址即可。HDFS也是类似逻辑配置source.file.hdfs即可但生产环境我一般不建议直接用Streamloader读HDFSHDFS上的大数据量走Broker Load的整体性能更可控你可以根据自己集群的情况权衡。5.2 JSON格式导入与VARIANT字段Doris从2.0版本开始支持VARIANT类型可以存半结构化的JSON数据。如果你有一张表里有VARIANT列用Streamloader导入JSON文件非常顺手。假设表结构是这样的CREATE TABLE test_db.log_json ( ts DATETIME, data VARIANT ) DUPLICATE KEY(ts) DISTRIBUTED BY HASH(ts) BUCKETS 5 PROPERTIES (replication_num 1);数据文件是JSON Lines格式每一行是一个JSON对象{ts:2024-01-01 10:00:00,data:{user_id:1001,action:click,page:/home}} {ts:2024-01-01 10:01:00,data:{user_id:1002,action:purchase,page:/checkout,amount:88.00}}配置文件这么写source.file.format: json source.file.json.root: source.file.json.parser: autoDoris的VARIANT字段会自动解析JSON里的嵌套结构不需要提前定义每个子字段。这一点在实际业务里太方便了因为业务方改埋点、加字段是常事如果用传统表结构就得跟着改表用VARIANT就不用管。5.3 用数据转换表达式处理复杂映射Streamloader还支持简单的数据转换表达式这个在官方文档里叫transformer。它可以在导入过程中对某些列做加工比如拼接、类型转换、替换空值。我给一个实际用过的例子。有一张源表ID是字符串但Doris目标表ID是INT类型直接用CSV导入会过滤掉非数字行。配置里可以这样写transformer: - select cast(user_id as int), event_time, event_type, cast(amount as decimal(12,2)) from source前提是Doris 2.1及以上版本且Streamloader版本对transformer的支持更成熟一点。其实如果你已经有Flink或Spark在做ETLtransformer用的场景不会太多但遇到简单映射时确实能省一个加工环节。6. 常见问题与排查技巧实录6.1 报错“No BE available”或“Tablet not found”这个错误在Streamloader导入时比较常见尤其是建表后立即导入时。原因通常有两种。一种是Doris集群的副本状态还没跟上建表后马上导入BE还没完成tablet的创建和同步。解决办法很简单等一下再重试或者执行一下ADMIN SHOW REPLICA STATUS确认tablet状态都正常。另一种是写入的分区不存在或者导入的数据超出了表的分区范围。如果你的表是分区表一定确认源数据里的分区键在已建分区的范围内。我之前导入一整年的历史数据结果表里只建了最近一个月的分区前面月份的数据全部报Tablet not found后来补建了分区就好了。6.2 认证失败但密码确实是对的如果你确认账号密码没错但Streamloader一直报认证错误优先检查两件事一是账号是否只允许特定IP访问Doris的CREATE USER支持host限制二是看FE日志里真正的报错信息。我在测试环境遇到过一种情况Doris开了LDAP认证SQL客户端能登录但Stream Load接口请求走的是HTTP Basic AuthLDAP用户没有在Doris内建实体导致报错。这种场景要么在Doris里建对应账号要么走代理认证。6.3 数据重复导入问题Streamloader默认用的是自动生成的label具备幂等性。也就是说同一个分片如果导入失败重试不会重复写入数据。但如果你手动指定了label就得小心同一个label在一个数据库里是唯一的重复使用会报冲突。另外要注意如果你导CSV时没有指定主键而Doris表是UNIQUE KEY模型重复导入相同数据会产生重复记录。Streamloader本身不做去重它只负责搬运。对于UNIQUE KEY模型的表大批量导入后建议手动触发一次compaction观察一下SHOW TABLET里的version数如果版本积压太多查询性能会明显下降。6.4 导入性能上不去怎么调如果你觉得导入速度没达到预期可以从以下几个方面排查先看Doris BE的CPU和内存。如果BE的CPU已经接近100%加客户端并发也没用。看FE的stream_load相关线程数。Doris默认的Stream Load并发上限是max_running_txn_num_per_db如果数据库里的导入事务太多后面的请求会排队。可以通过ADMIN SHOW FRONTEND CONFIG查看。文件太大且没有合理拆分可能导致单批次导入时间过长。把max.file.size调小一点比如500MB可以让任务更均匀地分发到不同BE。检查网络。Streamloader到FE、FE到BE的带宽都是瓶颈如果跨机房导入延迟会明显影响吞吐。我这边一个经验值在万兆内网、BE节点8台每台16核64GB的集群上Streamloader导入CSV数据能达到80MB/s左右的吞吐。如果远低于这个数值先怀疑网络或磁盘别急着加并发。6.5 日志怎么看现场排查三板斧Streamloader在终端上打印的是摘要日志详细日志会写到logs/目录下。你可以在执行时加--log_level debug查看更详细的调试信息./bin/streamloader --config import.yaml --log_level debug排查问题时我一般按这三步来先看终端有没有输出类的汇总信息比如文件数量、总大小、成功/过滤行数。再看logs/streamloader.log里的WARN和ERROR级别日志重点搜label、status、error关键词。最后去Doris FE的日志fe.log里搜对应的导入label看FE侧对这个导入任务的详细判定。如果怀疑是Doris端的问题多在FE日志里搜一下能得到比客户端更准确的报错。7. 从安装到生产落地我的几点体会最后说点实在的。我第一次接触Streamloader是在一个数据迁移项目上当时要往Doris里灌5年的历史订单数据总量大概1.2TB分布在几千个CSV文件里。最开始我写了一个Python脚本用并发调Stream Load接口跑了两天半中途各种报错、重试、数据校验头大得很。后来换成Streamloader配置好fe.nodes、source.file.path和concurrency一条命令跑完全量导入大概用了7个多小时而且中途BE做了次滚动重启它自己断点续传接着跑没丢一条数据。从那之后我所有往Doris导数据的活都优先考虑Streamloader。它不是万能的比如比较复杂的JSON嵌套转换、跨表JOIN后的导入它还是不如写代码灵活但90%的“把一堆文件倒进Doris”的场景它都是最优解。如果你打算在生产环境长期使用建议把Streamloader的二进制文件和配置文件纳入你的部署体系管理。它不依赖外部服务不需要常驻进程用完即走非常适合放到定时任务里。配合Doris的ALTER TABLE手动合并分区、慢查询优化这些日常运维手段整个数据导入链路会非常省心。下载和安装只是第一步真正值钱的是把配置调对、把数据导稳。上面这些参数和坑都是我实打实跑出来的你照着配置至少能少走一半弯路。
返回列表