本文梳理当前仓库里 Spark 与 StarRocks 的“任务日志”实现,重点说明日志从哪里产生、在哪里采集、如何落盘、如何查询和下载。
这里的“日志”主要分成三层:
Spark 和 StarRocks 的日志链路都遵循同一个思路:
对应的核心接口在:
Spark SQL 任务由后端 SparkAdapter 组装 SparkApplication spec,再提交到 orchestrator。
submitTask() 会把 SQL、执行集群、Spark 配置等组装成 SparkApplication。--sql 参数;否则走 Java runner 的 -e 参数。SparkAdapter 会强制注入 Polaris/Iceberg 相关默认配置,避免集群配置覆盖后端管理的 catalog URI。代码入口:
Spark 任务真正执行 SQL 的代码在 orchestrator 的 Python runner:
这个 runner 会在 Spark 执行完成后,主动打印一组标记,保证后续解析稳定:
WENSHU_SPARK_SQL_BEGINWENSHU_SPARK_SQL_RESULT_REFWENSHU_SPARK_SQL_RESULTWENSHU_SPARK_SQL_END如果结果写入对象存储失败,runner 还会打印:
WENSHU_SPARK_SQL_RESULT_REF_ERROR ...这套标记的目的,是让 orchestrator 后续能从 driver log 中精确切出 SQL、结果引用和结果正文。
orchestrator 的 GetSparkJobResultByOperation() 会:
对应实现:
解析结果会返回这些字段:
sqlresultresultReflogmessagetruncated其中 truncated 会在两种情况下为真:
后端的 TaskResultPersistenceService.persistTaskLog() 会在任务终态时调用 orchestrator,拿到 SparkJobResult,再决定是否持久化:
log 会写入 task-log/{taskInstanceId}result 会写入 task-result/{taskInstanceId}Spark 的完整日志优先级是:
sparkJobResult.log()sparkJobResult.message()sparkJobResult.result()sparkJobResult.resultRef()对应实现:
除了任务完成时的实时归档,Spark 还有一个单独的归档服务:
它的作用是:
归档逻辑里有一个重要策略:
这可以避免把真实 driver log 冲掉。
前端通过工作流执行页进入任务日志页面:
GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logsGET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs/download页面代码在:
前端显示逻辑是:
executionLogexecutionLog 时回退到 logContent也就是说,Spark 日志页面实际展示的是后端切分后的“执行日志段”,不是原始整段 driver log。
StarRocks SQL 任务由 StarRocksAdapter 提交。
它有两种执行模式:
ORCHESTRATORDIRECT通过 orchestrator 提交 StarRocks Job,由 K8s Job 去执行 mysql 客户端命令。
后端直接连 StarRocks FE 执行 JDBC SQL,并在本地拼装完整日志。
代码入口:
DIRECT 模式会把日志拆成两部分再合成一个完整 payload:
internal log 会包含这些字段:
executionMode=JDBCclusterendpointtaskIdstartedAtfinishedAtdurationMsstatus 或 error若有原生信息,还会附带:
queryIdget_query_profileshow_warningsshow_load_latestStarRocks 的结果块会用下面的标记包起来:
WENSHU_STARROCKS_INTERNAL_LOG_BEGINWENSHU_STARROCKS_INTERNAL_LOG_ENDWENSHU_STARROCKS_RESULT_BEGINWENSHU_STARROCKS_RESULT_END后端和前端都依赖这组标记来拆分“执行日志”和“执行结果”。
orchestrator 的 GetStarRocksJobStatusByOperation() 会同时读取:
然后把两部分拼成一个 log 字段返回。
对应实现:
拼接规则:
[FE Pod] pod-name(truncated)WENSHU_STARROCKS_INTERNAL_LOG_BEGIN/ENDWENSHU_STARROCKS_RESULT_BEGIN/ENDTaskResultPersistenceService.persistTaskLog() 对 StarRocks 的处理和 Spark 类似,但多了一个结果提取逻辑:
log 直接作为完整日志候选result 从 WENSHU_STARROCKS_RESULT_BEGIN/END 中截取如果完整日志已经存在,后续的 fallback 文本不会覆盖它。
完整日志优先级大致是:
starRocksJobStatus.log()starRocksJobStatus.message()执行结果优先从日志里解析:
WENSHU_STARROCKS_RESULT_BEGINWENSHU_STARROCKS_RESULT_END对应实现:
前端同样通过工作流任务详情页进入:
GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logsGET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/resultsGET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs/downloadGET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/results/download前端页面在:
页面逻辑和 Spark 一致:
executionLogexecutionResult当前任务日志和结果统一存到 OSS,key 规则是:
task-log/{taskInstanceId}task-result/{taskInstanceId}对应代码:
结果下载逻辑不仅支持纯文本,还支持一些对象存储中的结果目录:
.txt.csv.parquet下载时会自动做转换:
如果 polaris.oss 没有完整配置:
这意味着:
Spark 任务完成后,后端并不是立刻删掉 K8s 资源,而是由 scheduler / cleaner 分阶段处理:
SchedulerService 在任务成功、失败、超时等终态时会先调用 persistTaskLog()SparkApplicationCleanerReconciler 在 SparkApplication 超过 TTL 后,会再次尝试归档日志,然后删除 SparkApplication对应实现:
这个设计的核心目标是:
resultRef,可以把大结果落到 OSS 再回传引用。executionLog 和 executionResult,因为这是后端已经标准化后的内容。WENSHU_STARROCKS_RESULT_BEGIN/END 是否出现。task-log/{taskInstanceId} 和 task-result/{taskInstanceId} 是否存在。persistTaskLog() 是否在 scheduler / cleaner 路径上成功执行。