StarRocks `information_schema.task_runs` 详解:异步任务与物化视图刷新执行的观测指南
2026/9/17 23:53:42 网站建设 项目流程

StarRocksinformation_schema.task_runs详解:异步任务与物化视图刷新执行的观测指南

【免费下载链接】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

task_runs是 StarRocks 提供的一个 Information Schema 系统视图,用于记录异步任务(TaskRun)的执行元数据,覆盖异步 ETL(如SUBMIT TASK提交的 CTAS/INSERT)与异步物化视图刷新两类场景。阅读本文后,你将掌握task_runs的每个字段含义、EXTRA_MESSAGE中物化视图刷新细节的解析方法,并能借助状态机与底层调度源码(fe/fe-core/src/main/java/com/starrocks/scheduler/)快速定位刷新失败、刷新过慢、分区越界等常见问题。

什么是 task_runs

task_runs(即INFORMATION_SCHEMA.task_runs)为 StarRocks 提供异步任务执行的信息。每一条 TaskRun 记录都由以下两类语句之一产生:

  • SUBMIT TASK:将CREATE TABLE AS SELECTINSERTCACHE SELECT等 ETL 语句作为异步任务提交;
  • CREATE MATERIALIZED VIEW:创建异步物化视图后,系统按需/按周期触发的刷新任务。

从源码结构看(fe/fe-core/src/main/java/com/starrocks/scheduler/),StarRocks 的异步任务体系由Task(任务模板)、TaskRun(一次具体执行)两层构成:Task保存 ETL 语句定义与调度属性,每次触发执行会生成一个TaskRun,由TaskRunScheduler调度、TaskRunExecutor执行,而TaskRunManager负责管理与记录,最终把执行元数据写入task_runs视图。因此,查询task_runs实际上是查询 StarRocks FE 侧任务运行历史的最直接入口。

:::note 一个物化视图的刷新操作可能生成多个task run:每个 task run 代表一个按partition_refresh_number配置切分出来的刷新子任务。理解这一点对排查“一次刷新为什么产生多条记录”至关重要。 :::

字段说明

task_runs提供以下字段:

字段说明
QUERY_ID查询的 ID
TASK_NAME任务名称
CREATE_TIME任务创建时间
FINISH_TIME任务完成时间
STATE任务状态。有效值:PENDINGRUNNINGFAILEDSUCCESSMERGEDSKIPPED
CATALOG任务所属 Catalog
DATABASE任务所属 Database
DEFINITION任务的 SQL 定义
EXPIRE_TIME任务过期时间
ERROR_CODE任务的错误码
ERROR_MESSAGE任务的错误信息
PROGRESS任务进度
EXTRA_MESSAGE任务的附加信息,例如异步物化视图创建任务中的分区信息
PROPERTIES任务的属性
JOB_ID任务的 Job ID
PROCESS_TIME任务的处理时间
TASK_SOURCE提交任务的来源。有效值:CTASMVINSERTPIPEDATACACHE_SELECT。历史遗留记录(未记录来源)返回UNKNOWN

STATE 状态详解

STATE字段的六种取值与调度源码中的TaskRunState枚举一一对应(见 Constants.java),其转移关系如下:

PENDING -> RUNNING -> SUCCESS | |-----> FAILED | |-----> SKIPPED |-----> MERGED
  • PENDING:任务已进入等待队列,尚未开始执行;
  • RUNNING:任务正在执行;
  • FAILED:任务执行失败;
  • SUCCESS:任务执行成功;
  • MERGED:仅用于物化视图刷新任务。当新的刷新任务提交时,若旧任务仍停留在 PENDING 队列中,这两个任务会被合并(冗余刷新去重),合并后的任务保持原有优先级;
  • SKIPPED:仅用于物化视图刷新任务。当基表分区上没有检测到数据变化时,对应物化视图分区的刷新会被跳过。

Constants.java中的辅助方法可以进一步确认其语义:isSuccessState()SUCCESSMERGEDSKIPPED均视为终态成功(其中MERGED/SKIPPED属于“被合并/被跳过”的未实际执行成功态);而isFailedState()则覆盖FAILED状态。

TASK_SOURCE 来源说明

