资讯动态

StarRocks pipe_files 信息模式视图完全指南:通过 Pipe 加载的数据文件状态监控

发布时间:2026/9/17 20:01:21 来源:尧图企业网站定制
StarRocks pipe_files 信息模式视图完全指南通过 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/starrockspipe_files是 StarRocks 提供的一张 Information Schemainformation_schema系统视图用于查看通过指定 Pipe 加载的数据文件的详细状态包括文件的加载进度、时间节点与错误信息。本文以该视图为核心结合 FE 源码中系统表的列定义、PipeFileRecord记录模型与底层pipe_file_list存储表实现讲解每个字段的语义、LOAD_STATE状态机的迁移过程以及日常运维中的查询实践帮助你在使用 Pipe 进行持续数据导入时快速定位加载问题。一、为什么需要 pipe_filesPipe 持续导入的可见性StarRocks 的Pipe是一种面向对象存储如 S3、HDFS 等的持续数据导入机制只需定义一个指向存储路径的 PipeStarRocks 便会自动扫描该路径下新增或修改的数据文件并持续将其加载到目标表中。由于加载过程是异步、长时运行的用户需要一种方式来回答以下问题哪些数据文件已经被 Pipe 扫描到staged每个文件当前处于什么加载状态文件什么时候开始加载、什么时候加载完成如果加载失败失败原因是什么pipe_files视图正是为此而生的观测入口。它记录了 Pipe 所管理数据文件的逐文件状态支持从 StarRocks v3.2 版本开始使用。底层数据由 FE 在加载过程中写入一张内部表并对外暴露为系统视图因此你不需要关心内部存储细节直接通过标准的 SQL 查询即可获得文件级别的加载进度。二、字段详解逐列理解 pipe_files 的语义pipe_files视图提供以下字段完整覆盖了从文件被发现到加载结束的完整生命周期信息FieldDescriptionDATABASE_NAMEPipe 所在数据库的名称。PIPE_IDPipe 的唯一 ID。PIPE_NAMEPipe 的名称。FILE_NAME数据文件的名称包含其在对象存储中的路径信息。FILE_VERSION数据文件的摘要digest用于识别同一文件的版本变化。FILE_SIZE数据文件的大小单位为字节bytes。LAST_MODIFIED数据文件最后一次被修改的时间。格式yyyy-MM-dd HH:mm:ss。例如2023-07-24 14:58:58。LOAD_STATE数据文件的加载状态。有效值UNLOADED、LOADING、FINISHED、ERROR。STAGED_TIME数据文件首次被 Pipe 记录staged的日期与时间。格式yyyy-MM-dd HH:mm:ss。例如2023-07-24 14:58:58。START_LOAD_TIME数据文件开始加载的日期与时间。格式yyyy-MM-dd HH:mm:ss。例如2023-07-24 14:58:58。FINISH_LOAD_TIME数据文件加载完成的日期与时间。格式yyyy-MM-dd HH:mm:ss。例如2023-07-24 14:58:58。ERROR_MSG数据文件加载失败的详细信息。在 FE 源码中这张视图的列定义由 PipeFileSystemTable.java 实现并作为SCHEMA类型的系统表注册到information_schema数据库中见 InfoSchemaDb.java 中的registerTableUnlocked(PipeFileSystemTable.create())。其中各时间字段在源码中为VARCHAR(16)类型ERROR_MSG为VARCHAR(512)FILE_SIZE与PIPE_ID为BIGINT。需要说明的是源码中的PipeFileSystemTable还保留了若干 TODO 注释例如FILE_ROWS文件行数、ERROR_COUNT错误数、ERROR_LINE错误行号等列目前并未对外暴露仅存在于内部记录模型中。因此当你查询pipe_files时实际可见列以本文表格为准。三、LOAD_STATE 状态机从扫描到加载完成的全过程LOAD_STATE是这张视图中最核心的字段它标识每个数据文件当前的加载进度。文档列出的有效值为UNLOADED、LOADING、FINISHED和ERROR。从源码角度看状态定义位于 FileListRepo.java 中的PipeFileState枚举其完整定义还包括一个文档未提及的SKIPPED状态用于手动跳过某个文件public enum PipeFileState { UNLOADED, LOADING, FINISHED, SKIPPED, ERROR }各状态的含义与迁移路径如下UNLOADED待加载文件已被 Pipe 扫描并记录但尚未开始加载。这是文件进入 Pipe 文件清单后的初始状态。在 PipeFileRecord.java 中新记录创建时默认即为UNLOADED同时记录lastModified与stagedTime为当前时间。LOADING加载中Pipe 的加载任务已经取出该文件并开始导入此时会记录START_LOAD_TIME并写入对应的insert_label以便关联本次插入任务。FINISHED加载完成文件数据成功写入目标表此时记录FINISH_LOAD_TIME。ERROR加载失败文件加载过程中出错ERROR_MSG字段会携带具体的错误信息。SKIPPED已跳过源码支持文件被手动跳过不再参与加载。状态的迁移通过FileListRepo.updateFileState完成源码注释明确给出了四种更新场景LOADING开始加载任务、FINISHED成功完成加载、ERROR加载失败、SKIPPED手动跳过文件。在底层实现 FileListTableRepo.java 中这些更新分别对应不同的 SQL 语句开始加载时SET state LOADING, start_load now(), insert_label ...完成加载时SET state FINISHED, finish_load now()失败时SET state ERROR, error_info ...。四、查询实战如何监控 Pipe 的文件加载进度pipe_files属于 Information Schema 下的系统视图直接使用标准的SELECT语句查询即可无需任何权限或配置。以下是一些典型的运维查询场景。1. 查看某个 Pipe 管理的所有文件SELECT * FROM information_schema.pipe_files WHERE DATABASE_NAME your_db AND PIPE_NAME your_pipe;2. 查看当前仍处于加载中的文件SELECT PIPE_NAME, FILE_NAME, FILE_SIZE, START_LOAD_TIME FROM information_schema.pipe_files WHERE LOAD_STATE LOADING;3. 定位加载失败的文件与错误信息SELECT DATABASE_NAME, PIPE_NAME, FILE_NAME, FILE_VERSION, ERROR_MSG FROM information_schema.pipe_files WHERE LOAD_STATE ERROR;通过ERROR_MSG列可以拿到文件加载失败的具体原因这是排查导入问题最直接的入口。4. 监控文件的扫描与加载时延pipe_files提供了三个关键时间节点STAGED_TIME文件被 Pipe 发现、START_LOAD_TIME开始加载、FINISH_LOAD_TIME加载完成。通过它们可以计算发现到开始加载的排队耗时START_LOAD_TIME - STAGED_TIME单文件加载耗时FINISH_LOAD_TIME - START_LOAD_TIME。例如统计每个 Pipe 的平均文件加载耗时SELECT PIPE_NAME, AVG(TIMESTAMPDIFF(SECOND, START_LOAD_TIME, FINISH_LOAD_TIME)) AS avg_load_sec FROM information_schema.pipe_files WHERE LOAD_STATE FINISHED GROUP BY PIPE_NAME;5. 查看整体进度概览SELECT LOAD_STATE, COUNT(*) AS file_count FROM information_schema.pipe_files GROUP BY LOAD_STATE;该查询按状态汇总文件数量可以快速判断当前 Pipe 是否存在积压大量UNLOADED或失败大量ERROR的情况。五、源码原理pipe_files 背后的数据链路理解pipe_files的底层实现有助于更准确地解读视图中的每一列。整个数据链路可以分为三个层次第一层记录模型。PipeFileRecord.java 定义了 Pipe 文件清单中的一条记录字段与视图列一一对应pipeId、fileName、fileVersion、fileSize、loadState、lastModified、stagedTime、startLoadTime、finishLoadTime、errorMessage、insertLabel。值得注意的是FILE_VERSION的取值逻辑当底层文件系统支持 Etag如 S3 的EtagSource时取文件的 Etag 作为版本摘要否则取文件的修改时间modificationTime。这解释了为什么文档将FILE_VERSION描述为数据文件的 digest——它本质上是用于识别文件内容/版本变化的唯一标识。第二层持久化存储。记录并非直接暴露在内存中而是由FileListRepo体系持久化到一张内部表。默认实现是 FileListTableRepo.java 中的pipe_file_list表建表语句的关键结构如下CREATE TABLE IF NOT EXISTS db.pipe_file_list ( id bigint not null auto_increment, pipe_id bigint, file_name string, file_version string, file_size bigint, state string, last_modified datetime, staged_time datetime, start_load datetime, finish_load datetime, error_info string, insert_label string ) PRIMARY KEY(id) DISTRIBUTED BY HASH(id) BUCKETS 8 ORDER BY (pipe_id, file_name) properties(replication_num ...);可以看到内部表的列与视图字段几乎一一对应其中错误信息以 JSON 字符串存储在error_info列中。源码 PipeFileRecord.java 中定义了该 JSON 的结构包含errorMessage、errorCount、errorLine三个字段——也就是说除了视图暴露的ERROR_MSG外内部还保留了错误计数与错误行号等信息当前版本视图仅展示errorMessage部分。此外源码注释解释了为何不使用(pipe_id, file_name, file_version)作为主键当前主键实现将键长度限制为 128因此采用了自增id主键。第三层系统表映射。当 FE 启动或执行查询时information_schema中的pipe_files视图TSchemaTableType.SCH_PIPE_FILES会从上述内部表读取记录并转换为查询结果返回。这一层由 PipeFileSystemTable.java 负责列定义由PipeFileRecord提供 JSON 序列化/反序列化能力见其fromResultBatch与fromJson方法数据以{data: [col1, col2, col3, ...]}的 JSON 协议传输。此外单元测试 FileListRepoTest.java 覆盖了stageFiles、listFilesByState、updateFileState等核心链路包括将文件置为LOADING、FINISHED、ERROR等状态转换的验证PipeManagerTest.java 中也有对pipe_files查询结果非空的断言可作为理解该视图行为的参考。六、使用建议与注意事项结合文档定义与源码实现使用pipe_files时有以下几点值得留意版本要求pipe_files视图自 StarRocks v3.2 起支持使用时请确认集群版本满足要求。状态语义差异文档列出的LOAD_STATE有效值为UNLOADED、LOADING、FINISHED、ERROR而源码枚举中还包含SKIPPED状态用于表示被手动跳过的文件。如果你在查询中遇到了文档未列出的状态值可参照源码理解其含义。时间字段的可用性STAGED_TIME、START_LOAD_TIME、FINISH_LOAD_TIME分别对应文件生命周期中的不同节点未到达该节点时对应字段可能为空源码解析 JSON 时对 null/空串返回null。例如尚未开始加载的文件其START_LOAD_TIME与FINISH_LOAD_TIME为空。错误信息解读ERROR_MSG仅展示错误消息文本更细粒度的错误统计错误行号、错误数目前保留在内部记录模型与底层error_infoJSON 中尚未在视图中暴露。结合 Pipe 生命周期使用PIPE_ID是 Pipe 的唯一标识同一个数据库下 Pipe 名称唯一当 Pipe 被删除时其对应的文件记录也会随内部表清理见 FileListTableRepo.java 中的DELETE_BY_PIPE语句因此查询结果仅反映当前存在的 Pipe 及其文件状态。通过pipe_files视图你可以将 Pipe 这一黑盒式的持续导入过程透明化从文件被发现、排队、加载到最终完成或失败每一个环节都有据可查。当业务侧反馈数据未及时入库时先按LOAD_STATE分组查看状态分布再针对ERROR文件查看ERROR_MSG即可高效定位问题根因。【免费下载链接】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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价