spark-starrocks-logs.md 14 KB

Spark 和 StarRocks 日志链路说明

本文梳理当前仓库里 Spark 与 StarRocks 的“任务日志”实现,重点说明日志从哪里产生、在哪里采集、如何落盘、如何查询和下载。

这里的“日志”主要分成三层:

  1. 运行时日志:Spark driver pod、StarRocks Job pod、StarRocks FE pod 的输出。
  2. 任务日志:后端把运行时日志整理后,按任务实例存到 OSS。
  3. 执行结果:Spark SQL 结果、StarRocks SQL 结果,会单独存储并可下载。

总体链路

Spark 和 StarRocks 的日志链路都遵循同一个思路:

  1. 任务提交到 K8s orchestrator。
  2. orchestrator 从 K8s 里读取 Pod 日志或 Job 状态。
  3. 后端根据 orchestrator 返回的数据,提取“完整日志”和“结果文本”。
  4. 如果 OSS 已配置,后端把日志和结果归档到对象存储。
  5. 前端通过工作流实例页面查看日志,或下载日志/结果文件。

对应的核心接口在:

Spark 日志实现

1. 任务提交阶段

Spark SQL 任务由后端 SparkAdapter 组装 SparkApplication spec,再提交到 orchestrator。

  • submitTask() 会把 SQL、执行集群、Spark 配置等组装成 SparkApplication。
  • 若 runner 类型为 Python,会走 --sql 参数;否则走 Java runner 的 -e 参数。
  • SparkAdapter 会强制注入 Polaris/Iceberg 相关默认配置,避免集群配置覆盖后端管理的 catalog URI。

代码入口:

2. Spark SQL runner 的日志格式

Spark 任务真正执行 SQL 的代码在 orchestrator 的 Python runner:

这个 runner 会在 Spark 执行完成后,主动打印一组标记,保证后续解析稳定:

  • WENSHU_SPARK_SQL_BEGIN
  • SQL 原文
  • WENSHU_SPARK_SQL_RESULT_REF
  • 结果对象存储引用
  • WENSHU_SPARK_SQL_RESULT
  • 预览结果文本
  • WENSHU_SPARK_SQL_END

如果结果写入对象存储失败,runner 还会打印:

  • WENSHU_SPARK_SQL_RESULT_REF_ERROR ...

这套标记的目的,是让 orchestrator 后续能从 driver log 中精确切出 SQL、结果引用和结果正文。

3. orchestrator 采集 Spark 日志

orchestrator 的 GetSparkJobResultByOperation() 会:

  1. 根据 operation 找到对应 SparkApplication。
  2. 读取 driver pod 名称。
  3. 从 driver main container 中读取日志,最多 512 KB。
  4. 解析日志里的 Spark 标记。

对应实现:

解析结果会返回这些字段:

  • sql
  • result
  • resultRef
  • log
  • message
  • truncated

其中 truncated 会在两种情况下为真:

  • 读取日志超过限制。
  • 日志中缺少结束标记,说明内容不完整。

4. 后端归档 Spark 日志

后端的 TaskResultPersistenceService.persistTaskLog() 会在任务终态时调用 orchestrator,拿到 SparkJobResult,再决定是否持久化:

  • log 会写入 task-log/{taskInstanceId}
  • result 会写入 task-result/{taskInstanceId}

Spark 的完整日志优先级是:

  1. sparkJobResult.log()
  2. sparkJobResult.message()
  3. sparkJobResult.result()
  4. sparkJobResult.resultRef()
  5. 任务提交后的兜底提示文本

对应实现:

5. Spark 归档补偿任务

除了任务完成时的实时归档,Spark 还有一个单独的归档服务:

它的作用是:

  • 读取现有 task log/result。
  • 用 orchestrator 返回的日志内容补齐 OSS。
  • 避免用“提交中、等待日志”这类 fallback 文本覆盖已经存在的富日志。

归档逻辑里有一个重要策略:

  • 如果旧日志已经是富内容,新的 fallback 文本不会覆盖它。
  • 如果新日志更完整,则允许覆盖旧的 fallback 文本。

这可以避免把真实 driver log 冲掉。

6. 前端查看 Spark 日志

前端通过工作流执行页进入任务日志页面:

  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs
  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs/download

页面代码在:

前端显示逻辑是:

  • 优先展示 executionLog
  • 没有 executionLog 时回退到 logContent

也就是说,Spark 日志页面实际展示的是后端切分后的“执行日志段”,不是原始整段 driver log。

StarRocks 日志实现

1. 任务提交阶段

StarRocks SQL 任务由 StarRocksAdapter 提交。

它有两种执行模式:

  • ORCHESTRATOR
  • DIRECT

ORCHESTRATOR 模式

通过 orchestrator 提交 StarRocks Job,由 K8s Job 去执行 mysql 客户端命令。

DIRECT 模式

后端直接连 StarRocks FE 执行 JDBC SQL,并在本地拼装完整日志。

代码入口:

2. DIRECT 模式下的日志结构

DIRECT 模式会把日志拆成两部分再合成一个完整 payload:

  1. internal log
  2. execution result

internal log 会包含这些字段:

  • executionMode=JDBC
  • cluster
  • endpoint
  • taskId
  • startedAt
  • finishedAt
  • durationMs
  • statuserror
  • StarRocks 原生信息块