TASK_SOURCE标识了提交该任务的上游来源,源码中TaskSource枚举与之一致:

  • CTAS:来自异步CREATE TABLE AS SELECT任务;
  • MV:来自物化视图刷新任务;
  • INSERT:来自异步INSERT任务;
  • PIPE:来自 Pipe 导入管道(数据管道持续导入场景);
  • DATACACHE_SELECT:来自数据缓存预热(CACHE SELECT)任务;
  • UNKNOWN:遗留记录(写入时未记录来源)。

EXTRA_MESSAGE 字段

对于物化视图 task run,EXTRA_MESSAGE字段会包含物化视图 task run 的明细消息。你可以通过该字段获取刷新范围、计划/实际刷新分区、执行选项、优化器诊断信息等结构化数据。更完整的说明见 materialized_view_task_run_details。

EXTRA_MESSAGE 中的字段

EXTRA_MESSAGEMVTaskRunExtraMessage对象的 JSON 序列化结果(对应源码 MVTaskRunExtraMessage.java),主要包含以下字段:

字段类型说明
forceRefreshBoolean是否强制刷新。手动执行REFRESH MATERIALIZED VIEW ... FORCE触发强制全量刷新时返回true
partitionStartString本次刷新的起始分区边界(下界),例如"2024-01-01"
partitionEndString本次刷新的结束分区边界(上界),例如"2024-01-31"
mvPartitionsToRefreshSet<String>本次 task run 计划刷新的物化视图自身分区列表,例如["p20240101","p20240102"]。注意:不是基表分区
refBasePartitionsToRefreshMapMap<String, Set<String>>优化前计划扫描的基表分区映射{tableName -> Set<partitionName>}
basePartitionsToRefreshMapMap<String, Set<String>>物化视图版本映射提交后实际扫描的基表分区映射,反映优化器真实使用的分区集合
nextPartitionStart/nextPartitionEndString下一次增量刷新的起始/结束边界。当一次刷新因资源限制或数据量过大被拆分为多个 task run 时,用于标识剩余待刷分区范围
nextPartitionValuesString下一次刷新的序列化分区值,用于列表分区或复杂分区方案,例如"('US', 'ACTIVE'), ('UK', 'ACTIVE')"
processStartTimeInteger(毫秒时间戳)task run 实际开始处理的时间(不含排队等待时间)
executeOptionExecuteOption 对象任务执行选项。默认Priority = LOWESTisMergeRedundant = false
planBuilderMessageMap<String, String>查询计划构建器的诊断消息与元数据,包含查询规划、优化决策与潜在问题
refreshModeString刷新模式:"COMPLETE"(全量)、"PARTIAL"(增量)、"FORCE"(强制)、""(默认/未指定)
adaptivePartitionRefreshNumberInteger自适应分区刷新时每轮迭代刷新的分区数,默认-1(未启用自适应)

executeOption 与任务优先级

executeOption对象包含以下字段(对应源码 ExecuteOption.java):

  • priority:任务执行优先级,取值为Constants.TaskRunPriority。注意:数值越大优先级越高,与直觉相反。文档侧给出的枚举值(HIGHEST: 0 →LOWEST: 127)为执行顺序排列,实际调度语义应以 Constants.java 为准——其中LOWEST(0)LOW(20)NORMAL(50)HIGH(80)HIGHER(90)HIGHEST(100),默认优先级为LOWEST
  • isMergeRedundant:是否合并冗余刷新操作(布尔值);
  • properties:附加执行属性,格式为Map<String, String>

refBasePartitionsToRefreshMap 与 basePartitionsToRefreshMap 的区别

两者容易混淆,官方文档强调如下:

  • refBasePartitionsToRefreshMap优化前的计划分区(通常针对主引用基表);
  • basePartitionsToRefreshMap优化后的实际分区(包含所有表与优化后的分区集合)。

对比两个映射,可以判断查询优化器是否改变了分区裁剪计划,是排查“刷新了意外分区”的核心手段。

分区数量截断

为避免元数据存储无限膨胀,mvPartitionsToRefreshrefBasePartitionsToRefreshMapbasePartitionsToRefreshMapplanBuilderMessage中的分区数量会被自动截断到 FE 配置项max_mv_task_run_meta_message_values_length(默认 100)以内。源码 MVTaskRunExtraMessage.java 在写入这些字段时统一使用该配置做长度限制。

