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 SELECT、INSERT、CACHE 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 | 任务状态。有效值:PENDING、RUNNING、FAILED、SUCCESS、MERGED、SKIPPED |
| CATALOG | 任务所属 Catalog |
| DATABASE | 任务所属 Database |
| DEFINITION | 任务的 SQL 定义 |
| EXPIRE_TIME | 任务过期时间 |
| ERROR_CODE | 任务的错误码 |
| ERROR_MESSAGE | 任务的错误信息 |
| PROGRESS | 任务进度 |
| EXTRA_MESSAGE | 任务的附加信息,例如异步物化视图创建任务中的分区信息 |
| PROPERTIES | 任务的属性 |
| JOB_ID | 任务的 Job ID |
| PROCESS_TIME | 任务的处理时间 |
| TASK_SOURCE | 提交任务的来源。有效值:CTAS、MV、INSERT、PIPE、DATACACHE_SELECT。历史遗留记录(未记录来源)返回UNKNOWN |
STATE 状态详解
STATE字段的六种取值与调度源码中的TaskRunState枚举一一对应(见 Constants.java),其转移关系如下:
PENDING -> RUNNING -> SUCCESS | |-----> FAILED | |-----> SKIPPED |-----> MERGEDPENDING:任务已进入等待队列,尚未开始执行;RUNNING:任务正在执行;FAILED:任务执行失败;SUCCESS:任务执行成功;MERGED:仅用于物化视图刷新任务。当新的刷新任务提交时,若旧任务仍停留在 PENDING 队列中,这两个任务会被合并(冗余刷新去重),合并后的任务保持原有优先级;SKIPPED:仅用于物化视图刷新任务。当基表分区上没有检测到数据变化时,对应物化视图分区的刷新会被跳过。
从Constants.java中的辅助方法可以进一步确认其语义:isSuccessState()将SUCCESS、MERGED、SKIPPED均视为终态成功(其中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_MESSAGE是MVTaskRunExtraMessage对象的 JSON 序列化结果(对应源码 MVTaskRunExtraMessage.java),主要包含以下字段:
| 字段 | 类型 | 说明 |
|---|---|---|
forceRefresh | Boolean | 是否强制刷新。手动执行REFRESH MATERIALIZED VIEW ... FORCE触发强制全量刷新时返回true |
partitionStart | String | 本次刷新的起始分区边界(下界),例如"2024-01-01" |
partitionEnd | String | 本次刷新的结束分区边界(上界),例如"2024-01-31" |
mvPartitionsToRefresh | Set<String> | 本次 task run 计划刷新的物化视图自身分区列表,例如["p20240101","p20240102"]。注意:不是基表分区 |
refBasePartitionsToRefreshMap | Map<String, Set<String>> | 优化前计划扫描的基表分区映射{tableName -> Set<partitionName>} |
basePartitionsToRefreshMap | Map<String, Set<String>> | 物化视图版本映射提交后实际扫描的基表分区映射,反映优化器真实使用的分区集合 |
nextPartitionStart/nextPartitionEnd | String | 下一次增量刷新的起始/结束边界。当一次刷新因资源限制或数据量过大被拆分为多个 task run 时,用于标识剩余待刷分区范围 |
nextPartitionValues | String | 下一次刷新的序列化分区值,用于列表分区或复杂分区方案,例如"('US', 'ACTIVE'), ('UK', 'ACTIVE')" |
processStartTime | Integer(毫秒时间戳) | task run 实际开始处理的时间(不含排队等待时间) |
executeOption | ExecuteOption 对象 | 任务执行选项。默认Priority = LOWEST、isMergeRedundant = false |
planBuilderMessage | Map<String, String> | 查询计划构建器的诊断消息与元数据,包含查询规划、优化决策与潜在问题 |
refreshMode | String | 刷新模式:"COMPLETE"(全量)、"PARTIAL"(增量)、"FORCE"(强制)、""(默认/未指定) |
adaptivePartitionRefreshNumber | Integer | 自适应分区刷新时每轮迭代刷新的分区数,默认-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:优化后的实际分区(包含所有表与优化后的分区集合)。
对比两个映射,可以判断查询优化器是否改变了分区裁剪计划,是排查“刷新了意外分区”的核心手段。
分区数量截断
为避免元数据存储无限膨胀,mvPartitionsToRefresh、refBasePartitionsToRefreshMap、basePartitionsToRefreshMap、planBuilderMessage中的分区数量会被自动截断到 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_second | 86400 | Task(一次性任务)的有效期,超过后被删除,单位秒 |
task_check_interval_second | 3600 | 清理无效 Task 的间隔,单位秒 |
task_runs_ttl_second | 86400 | TaskRun 的有效期,超过后自动删除;FAILED与SUCCESS状态的记录也会被自动清理 |
task_runs_concurrency | 4 | 可并行执行的 TaskRun 最大数量 |
task_runs_queue_length | 500 | 等待执行的 TaskRun 最大排队数量,超过后新任务将被挂起 |
task_runs_max_history_number | 10000 | 保留的 TaskRun 记录最大条数 |
task_min_schedule_interval_s | 10 | 任务执行的最小调度间隔,单位秒 |
这些配置解释了task_runs记录为何不是永久保留:默认 TTL 为 86400 秒(24 小时),历史记录上限 10000 条,超出后旧记录会被清理。
物化视图刷新的性能分析与问题排查
结合 materialized_view_task_run_details 中的最佳实践与故障排查建议,EXTRA_MESSAGE可用于以下场景。
监控刷新性能
- 对比
processStartTime与FINISH_TIME:若两者差距大而CREATE_TIME与processStartTime差距也大,说明 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配置。
排查典型问题
问题一:刷新耗时过长。依次检查:
processStartTime——与创建时间差距大说明任务长时间排队;basePartitionsToRefreshMap——分区数量过大说明扫描了过多分区;adaptivePartitionRefreshNumber——可能需要调整负载或批量参数。
问题二:刷新了意外的分区。依次检查:
forceRefresh——为true说明执行了强制全量刷新;refBasePartitionsToRefreshMap——计划分区;basePartitionsToRefreshMap——优化后的实际分区;- 对比两个映射,确认优化器是否改变了执行计划。
问题三:刷新卡在多轮迭代。依次检查:
nextPartitionStart/nextPartitionEnd——反映未完成的刷新状态;adaptivePartitionRefreshNumber——可能需要调整负载;- 考虑增大批处理大小或减小分区粒度。
配置项:max_mv_task_run_meta_message_values_length
- 类型:Integer
- 默认值:100
- 作用域:FE 配置
- 说明:限制 Set 或 Map 字段中存储的最大条目数,防止元数据过度增长。约束的对象包括
mvPartitionsToRefresh、refBasePartitionsToRefreshMap、basePartitionsToRefreshMap与planBuilderMessage。
小结
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),仅供参考