若有原生信息,还会附带:

  • queryId
  • get_query_profile
  • show_warnings
  • show_load_latest

StarRocks 的结果块会用下面的标记包起来:

  • WENSHU_STARROCKS_INTERNAL_LOG_BEGIN
  • WENSHU_STARROCKS_INTERNAL_LOG_END
  • WENSHU_STARROCKS_RESULT_BEGIN
  • WENSHU_STARROCKS_RESULT_END

后端和前端都依赖这组标记来拆分“执行日志”和“执行结果”。

3. orchestrator 采集 StarRocks 日志

orchestrator 的 GetStarRocksJobStatusByOperation() 会同时读取:

  1. StarRocks Job 对应的 pod 日志
  2. FE pod 日志

然后把两部分拼成一个 log 字段返回。

对应实现:

拼接规则:

  • FE 日志会按 pod 单独加前缀:[FE Pod] pod-name
  • 如果内容被 limit 截断,会在标题上标记 (truncated)
  • 最终整段日志会被包装进:
    • WENSHU_STARROCKS_INTERNAL_LOG_BEGIN/END
    • WENSHU_STARROCKS_RESULT_BEGIN/END

4. 后端保存 StarRocks 日志

TaskResultPersistenceService.persistTaskLog() 对 StarRocks 的处理和 Spark 类似,但多了一个结果提取逻辑:

  • log 直接作为完整日志候选
  • resultWENSHU_STARROCKS_RESULT_BEGIN/END 中截取

如果完整日志已经存在,后续的 fallback 文本不会覆盖它。

完整日志优先级大致是:

  1. starRocksJobStatus.log()
  2. starRocksJobStatus.message()
  3. 任务提交后的兜底提示文本

执行结果优先从日志里解析:

  • WENSHU_STARROCKS_RESULT_BEGIN
  • WENSHU_STARROCKS_RESULT_END

对应实现:

5. 前端查看 StarRocks 日志和结果

前端同样通过工作流任务详情页进入:

  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs
  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/results
  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/logs/download
  • GET /api/task-exec/workflows/instances/{workflowInstanceId}/tasks/{taskInstanceId}/results/download

前端页面在:

页面逻辑和 Spark 一致:

  • 日志页面展示 executionLog
  • 结果页面展示 executionResult

存储与下载

OSS key 规范

当前任务日志和结果统一存到 OSS,key 规则是:

  • 日志:task-log/{taskInstanceId}
  • 结果:task-result/{taskInstanceId}

对应代码:

支持的结果格式

结果下载逻辑不仅支持纯文本,还支持一些对象存储中的结果目录:

  • 单个 .txt
  • 单个 .csv
  • 单个 .parquet
  • 多个 parquet 文件目录
  • 多个 csv / csv.gz 文件目录

下载时会自动做转换:

  • parquet 目录会被转成 csv
  • csv 会补 UTF-8 BOM,方便 Excel 打开
  • 多文件目录会打成 zip

OSS 未配置时的行为

如果 polaris.oss 没有完整配置:

  • 日志和结果不会落 OSS
  • 查询接口仍会尽量返回当前阶段拼出来的 fallback 文本
  • 下载接口可能返回空或者失败

这意味着:

  • 能看到页面内容,不代表已经完成归档
  • 生产环境建议确保 OSS 配置完整,否则历史日志不可追溯

归档和清理

Spark 任务完成后,后端并不是立刻删掉 K8s 资源,而是由 scheduler / cleaner 分阶段处理:

  • SchedulerService 在任务成功、失败、超时等终态时会先调用 persistTaskLog()
  • SparkApplicationCleanerReconciler 在 SparkApplication 超过 TTL 后,会再次尝试归档日志,然后删除 SparkApplication

对应实现:

这个设计的核心目标是:

  • 先拿到日志,再删 K8s 资源
  • 如果首次归档失败,允许在 grace period 内重试

当前实现的几个特性和限制

特性

  • Spark 和 StarRocks 都有“完整日志 + 执行结果”两层数据。
  • 任务详情页可以直接查看日志,也可以下载原始文件。
  • 归档逻辑会尽量避免用 fallback 文本覆盖真实日志。
  • Spark 结果支持 resultRef,可以把大结果落到 OSS 再回传引用。

限制

  • Spark driver log 读取有 512 KB 限制,超过后会截断。
  • StarRocks Job/FE 日志读取也有字节上限,FE 日志合并后可能被截断。
  • 结果解析依赖固定标记,如果运行时输出被其它内容污染,解析会退化。
  • OSS 未配置时,历史日志和结果无法稳定保留。

排障建议

  1. 优先看工作流任务详情页里的 executionLogexecutionResult,因为这是后端已经标准化后的内容。
  2. Spark 任务如果只看到“waiting for driver logs”之类内容,通常说明实际 driver log 还没可读,或归档还没完成。
  3. StarRocks 任务如果只看到内部日志没有结果,重点检查 WENSHU_STARROCKS_RESULT_BEGIN/END 是否出现。
  4. 如果怀疑对象存储问题,直接检查 task-log/{taskInstanceId}task-result/{taskInstanceId} 是否存在。
  5. 如果前端页面空白但任务已完成,先确认 persistTaskLog() 是否在 scheduler / cleaner 路径上成功执行。

相关文件