查询 task_runs

task_runs是 Information Schema 下的系统视图,直接使用SELECT即可查询:

-- 查看所有异步任务运行记录 SELECT * FROM INFORMATION_SCHEMA.task_runs; -- 按任务名精确过滤 SELECT * FROM information_schema.task_runs WHERE task_name = '<task_name>';

查询物化视图刷新细节

由于物化视图任务名通常以mv-为前缀,可以按名称过滤并展开EXTRA_MESSAGE

SELECT TASK_NAME, CREATE_TIME, FINISH_TIME, STATE, EXTRA_MESSAGE FROM information_schema.task_runs WHERE TASK_NAME LIKE 'mv-%' ORDER BY CREATE_TIME DESC LIMIT 10;

EXTRA_MESSAGE列存放MVTaskRunExtraMessage的 JSON 表示,可使用 JSON 函数解析出可读字段:

SELECT TASK_NAME, CREATE_TIME, get_json_string(EXTRA_MESSAGE, '$.refreshMode') AS refresh_mode, get_json_string(EXTRA_MESSAGE, '$.forceRefresh') AS force_refresh, get_json_string(EXTRA_MESSAGE, '$.mvPartitionsToRefresh') AS mv_partitions, get_json_int(EXTRA_MESSAGE, '$.processStartTime') AS process_start_ms, get_json_int(EXTRA_MESSAGE, '$.adaptivePartitionRefreshNumber') AS adaptive_batch_size FROM information_schema.task_runs WHERE TASK_NAME = 'mv-12345' ORDER BY CREATE_TIME DESC;

计算实际处理时间

processStartTime排除了排队时间,因此实际处理时长的计算公式为:

SELECT TASK_NAME, FINISH_TIME, get_json_bigint(EXTRA_MESSAGE, '$.processStartTime') AS process_start_time, (unix_timestamp(FINISH_TIME) * 1000 - get_json_bigint(EXTRA_MESSAGE, '$.processStartTime')) / 1000 AS processing_seconds FROM information_schema.task_runs WHERE TASK_NAME LIKE 'mv-%' AND STATE = 'SUCCESS';

分析分区刷新模式

SELECT TASK_NAME, CREATE_TIME, get_json_string(EXTRA_MESSAGE, '$.partitionStart') AS start_partition, get_json_string(EXTRA_MESSAGE, '$.partitionEnd') AS end_partition, get_json_string(EXTRA_MESSAGE, '$.nextPartitionStart') AS next_start, get_json_string(EXTRA_MESSAGE, '$.nextPartitionEnd') AS next_end FROM information_schema.task_runs WHERE TASK_NAME = 'mv-12345' ORDER BY CREATE_TIME DESC;

与 SUBMIT TASK 异步任务的联动

task_runs记录了SUBMIT TASK提交的异步 ETL 任务的执行历史。执行SUBMIT TASK会创建一个Task(任务模板,可含SCHEDULE EVERY(INTERVAL ...)周期性调度),每次触发执行则生成一个TaskRun

-- 提交一个异步 CTAS 任务 SUBMIT TASK etl0 AS CREATE TABLE tbl1 AS SELECT * FROM src_tbl; -- 提交一个周期性执行的 INSERT OVERWRITE 任务(每 1 分钟执行一次) SUBMIT TASK SCHEDULE EVERY(INTERVAL 1 MINUTE) AS INSERT OVERWRITE insert_wiki_edit SELECT dt, user_id, count(*) FROM source_wiki_edit GROUP BY dt, user_id;

任务模板信息查询 tasks 视图,任务执行历史则查询本文介绍的task_runs视图:

SELECT * FROM INFORMATION_SCHEMA.tasks WHERE task_name = '<task_name>'; SELECT * FROM information_schema.task_runs WHERE task_name = '<task_name>';

相关 FE 配置项

SUBMIT TASK异步任务的运行行为受以下 FE 配置控制(详见 SUBMIT_TASK),它们直接决定了task_runs中记录的生命周期与并发规模:

