
StarRocksinformation_schema.task_runs详解异步任务与物化视图刷新执行的观测指南【免费下载链接】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/starrockstask_runs是 StarRocks 提供的一个 Information Schema 系统视图用于记录异步任务TaskRun的执行元数据覆盖异步 ETL如SUBMIT TASK提交的 CTAS/INSERT与异步物化视图刷新两类场景。阅读本文后你将掌握task_runs的每个字段含义、EXTRA_MESSAGE中物化视图刷新细节的解析方法并能借助状态机与底层调度源码fe/fe-core/src/main/java/com/starrocks/scheduler/快速定位刷新失败、刷新过慢、分区越界等常见问题。什么是 task_runstask_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查询的 IDTASK_NAME任务名称CREATE_TIME任务创建时间FINISH_TIME任务完成时间STATE任务状态。有效值PENDING、RUNNING、FAILED、SUCCESS、MERGED、SKIPPEDCATALOG任务所属 CatalogDATABASE任务所属 DatabaseDEFINITION任务的 SQL 定义EXPIRE_TIME任务过期时间ERROR_CODE任务的错误码ERROR_MESSAGE任务的错误信息PROGRESS任务进度EXTRA_MESSAGE任务的附加信息例如异步物化视图创建任务中的分区信息PROPERTIES任务的属性JOB_ID任务的 Job IDPROCESS_TIME任务的处理时间TASK_SOURCE提交任务的来源。有效值CTAS、MV、INSERT、PIPE、DATACACHE_SELECT。历史遗留记录未记录来源返回UNKNOWNSTATE 状态详解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 runEXTRA_MESSAGE字段会包含物化视图 task run 的明细消息。你可以通过该字段获取刷新范围、计划/实际刷新分区、执行选项、优化器诊断信息等结构化数据。更完整的说明见 materialized_view_task_run_details。EXTRA_MESSAGE 中的字段EXTRA_MESSAGE是MVTaskRunExtraMessage对象的 JSON 序列化结果对应源码 MVTaskRunExtraMessage.java主要包含以下字段字段类型说明forceRefreshBoolean是否强制刷新。手动执行REFRESH MATERIALIZED VIEW ... FORCE触发强制全量刷新时返回truepartitionStartString本次刷新的起始分区边界下界例如2024-01-01partitionEndString本次刷新的结束分区边界上界例如2024-01-31mvPartitionsToRefreshSetString本次 task run 计划刷新的物化视图自身分区列表例如[p20240101,p20240102]。注意不是基表分区refBasePartitionsToRefreshMapMapString, SetString优化前计划扫描的基表分区映射{tableName - SetpartitionName}basePartitionsToRefreshMapMapString, SetString物化视图版本映射提交后实际扫描的基表分区映射反映优化器真实使用的分区集合nextPartitionStart/nextPartitionEndString下一次增量刷新的起始/结束边界。当一次刷新因资源限制或数据量过大被拆分为多个 task run 时用于标识剩余待刷分区范围nextPartitionValuesString下一次刷新的序列化分区值用于列表分区或复杂分区方案例如(US, ACTIVE), (UK, ACTIVE)processStartTimeInteger毫秒时间戳task run 实际开始处理的时间不含排队等待时间executeOptionExecuteOption 对象任务执行选项。默认Priority LOWEST、isMergeRedundant falseplanBuilderMessageMapString, String查询计划构建器的诊断消息与元数据包含查询规划、优化决策与潜在问题refreshModeString刷新模式COMPLETE全量、PARTIAL增量、FORCE强制、默认/未指定adaptivePartitionRefreshNumberInteger自适应分区刷新时每轮迭代刷新的分区数默认-1未启用自适应executeOption 与任务优先级executeOption对象包含以下字段对应源码 ExecuteOption.javapriority任务执行优先级取值为Constants.TaskRunPriority。注意数值越大优先级越高与直觉相反。文档侧给出的枚举值HIGHEST: 0 →LOWEST: 127为执行顺序排列实际调度语义应以 Constants.java 为准——其中LOWEST(0)、LOW(20)、NORMAL(50)、HIGH(80)、HIGHER(90)、HIGHEST(100)默认优先级为LOWESTisMergeRedundant是否合并冗余刷新操作布尔值properties附加执行属性格式为MapString, String。refBasePartitionsToRefreshMap 与 basePartitionsToRefreshMap 的区别两者容易混淆官方文档强调如下refBasePartitionsToRefreshMap优化前的计划分区通常针对主引用基表basePartitionsToRefreshMap优化后的实际分区包含所有表与优化后的分区集合。对比两个映射可以判断查询优化器是否改变了分区裁剪计划是排查“刷新了意外分区”的核心手段。分区数量截断为避免元数据存储无限膨胀mvPartitionsToRefresh、refBasePartitionsToRefreshMap、basePartitionsToRefreshMap、planBuilderMessage中的分区数量会被自动截断到 FE 配置项max_mv_task_run_meta_message_values_length默认 100以内。源码 MVTaskRunExtraMessage.java 在写入这些字段时统一使用该配置做长度限制。查询 task_runstask_runs是 Information Schema 下的系统视图直接使用SELECT即可查询-- 查看所有异步任务运行记录 SELECT * FROM INFORMATION_SCHEMA.task_runs; -- 按任务名精确过滤 SELECT * FROM information_schema.task_runs WHERE task_name task_name;查询物化视图刷新细节由于物化视图任务名通常以mv-为前缀可以按名称过滤并展开EXTRA_MESSAGESELECT 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 的有效期超过后自动删除FAILED与SUCCESS状态的记录也会被自动清理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可用于以下场景。监控刷新性能对比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 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),仅供参考