StarRocks 从 HDFS 加载数据全指南:INSERT+FILES()、Broker Load 与 Pipe 三种方案对比与实战
【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks
本指南以 docs/en/loading/hdfs_load.md 为核心脉络,系统讲解 StarRocks 从 HDFS 加载数据的三种主流方式:同步加载INSERT+FILES()、异步加载 Broker Load、持续异步加载 Pipe。读完本文,你将能够根据数据规模、文件格式与实时性要求选择最合适的加载方案,并掌握从建库建表、发起加载、查看进度到排查失败(如listPath failed)的完整实战流程,同时了解这些能力背后的源码实现(FE/BE 中的加载任务调度、文件唯一性校验等)。
一、三种加载方式总览与选型
StarRocks 为从 HDFS 加载数据提供了三种方式,它们的核心差异在于同步/异步与持续加载能力:
| 加载方式 | 同步/异步 | 支持格式 | 适用场景 |
|---|---|---|---|
INSERT+FILES() | 同步 | Parquet、ORC、CSV(CSV 从 v3.3.0 起) | 小批量数据、交互式加载,使用最简单 |
| Broker Load | 异步 | Parquet、ORC、CSV、JSON(JSON 从 v3.2.3 起) | 长时间运行的大任务,支持加载过程中执行数据变更(如 DELETE) |
| Pipe | 持续异步 | Parquet、ORC(从 v3.2 起) | 大规模批量加载、持续增量加载 |
选型建议如下:
- 大多数场景推荐
INSERT+FILES(),它无需部署额外组件,使用门槛最低。 - 如果需要加载 JSON 等
FILES()不支持的格式,或在加载过程中执行 DELETE 等数据变更,请改用Broker Load。 - 如果需要加载大量数据文件且总数据量很大(例如超过 100 GB 甚至 1 TB),推荐使用Pipe:Pipe 会按文件数量或大小自动拆分任务,将大任务分解为小批次顺序执行,从而保证单个文件的错误不会拖垮整个任务,并最大限度减少因数据错误导致的重复加载成本。
二、开始之前:三个前置条件
1. 准备好源数据
确保待加载的源数据已经正确存放在 HDFS 集群中。本文示例统一使用 HDFS 上的数据文件:
/user/amber/user_behavior_ten_million_rows.parquet2. 检查权限
只有对目标 StarRocks 表拥有INSERT 权限的用户才能加载数据。若当前用户没有该权限,可参照 GRANT 文档执行授权,语法如下:
GRANT INSERT ON TABLE <table_name> IN DATABASE <database_name> TO { ROLE <role_name> | USER <user_identity>}3. 收集认证信息
连接 HDFS 集群可使用simple 认证(简单认证)方式,此时需要准备用于访问 HDFS NameNode 的账号和密码。对应的三个连接参数为:
| Key | 必填 | 说明 |
|---|---|---|
hadoop.security.authentication | 否 | 认证方式,取值为simple(默认值),表示简单认证 |
username | 是 | 访问 HDFS NameNode 的账号用户名 |
password | 是 | 访问 HDFS NameNode 的账号密码 |
三、方式一:使用 INSERT+FILES() 同步加载
该方式自v3.1起可用,当前仅支持Parquet、ORC 和 CSV(自 v3.3.0 起)三种文件格式。
3.1 FILES() 表函数的能力
FILES()表函数可以根据你指定的路径相关属性读取云存储/分布式文件系统中的文件,自动推断文件中的数据表结构,并将文件数据以数据行的形式返回。基于FILES()你可以:
- 使用 SELECT 直接查询 HDFS 中的数据;
- 使用 CREATE TABLE AS SELECT(CTAS)建表并加载数据;
- 使用 INSERT 将数据加载进已有表。
3.2 用 SELECT+FILES() 直接预览 HDFS 数据
在正式建表之前,直接查询 HDFS 中的文件可以很好地预览数据集内容,例如:在不存储数据的前提下预览数据集、查询某列的 min/max 值以决定数据类型、检查是否存在NULL值。
SELECT * FROM FILES ( "path" = "hdfs://<hdfs_ip>:<hdfs_port>/user/amber/user_behavior_ten_million_rows.parquet", "format" = "parquet", "hadoop.security.authentication" = "simple", "username" = "<hdfs_username>", "password" = "<hdfs_password>" ) LIMIT 3;返回结果如下(注意:返回的列名由 Parquet 文件本身提供):
+--------+---------+------------+--------------+---------------------+ | UserID | ItemID | CategoryID | BehaviorType | Timestamp | +--------+---------+------------+--------------+---------------------+ | 543711 | 829192 | 2355072 | pv | 2017-11-27 08:22:37 | | 543711 | 2056618 | 3645362 | pv | 2017-11-27 10:16:46 | | 543711 | 1165492 | 3645362 | pv | 2017-11-27 10:17:00 | +--------+---------+------------+--------------+---------------------+3.3 用 CTAS 自动建表并加载
将上面的查询包装进 CREATE TABLE AS SELECT(CTAS),即可利用 schema 推断自动完成建表和加载:StarRocks 会推断表结构、创建目标表并写入数据。使用 Parquet 文件时无需显式声明列名和类型,因为 Parquet 格式本身就包含列名。
注意:使用 schema 推断的 CREATE TABLE 语法不允许设置副本数,需在建表前通过 FE 配置设定。以下示例适用于三副本系统:
ADMIN SET FRONTEND CONFIG ('default_replication_num' = "3");
创建数据库并切换:
CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase;使用 CTAS 建表并加载/user/amber/user_behavior_ten_million_rows.parquet:
CREATE TABLE user_behavior_inferred AS SELECT * FROM FILES ( "path" = "hdfs://<hdfs_ip>:<hdfs_port>/user/amber/user_behavior_ten_million_rows.parquet", "format" = "parquet", "hadoop.security.authentication" = "simple", "username" = "<hdfs_username>", "password" = "<hdfs_password>" );建表后使用 DESCRIBE 查看自动推断出的表结构:
DESCRIBE user_behavior_inferred;+--------------+-----------+------+-------+---------+-------+ | Field | Type | Null | Key | Default | Extra | +--------------+-----------+------+-------+---------+-------+ | UserID | bigint | YES | true | NULL | | | ItemID | bigint | YES | true | NULL | | | CategoryID | bigint | YES | true | NULL | | | BehaviorType | varbinary | YES | false | NULL | | | Timestamp | varbinary | YES | false | NULL | | +--------------+-----------+------+-------+---------+-------+查询表验证加载结果:
SELECT * from user_behavior_inferred LIMIT 3;+--------+--------+------------+--------------+---------------------+ | UserID | ItemID | CategoryID | BehaviorType | Timestamp | +--------+--------+------------+--------------+---------------------+ | 84 | 56257 | 1879194 | pv | 2017-11-26 05:56:23 | | 84 | 108021 | 2982027 | pv | 2017-12-02 05:43:00 | | 84 | 390657 | 1879194 | pv | 2017-11-28 11:20:30 | +--------+--------+------------+--------------+---------------------+3.4 用 INSERT 加载进手动建的表
CTAS 的表结构由文件推断,而你可能希望自定义目标表,例如指定:列的数据类型、是否允许 NULL、默认值;键类型与键列;数据分区与分桶方式。
构建最高效的表结构需要了解数据的使用方式与列内容,本文不展开表设计,请参考 Table types。
基于对 HDFS 数据的预查询(见 3.2 节),可以做出如下建表决策:
Timestamp列数据与 VARBINARY 类型匹配,DDL 中指定为varbinary;- 数据集中不存在
NULL值,因此 DDL 不设置任何列为可空; - 根据预期的查询类型,排序键与分桶列设置为
UserID(你的场景可能不同,也可以使用ItemID作为排序键)。
CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase; CREATE TABLE user_behavior_declared ( UserID int(11), ItemID int(11), CategoryID int(11), BehaviorType varchar(65533), Timestamp varbinary ) ENGINE = OLAP DUPLICATE KEY(UserID) DISTRIBUTED BY HASH(UserID);查看表结构并与FILES()推断出的结构对比:
DESCRIBE user_behavior_declared;+--------------+----------------+------+-------+---------+-------+ | Field | Type | Null | Key | Default | Extra | +--------------+----------------+------+-------+---------+-------+ | UserID | int | NO | true | NULL | | | ItemID | int | NO | false | NULL | | | CategoryID | int | NO | false | NULL | | | BehaviorType | varchar(65533) | NO | false | NULL | | | Timestamp | varbinary | NO | false | NULL | | +--------------+----------------+------+-------+---------+-------+ 5 rows in set (0.00 sec)对比两张表的差异可以重点看三个维度:数据类型、是否可空、键字段。在生产环境中,为了更精确地控制目标表结构并获得更好的查询性能,推荐手动指定表结构。
使用INSERT INTO SELECT FROM FILES()加载数据:
INSERT INTO user_behavior_declared SELECT * FROM FILES ( "path" = "hdfs://<hdfs_ip>:<hdfs_port>/user/amber/user_behavior_ten_million_rows.parquet", "format" = "parquet", "hadoop.security.authentication" = "simple", "username" = "<hdfs_username>", "password" = "<hdfs_password>" );加载完成后验证数据:
SELECT * from user_behavior_declared LIMIT 3;+--------+---------+------------+--------------+---------------------+ | UserID | ItemID | CategoryID | BehaviorType | Timestamp | +--------+---------+------------+--------------+---------------------+ | 107 | 1568743 | 4476428 | pv | 2017-11-25 14:29:53 | | 107 | 470767 | 1020087 | pv | 2017-11-25 14:32:31 | | 107 | 358238 | 1817004 | pv | 2017-11-25 14:43:23 | +--------+---------+------------+--------------+---------------------+3.5 查看 INSERT 加载进度
从v3.1起,可以从 StarRocks Information Schema 的loads视图中查询 INSERT 任务的进度:
SELECT * FROM information_schema.loads ORDER BY JOB_ID DESC;如果提交了多个加载任务,可以按任务关联的LABEL过滤。例如:
SELECT * FROM information_schema.loads WHERE LABEL = 'insert_0d86c3f9-851f-11ee-9c3e-00163e044958' \G *************************** 1. row *************************** JOB_ID: 10214 LABEL: insert_0d86c3f9-851f-11ee-9c3e-00163e044958 DATABASE_NAME: mydatabase STATE: FINISHED PROGRESS: ETL:100%; LOAD:100% TYPE: INSERT PRIORITY: NORMAL SCAN_ROWS: 10000000 FILTERED_ROWS: 0 UNSELECTED_ROWS: 0 SINK_ROWS: 10000000 ETL_INFO: TASK_INFO: resource:N/A; timeout(s):300; max_filter_ratio:0.0 CREATE_TIME: 2023-11-17 15:58:14 ETL_START_TIME: 2023-11-17 15:58:14 ETL_FINISH_TIME: 2023-11-17 15:58:14 LOAD_START_TIME: 2023-11-17 15:58:14 LOAD_FINISH_TIME: 2023-11-17 15:58:18 JOB_DETAILS: {"All backends":{"0d86c3f9-851f-11ee-9c3e-00163e044958":[10120]},"FileNumber":0,"FileSize":0,"InternalTableLoadBytes":311710786,"InternalTableLoadRows":10000000,"ScanBytes":581574034,"ScanRows":10000000,"TaskNumber":1,"Unfinished backends":{"0d86c3f9-851f-11ee-9c3e-00163e044958":[]}} ERROR_MSG: NULL TRACKING_URL: NULL TRACKING_SQL: NULL REJECTED_RECORD_PATH: NULL注意:INSERT 是同步命令。如果 INSERT 任务仍在运行,需要另开一个会话查询其执行状态。
3.6 深入:FILES() 的语法、路径与认证参数
FILES()表函数的完整语法为:
FILES( data_location , [data_format] [, schema_detect ] [, StorageCredentialParams ] [, columns_from_path ] [, list_files_only ] [, list_recursively])关键参数说明:
data_location(即path):访问文件的 URI,可以指向单个文件,也可以使用通配符?、*、[]、^指向多个文件。访问 HDFS 时格式为"hdfs://<hdfs_host>:<hdfs_port>/<hdfs_path>"(例如"hdfs://127.0.0.1:9000/path/file.parquet")。通配符同样可用于中间路径,例如"hdfs://<hdfs_host>:<hdfs_port>/user/data/tablename/dt=202104*/*"可匹配所有202104分区下的文件。data_format:文件格式,取值parquet、orc(v3.3 起)、csv(v3.3 起)、avro(v3.4.4 起,仅用于加载)。CSV 格式还支持csv.column_separator(默认\t,Hive 文件需写\\x01)、csv.enclose(默认NONE)、csv.skip_header(默认0)、csv.escape(默认NONE)等细分参数;CSV 中空值用\N表示。StorageCredentialParams:访问存储系统的认证信息。HDFS 支持 simple 认证(参数见本文第二节表格);也支持通过放置于fe/conf、be/conf、cn/conf目录下的hdfs-site.xml配置 Kerberos 认证与 HA 模式(Kerberos 还需在fe.conf/be.conf/cn.conf的JAVA_OPTS中追加-Djava.security.krb5.conf=<path_to_kerberos_conf_file>,并定期执行kinit -kt <keytab路径> <principal>刷新票据)。- schema 检测相关参数(v3.2 起):
auto_detect_sample_files(每批随机采样文件数,默认2)、auto_detect_sample_rows(每个采样文件扫描行数,默认500)、auto_detect_types(CSV 类型推断开关,默认true)。此外 v3.4.0 起还支持fill_mismatch_column_with(值为none/null)以应对不同分区 schema 不一致的情况。 columns_from_path(v3.2 起):从文件路径中的 key/value 对提取列值,例如"columns_from_path" = "country, city"。list_files_only/list_recursively(v3.4.0 起):仅列出文件元信息(返回PATH、SIZE、IS_DIR、MODIFICATION_TIME),list_recursively仅在list_files_only=true时生效。
四、方式二:使用 Broker Load 异步加载
Broker Load 是一种异步加载方式:任务提交后由后台进程完成 HDFS 连接、数据拉取与入库,客户端无需保持连接。
支持的格式:Parquet、ORC、CSV、JSON(JSON 从 v3.2.3 起)。
4.1 优点
- Broker Load 在后台运行,客户端无需一直保持连接,任务也能继续执行;
- 适合长时间运行的任务,默认超时长达4 小时;
- 除 Parquet、ORC 外,还支持 CSV 与 JSON(JSON 自 v3.2.3 起),覆盖更多文件格式。
4.2 数据流
- 用户创建加载任务;
- 前端节点(FE)生成查询计划并分发给后端节点(BE)或计算节点(CN);
- BE/CN 从数据源拉取数据并加载进 StarRocks。
4.3 典型示例:建库建表并发起加载
创建数据库并切换到目标库,然后手动建表(推荐表结构与待加载的 Parquet 文件保持一致):
CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase; CREATE TABLE user_behavior ( UserID int(11), ItemID int(11), CategoryID int(11), BehaviorType varchar(65533), Timestamp varbinary ) ENGINE = OLAP DUPLICATE KEY(UserID) DISTRIBUTED BY HASH(UserID);启动 Broker Load 任务,将 HDFS 上的/user/amber/user_behavior_ten_million_rows.parquet加载到user_behavior表:
LOAD LABEL user_behavior ( DATA INFILE("hdfs://<hdfs_ip>:<hdfs_port>/user/amber/user_behavior_ten_million_rows.parquet") INTO TABLE user_behavior FORMAT AS "parquet" ) WITH BROKER ( "hadoop.security.authentication" = "simple", "username" = "<hdfs_username>", "password" = "<hdfs_password>" ) PROPERTIES ( "timeout" = "72000" );该任务由四个主要部分组成:
LABEL:用于查询加载任务状态的字符串(label 在同一数据库内唯一);LOAD声明:源数据 URI、源数据格式与目标表名;BROKER:数据源的连接信息;PROPERTIES:超时时间等应用于加载任务的其他属性。
更详细的语法与参数说明见 BROKER LOAD。
4.4 查看加载进度
从v3.1起,可以从loads视图查询 Broker Load 任务进度:
SELECT * FROM information_schema.loads;多个任务时可按LABEL过滤:
SELECT * FROM information_schema.loads WHERE LABEL = 'user_behavior';下面的输出中,任务user_behavior有两条记录:第一条STATE为CANCELLED,查看ERROR_MSG可知任务因listPath failed失败;第二条STATE为FINISHED,表示任务成功:
JOB_ID|LABEL |DATABASE_NAME|STATE |PROGRESS |TYPE |PRIORITY|SCAN_ROWS|FILTERED_ROWS|UNSELECTED_ROWS|SINK_ROWS|ETL_INFO|TASK_INFO |CREATE_TIME |ETL_START_TIME |ETL_FINISH_TIME |LOAD_START_TIME |LOAD_FINISH_TIME |JOB_DETAILS |ERROR_MSG |TRACKING_URL|TRACKING_SQL|REJECTED_RECORD_PATH| ------+-------------------------------------------+-------------+---------+-------------------+------+--------+---------+-------------+---------------+---------+--------+----------------------------------------------------+-------------------+-------------------+-------------------+-------------------+-------------------+---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------+--------------------------------------+------------+------------+--------------------+ 10121|user_behavior |mydatabase |CANCELLED|ETL:N/A; LOAD:N/A |BROKER|NORMAL | 0| 0| 0| 0| |resource:N/A; timeout(s):72000; max_filter_ratio:0.0|2023-08-10 14:59:30| | | |2023-08-10 14:59:34|{"All backends":{},"FileNumber":0,"FileSize":0,"InternalTableLoadBytes":0,"InternalTableLoadRows":0,"ScanBytes":0,"ScanRows":0,"TaskNumber":0,"Unfinished backends":{}} |type:ETL_RUN_FAIL; msg:listPath failed| | | | 10106|user_behavior |mydatabase |FINISHED |ETL:100%; LOAD:100%|BROKER|NORMAL | 86953525| 0| 0| 86953525| |resource:N/A; timeout(s):72000; max_filter_ratio:0.0|2023-08-10 14:50:15|2023-08-10 14:50:19|2023-08-10 14:50:19|2023-08-10 14:50:19|2023-08-10 14:55:10|{"All backends":{"a5fe5e1d-d7d0-4826-ba99-c7348f9a5f2f":[10004]},"FileNumber":1,"FileSize":1225637388,"InternalTableLoadBytes":2710603082,"InternalTableLoadRows":86953525,"ScanBytes":1225637388,"ScanRows":86953525,"TaskNumber":1,"Unfinished backends":{"a5| | | | |listPath failed通常意味着 BE/CN 无法列出指定 HDFS 路径下的文件,常见原因包括路径不存在、认证信息错误或网络不通,需要检查DATA INFILE路径与 BROKER 认证参数。
确认加载完成后,查询目标表子集验证数据:
SELECT * from user_behavior LIMIT 3;+--------+---------+------------+--------------+---------------------+ | UserID | ItemID | CategoryID | BehaviorType | Timestamp | +--------+---------+------------+--------------+---------------------+ | 142 | 2869980 | 2939262 | pv | 2017-11-25 03:43:22 | | 142 | 2522236 | 1669167 | pv | 2017-11-25 15:14:12 | | 142 | 3031639 | 3607361 | pv | 2017-11-25 15:19:25 | +--------+---------+------------+--------------+---------------------+4.5 深入:Broker Load 的底层执行与标签语义
从 FE 源码结构看,Broker Load 由 fe-core/src/main/java/com/starrocks/load/loadv2/BrokerLoadJob.java 与 BrokerLoadPendingTask.java 等类驱动:FE 侧创建任务、生成执行计划,BE/CN 侧执行拉取与写入。
值得关注的语义特性:
- 一个任务可加载多个数据文件:一个
LOAD语句中可用多个data_desc声明多个文件,也可用通配符?、*、[]、{}、^声明一个路径下的所有文件;一个任务内多个文件的加载具备事务原子性——要么全部成功、要么全部失败,不会出现部分成功。 - Label 与 Exactly-Once:每个加载任务有一个在整个数据库内唯一的 label。任务进入
FINISHED状态后 label 不可复用;只有CANCELLED状态的 label 可复用。通常复用 label 重试同一批数据即可实现 Exactly-Once 语义,防止重复加载。 loads视图字段:包含LABEL、DB_NAME、TABLE_NAME、STATE(PENDING/BEGIN、QUEUEING/BEFORE_LOAD、LOADING、PREPARING、PREPARED、COMMITED、FINISHED、CANCELLED)、PROGRESS、TYPE(Broker Load 为BROKER)、PRIORITY、SCAN_ROWS、FILTERED_ROWS、UNSELECTED_ROWS、SINK_ROWS、ERROR_MSG、TRACKING_SQL、REJECTED_RECORD_PATH等;REJECTED_RECORD_PATH可用来获取被过滤的不合格数据行。
五、方式三:使用 Pipe 持续加载
自v3.2起,StarRocks 提供 Pipe 加载方式,目前仅支持Parquet 和 ORC文件格式。
5.1 优点
- 微批次大规模加载,降低错误重试成本:Pipe 按文件数量或大小自动将大任务拆分为多个小的顺序子任务,单个文件的错误不会影响整个加载任务;Pipe 会记录每个文件的加载状态,便于定位和修复出错文件,从而显著降低因数据错误导致的重复加载成本。
- 持续加载,降低人力成本:创建 Pipe 任务时指定
"AUTO_INGEST" = "TRUE",Pipe 会持续监控指定路径下数据文件的变化,自动将新增或更新的文件数据加载进目标表。 - 文件唯一性校验,防止重复加载:加载过程中,Pipe 根据文件名 + digest校验每个文件的唯一性;如果某个文件名与 digest 组合已被处理过,Pipe 会跳过后续所有相同文件名与 digest 的文件。注意:HDFS 以
LastModifiedTime(最后修改时间)作为文件 digest。 - 状态可观测:每个数据文件的加载状态记录并保存在
information_schema.pipe_files视图中;删除 Pipe 任务后,其加载文件记录也会一并删除。
5.2 数据流
数据从 HDFS/S3 等源端经 Pipe 管道,以微批次方式持续流入 StarRocks。
5.3 Pipe 与 INSERT+FILES() 的差异
Pipe 任务会根据每个数据文件的大小与行数拆分为一个或多个事务,加载过程中用户可以查询中间结果;而 INSERT+FILES()任务作为单个事务整体执行,加载过程中用户无法看到中间数据。
5.4 文件加载顺序
对每个 Pipe 任务,StarRocks 维护一个文件队列,按微批次取文件加载。Pipe不保证文件按上传顺序加载,因此较新的数据可能先于较旧的数据被加载。
5.5 典型示例
建库、建表(推荐表结构与 Parquet 文件保持一致):
CREATE DATABASE IF NOT EXISTS mydatabase; USE mydatabase; CREATE TABLE user_behavior_replica ( UserID int(11), ItemID int(11), CategoryID int(11), BehaviorType varchar(65533), Timestamp varbinary ) ENGINE = OLAP DUPLICATE KEY(UserID) DISTRIBUTED BY HASH(UserID);创建 Pipe 任务user_behavior_replica,将 HDFS 数据文件加载到user_behavior_replica表:
CREATE PIPE user_behavior_replica PROPERTIES ( "AUTO_INGEST" = "TRUE" ) AS INSERT INTO user_behavior_replica SELECT * FROM FILES ( "path" = "hdfs://<hdfs_ip>:<hdfs_port>/user/amber/user_behavior_ten_million_rows.parquet", "format" = "parquet", "hadoop.security.authentication" = "simple", "username" = "<hdfs_username>", "password" = "<hdfs_password>" );任务包含四个主要部分:
pipe_name:Pipe 名称,在所属数据库内必须唯一;INSERT_SQL:用于从指定源数据文件加载数据到目标表的INSERT INTO SELECT FROM FILES语句;PROPERTIES:一组可选参数,控制 Pipe 的执行方式,包括AUTO_INGEST、POLL_INTERVAL、BATCH_SIZE、BATCH_FILES,均以"key" = "value"格式指定。
CREATE PIPE 的PROPERTIES参数表如下:
| 属性 | 默认值 | 说明 |
|---|---|---|
AUTO_INGEST | TRUE | 是否开启自动增量加载。TRUE开启自动增量;FALSE只加载任务创建时指定的源数据文件内容,后续新增或更新的文件不会被加载(批量一次性加载可设为FALSE) |
POLL_INTERVAL | 300秒 | 自动增量加载的轮询间隔 |
BATCH_SIZE | 1GB | 每个批次加载的数据量(不带单位时默认按字节计) |
BATCH_FILES | 256 | 每个批次加载的源数据文件数 |
5.6 查看加载进度
方式一:使用 SHOW PIPES:
SHOW PIPES;多个任务时可按NAME过滤:
SHOW PIPES WHERE NAME = 'user_behavior_replica' \G *************************** 1. row *************************** DATABASE_NAME: mydatabase PIPE_ID: 10252 PIPE_NAME: user_behavior_replica STATE: RUNNING TABLE_NAME: mydatabase.user_behavior_replica LOAD_STATUS: {"loadedFiles":1,"loadedBytes":132251298,"loadingFiles":0,"lastLoadedTime":"2023-11-17 16:13:22"} LAST_ERROR: NULL CREATED_TIME: 2023-11-17 16:13:15 1 row in set (0.00 sec)STATE可取RUNNING、FINISHED、SUSPENDED、ERROR;LOAD_STATUS中的loadedFiles、loadedBytes反映整体加载进度。
方式二:查询 Information Schema 的pipes视图:
SELECT * FROM information_schema.pipes;SELECT * FROM information_schema.pipes WHERE pipe_name = 'user_behavior_replica' \G5.7 查看单个文件的加载状态
从pipe_files视图查询经 Pipe 加载的文件的详细状态:
SELECT * FROM information_schema.pipe_files;按PIPE_NAME过滤:
SELECT * FROM information_schema.pipe_files WHERE pipe_name = 'user_behavior_replica' \G *************************** 1. row *************************** DATABASE_NAME: mydatabase PIPE_ID: 10252 PIPE_NAME: user_behavior_replica FILE_NAME: hdfs://172.26.195.67:9000/user/amber/user_behavior_ten_million_rows.parquet FILE_VERSION: 1700035418838 FILE_SIZE: 132251298 LAST_MODIFIED: 2023-11-15 08:03:38 LOAD_STATE: FINISHED STAGED_TIME: 2023-11-17 16:13:16 START_LOAD_TIME: 2023-11-17 16:13:17 FINISH_LOAD_TIME: 2023-11-17 16:13:22 ERROR_MSG: 1 row in set (0.02 sec)pipe_files视图的关键字段:FILE_NAME(数据文件名)、FILE_VERSION(文件 digest,HDFS 上即LastModifiedTime)、FILE_SIZE(字节)、LOAD_STATE(取值UNLOADED、LOADING、FINISHED、ERROR)、STAGED_TIME(首次被 Pipe 记录的时间)、START_LOAD_TIME/FINISH_LOAD_TIME(加载起止时间)、ERROR_MSG(错误详情)。注意,HDFS 场景下FILE_VERSION为LastModifiedTime(毫秒时间戳),这正是 Pipe 做文件唯一性校验的 digest 依据。
5.8 管理 Pipe
Pipe 支持修改、暂停/恢复、删除、查询,以及对指定数据文件重试加载:
- ALTER PIPE:修改 Pipe 属性;
- SUSPEND or RESUME PIPE:暂停或恢复 Pipe;
- DROP PIPE:删除 Pipe(其
pipe_files加载记录一并删除); - SHOW PIPES:查询 Pipe 状态;
- RETRY FILE:对失败的数据文件重试加载。
5.9 深入:Pipe 的源码实现
从 FE 源码可以印证上述行为:
- 轮询与微批次:在 fe-core/src/main/java/com/starrocks/load/pipe/FilePipeSource.java 的
poll()方法中,Pipe 通过HdfsUtil.listFileMeta()拉取路径下的文件元数据,转换为PipeFileRecord记录后交给文件仓库(FileListRepo)暂存;当autoIngest为false时,一次性 Pipe 会在加载完首轮文件后进入eos(end-of-source)状态,不再拉取新文件——这与文档中AUTO_INGEST=FALSE的语义一致。该类的batchSize、batchFiles字段则对应BATCH_SIZE、BATCH_FILES属性。 - 文件唯一性校验:在 fe-core/src/main/java/com/starrocks/load/pipe/PipeFileRecord.java 中,
equals()/hashCode()以pipeId + fileName + fileVersion三元组作为记录相等性依据,即"文件名 + digest"唯一性校验的代码级实现(HDFS 的 digest 为LastModifiedTime)。
六、总结:如何选择加载方案
| 需求场景 | 推荐方案 | 理由 |
|---|---|---|
| 小批量、交互式加载,追求简单 | INSERT+FILES() | 一条 SQL 即可完成查询/建表/加载,支持 schema 自动推断 |
| 加载 JSON 格式,或加载中需要 DELETE 等数据变更 | Broker Load | 支持 Parquet/ORC/CSV/JSON,异步运行、默认超时 4 小时 |
| 总数据量超百 GB/TB 级的批量导入 | Pipe | 微批次拆分、文件级错误隔离、支持断点重试 |
| 新文件持续产生,需要自动增量入库 | Pipe(AUTO_INGEST=TRUE) | 持续轮询目录,按"文件名+digest"去重,避免重复加载 |
无论选择哪种方式,加载后都可以通过information_schema.loads(INSERT 与 Broker Load)、SHOW PIPES/information_schema.pipes/information_schema.pipe_files(Pipe)等视图与命令追踪任务状态、定位失败原因,再结合本指南中的参数说明与源码实现理解其背后的执行机制。
【免费下载链接】starrocksThe world's fastest open query engine for sub-second analytics both on and off the data lakehouse. With the flexibility to support nearly any scenario, StarRocks provides best-in-class performance for multi-dimensional analytics, real-time analytics, and ad-hoc queries. A Linux Foundation project.项目地址: https://gitcode.com/GitHub_Trending/st/starrocks
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考