配置项默认值说明
task_ttl_second86400Task(一次性任务)的有效期,超过后被删除,单位秒
task_check_interval_second3600清理无效 Task 的间隔,单位秒
task_runs_ttl_second86400TaskRun 的有效期,超过后自动删除;FAILEDSUCCESS状态的记录也会被自动清理
task_runs_concurrency4可并行执行的 TaskRun 最大数量
task_runs_queue_length500等待执行的 TaskRun 最大排队数量,超过后新任务将被挂起
task_runs_max_history_number10000保留的 TaskRun 记录最大条数
task_min_schedule_interval_s10任务执行的最小调度间隔,单位秒

这些配置解释了task_runs记录为何不是永久保留:默认 TTL 为 86400 秒(24 小时),历史记录上限 10000 条,超出后旧记录会被清理。

物化视图刷新的性能分析与问题排查

结合 materialized_view_task_run_details 中的最佳实践与故障排查建议,EXTRA_MESSAGE可用于以下场景。

监控刷新性能

  • 对比processStartTimeFINISH_TIME:若两者差距大而CREATE_TIMEprocessStartTime差距也大,说明 task run 在队列中等待较久;
  • 使用adaptivePartitionRefreshNumber优化批量大小(对应源码MVTaskRunExtraMessage.adaptivePartitionRefreshNumber,默认-1表示未启用自适应刷新,见 MVTaskRunExtraMessage.java)。

调试失败的刷新

  • 检查planBuilderMessage中是否有优化器相关的问题;
  • 对比refBasePartitionsToRefreshMap(计划分区)与basePartitionsToRefreshMap(实际分区),定位分区裁剪异常。

优化增量刷新

  • 监控nextPartitionStart/nextPartitionEnd,理解多轮迭代刷新模式;若刷新频繁横跨多个 task run,可调整分区粒度(如调大partition_refresh_number,见源码 MVPCTRefreshPartitioner.java 中对partition_refresh_number的读取逻辑);
  • partition_refresh_number未显式设置时,会回退到default_mv_partition_refresh_number配置。

排查典型问题

问题一:刷新耗时过长。依次检查:

  1. processStartTime——与创建时间差距大说明任务长时间排队;
  2. basePartitionsToRefreshMap——分区数量过大说明扫描了过多分区;
  3. adaptivePartitionRefreshNumber——可能需要调整负载或批量参数。

问题二:刷新了意外的分区。依次检查:

  1. forceRefresh——为true说明执行了强制全量刷新;
  2. refBasePartitionsToRefreshMap——计划分区;
  3. basePartitionsToRefreshMap——优化后的实际分区;
  4. 对比两个映射,确认优化器是否改变了执行计划。

问题三:刷新卡在多轮迭代。依次检查:

  1. nextPartitionStart/nextPartitionEnd——反映未完成的刷新状态;
  2. adaptivePartitionRefreshNumber——可能需要调整负载;
  3. 考虑增大批处理大小或减小分区粒度。

配置项:max_mv_task_run_meta_message_values_length

  • 类型:Integer
  • 默认值:100
  • 作用域:FE 配置
  • 说明:限制 Set 或 Map 字段中存储的最大条目数,防止元数据过度增长。约束的对象包括mvPartitionsToRefreshrefBasePartitionsToRefreshMapbasePartitionsToRefreshMapplanBuilderMessage

小结

information_schema.task_runs是 StarRocks 异步任务体系的观测窗口:对于SUBMIT TASK异步 ETL,它记录每次运行的六态状态机与生命周期;对于异步物化视图刷新,它的EXTRA_MESSAGE承载了从计划分区、实际分区、强制刷新标记到自适应批量大小、下一轮刷新边界的全量刷新细节。结合 Constants.java 的状态枚举与 MVTaskRunExtraMessage.java 的序列化实现,开发者可以在不改动任何配置的情况下,用标准 SQL 完成刷新性能监控、失败诊断与增量策略调优。

延伸阅读

  • materialized_view_task_run_details:物化视图 task run 明细字段的完整说明
  • SUBMIT TASK:异步 ETL 任务提交语法与 FE 配置
  • CREATE MATERIALIZED VIEW:异步物化视图的创建语法
  • REFRESH MATERIALIZED VIEW:物化视图手动刷新语法

【免费下载链接】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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询