StarRocks 从 HDFS 加载数据全指南:INSERT+FILES()、Broker Load 与 Pipe 三种方案对比与实战

发布时间:2026/9/17 0:06:04
StarRocks 从 HDFS 加载数据全指南:INSERT+FILES()、Broker Load 与 Pipe 三种方案对比与实战 StarRocks 从 HDFS 加载数据全指南INSERTFILES()、Broker Load 与 Pipe 三种方案对比与实战【免费下载链接】starrocksThe worlds 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 加载数据的三种主流方式同步加载INSERTFILES()、异步加载 Broker Load、持续异步加载 Pipe。读完本文你将能够根据数据规模、文件格式与实时性要求选择最合适的加载方案并掌握从建库建表、发起加载、查看进度到排查失败如listPath failed的完整实战流程同时了解这些能力背后的源码实现FE/BE 中的加载任务调度、文件唯一性校验等。一、三种加载方式总览与选型StarRocks 为从 HDFS 加载数据提供了三种方式它们的核心差异在于同步/异步与持续加载能力加载方式同步/异步支持格式适用场景INSERTFILES()同步Parquet、ORC、CSVCSV 从 v3.3.0 起小批量数据、交互式加载使用最简单Broker Load异步Parquet、ORC、CSV、JSONJSON 从 v3.2.3 起长时间运行的大任务支持加载过程中执行数据变更如 DELETEPipe持续异步Parquet、ORC从 v3.2 起大规模批量加载、持续增量加载选型建议如下大多数场景推荐INSERTFILES()它无需部署额外组件使用门槛最低。如果需要加载 JSON 等FILES()不支持的格式或在加载过程中执行 DELETE 等数据变更请改用Broker Load。如果需要加载大量数据文件且总数据量很大例如超过 100 GB 甚至 1 TB推荐使用PipePipe 会按文件数量或大小自动拆分任务将大任务分解为小批次顺序执行从而保证单个文件的错误不会拖垮整个任务并最大限度减少因数据错误导致的重复加载成本。二、开始之前三个前置条件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 的账号密码三、方式一使用 INSERTFILES() 同步加载该方式自v3.1起可用当前仅支持Parquet、ORC 和 CSV自 v3.3.0 起三种文件格式。3.1 FILES() 表函数的能力FILES()表函数可以根据你指定的路径相关属性读取云存储/分布式文件系统中的文件自动推断文件中的数据表结构并将文件数据以数据行的形式返回。基于FILES()你可以使用 SELECT 直接查询 HDFS 中的数据使用 CREATE TABLE AS SELECTCTAS建表并加载数据使用 INSERT 将数据加载进已有表。3.2 用 SELECTFILES() 直接预览 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 SELECTCTAS即可利用 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.parquetCREATE 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/dt202104*/*可匹配所有202104分区下的文件。data_format文件格式取值parquet、orcv3.3 起、csvv3.3 起、avrov3.4.4 起仅用于加载。CSV 格式还支持csv.column_separator默认\tHive 文件需写\\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.confpath_to_kerberos_conf_file并定期执行kinit -kt keytab路径 principal刷新票据。schema 检测相关参数v3.2 起auto_detect_sample_files每批随机采样文件数默认2、auto_detect_sample_rows每个采样文件扫描行数默认500、auto_detect_typesCSV 类型推断开关默认true。此外 v3.4.0 起还支持fill_mismatch_column_with值为none/null以应对不同分区 schema 不一致的情况。columns_from_pathv3.2 起从文件路径中的 key/value 对提取列值例如columns_from_path country, city。list_files_only/list_recursivelyv3.4.0 起仅列出文件元信息返回PATH、SIZE、IS_DIR、MODIFICATION_TIMElist_recursively仅在list_files_onlytrue时生效。四、方式二使用 Broker Load 异步加载Broker Load 是一种异步加载方式任务提交后由后台进程完成 HDFS 连接、数据拉取与入库客户端无需保持连接。支持的格式Parquet、ORC、CSV、JSONJSON 从 v3.2.3 起。4.1 优点Broker Load 在后台运行客户端无需一直保持连接任务也能继续执行适合长时间运行的任务默认超时长达4 小时除 Parquet、ORC 外还支持 CSV 与 JSONJSON 自 v3.2.3 起覆盖更多文件格式。4.2 数据流用户创建加载任务前端节点FE生成查询计划并分发给后端节点BE或计算节点CNBE/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、STATEPENDING/BEGIN、QUEUEING/BEFORE_LOAD、LOADING、PREPARING、PREPARED、COMMITED、FINISHED、CANCELLED、PROGRESS、TYPEBroker 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 TRUEPipe 会持续监控指定路径下数据文件的变化自动将新增或更新的文件数据加载进目标表。文件唯一性校验防止重复加载加载过程中Pipe 根据文件名 digest校验每个文件的唯一性如果某个文件名与 digest 组合已被处理过Pipe 会跳过后续所有相同文件名与 digest 的文件。注意HDFS 以LastModifiedTime最后修改时间作为文件 digest。状态可观测每个数据文件的加载状态记录并保存在information_schema.pipe_files视图中删除 Pipe 任务后其加载文件记录也会一并删除。5.2 数据流数据从 HDFS/S3 等源端经 Pipe 管道以微批次方式持续流入 StarRocks。5.3 Pipe 与 INSERTFILES() 的差异Pipe 任务会根据每个数据文件的大小与行数拆分为一个或多个事务加载过程中用户可以查询中间结果而 INSERTFILES()任务作为单个事务整体执行加载过程中用户无法看到中间数据。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_namePipe 名称在所属数据库内必须唯一INSERT_SQL用于从指定源数据文件加载数据到目标表的INSERT INTO SELECT FROM FILES语句PROPERTIES一组可选参数控制 Pipe 的执行方式包括AUTO_INGEST、POLL_INTERVAL、BATCH_SIZE、BATCH_FILES均以key value格式指定。CREATE PIPE 的PROPERTIES参数表如下属性默认值说明AUTO_INGESTTRUE是否开启自动增量加载。TRUE开启自动增量FALSE只加载任务创建时指定的源数据文件内容后续新增或更新的文件不会被加载批量一次性加载可设为FALSEPOLL_INTERVAL300秒自动增量加载的轮询间隔BATCH_SIZE1GB每个批次加载的数据量不带单位时默认按字节计BATCH_FILES256每个批次加载的源数据文件数5.6 查看加载进度方式一使用 SHOW PIPESSHOW 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、ERRORLOAD_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文件 digestHDFS 上即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 管理 PipePipe 支持修改、暂停/恢复、删除、查询以及对指定数据文件重试加载ALTER PIPE修改 Pipe 属性SUSPEND or RESUME PIPE暂停或恢复 PipeDROP 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 会在加载完首轮文件后进入eosend-of-source状态不再拉取新文件——这与文档中AUTO_INGESTFALSE的语义一致。该类的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。六、总结如何选择加载方案需求场景推荐方案理由小批量、交互式加载追求简单INSERTFILES()一条 SQL 即可完成查询/建表/加载支持 schema 自动推断加载 JSON 格式或加载中需要 DELETE 等数据变更Broker Load支持 Parquet/ORC/CSV/JSON异步运行、默认超时 4 小时总数据量超百 GB/TB 级的批量导入Pipe微批次拆分、文件级错误隔离、支持断点重试新文件持续产生需要自动增量入库PipeAUTO_INGESTTRUE持续轮询目录按文件名digest去重避免重复加载无论选择哪种方式加载后都可以通过information_schema.loadsINSERT 与 Broker Load、SHOW PIPES/information_schema.pipes/information_schema.pipe_filesPipe等视图与命令追踪任务状态、定位失败原因再结合本指南中的参数说明与源码实现理解其背后的执行机制。【免费下载链接】starrocksThe worlds 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),仅供参考