开发者概览

指标框架

qianmoQqianmoQ· 更新于 2026-10-02· 阅读 15 分钟· 0 次阅读

登录后可跨设备保存划线和私人笔记登录

指标框架

本文档说明 Velox 算子指标如何映射回 Gluten Spark SQL 指标。该映射包含三个有序步骤:

  1. 原生代码将 Velox 计划树化为 orderedNodeIds。
  2. Scala 将 JSON 负载解析为扁平结构 JList[OperatorMetrics]。
  3. Scala 将 Spark 计划树化为 MetricsUpdaterTree,并按该顺序消费扁平化后的原生指标。

JSON 传输方式使 JNI 边界保持小巧且稳定。C++ 以确定的顺序报告具名 Velox 统计信息;由 Scala 负责将这些统计信息映射为 Gluten 算子指标。

映射概览

指标映射会关联同一次执行的三个视图:

  • 原生执行产生的 Velox 计划节点 id 与任务统计信息。
  • 记录在 operatorToRelsMap 中的 Substrait rel id。
  • Spark 物理算子,每个算子都有一个 MetricsUpdater。

在规划阶段,Gluten 会为每个转换算子分配一个算子 id,并记录为该算子生成的 Substrait rel id:

operatorToRelsMap: Spark operator id -> Substrait rel ids

执行完成后,原生代码会按照 orderedNodeIds 顺序序列化 Velox 统计信息。Scala 将该 JSON 解析为:

Velox JSON node stats -> JList[OperatorMetrics]

最后,MetricsUtil.updateTransformerMetricsInternal 会同时遍历 MetricsUpdaterTree、operatorToRelsMap 以及原生指标列表。当前实现从末尾同时消费 Spark 算子 id 和原生指标索引:

operatorIdx = relMap.size() - 1
metricsIdx  = nativeMetrics.size() - 1

每个 Spark 算子都会消费与其 Substrait rel id 相对应的原生指标套件,对其进行合并或解释,然后将最终数值写入 Spark SQLMetrics。

第 1 步:原生 treefy 转换为 orderedNodeIds

orderedNodeIds 由 WholeStageResultIterator::getOrderedNodeIds 生成。这是对 Velox 计划执行的一次原生 treefy 过程,它将 Velox 计划树转换为一个确定性的节点 id 列表,Scala 之后可以将其展平为 OperatorMetrics。

此步骤是必需的,因为 Velox 任务统计信息是以计划节点 id 为键的。统计信息的映射结构无法保留 Gluten 指标更新器树所需的遍历顺序,必须由原生代码显式提供该顺序。

对于普通的 Velox 节点,getOrderedNodeIds 执行后序遍历:

visit all sources first
then append current node id

这样 Scala 就得到了一个与 Substrait 关系顺序相匹配的列表,可以通过对 metricsIdx 递减来反向消费。

getOrderedNodeIds 还编码了仅凭任务统计 Scala 无法可靠推断的 Velox 特有的计划形状调整:

  • 对于 Project 节点,它先访问 source,然后再访问 Project 节点。
  • 如果 Project 的 source 是 Filter 节点,Velox 已将 Filter 之于 Project 的关系映射为一个 FilterProject 算子。Filter 节点没有独立的统计信息,因此原生代码会把 Filter 的 id 记录到 omittedNodeIds 中。
  • 对于 LocalPartitionNode,Velox 可能会插入本地 exchange/分区节点以及可选的投影子节点。原生代码在存在投影子节点时遍历这些投影子节点,否则直接遍历 source。当节点有两个 source 时,它会将 LocalPartitionNode 的 id 记录为具体的 Spark native union transformer。

原生的 treefy 结果成为框架其余部分的排序契约:

Velox plan tree
  -> getOrderedNodeIds
  -> orderedNodeIds + omittedNodeIds

没有 orderedNodeIds,Scala 就不得不依赖 Velox 的 planStats 映射的迭代顺序,或者在事后重新构建 Velox 特有的计划重写。这两种方案都会使指标分配变得脆弱,尤其是对于那些被融合、被省略或被展开为多个 Velox 节点的算子。

步骤 2:JSON 转换为 OperatorMetrics

Velox 迭代器执行完毕后,C++ 会读取 Velox 任务统计信息,并将包含以下字段的 JSON 负载序列化输出:

  • orderedNodeIds:按原生 treefy 顺序排列的 Velox 计划节点 id。
  • omittedNodeIds:预期存在但没有 Velox 统计信息的节点。
  • nodeStats:每个节点的 Velox 算子统计信息。
  • loadLazyVectorTime:Gluten 惰性向量加载耗时。

Scala 在 MetricsUtil.parseNativeOperatorMetrics 中解析该负载:

  1. 遍历 orderedNodeIds。
  2. 在 nodeStats 中查找每个节点 id。
  3. 将该节点的每个 Velox 算子统计信息转换为 OperatorMetrics。
  4. 如果该节点 id 存在于 omittedNodeIds 中且没有统计信息,则插入一个空的 OperatorMetrics 占位符。
  5. 将 loadLazyVectorTime 附加到最新的已扁平化原生指标集合。
  6. 校验解析出的数量是否与 Metrics.numMetrics 一致。

这样便生成了 Spark 更新器树将要消费的扁平化 JList[OperatorMetrics]。

orderedNodeIds: [n0, n1, n2]
nodeStats:
  n0 -> [stat0]
  n1 -> [stat1, stat2]
  n2 -> [stat3]

flattened nativeMetrics:
  [OperatorMetrics(stat0),
   OperatorMetrics(stat1),
   OperatorMetrics(stat2),
   OperatorMetrics(stat3)]

如果某个节点被省略,Scala 仍会插入一个零值占位符,以保证原生指标的索引与更新器遍历顺序持续对齐。

步骤 3:Spark 计划转换为 MetricsUpdaterTree

MetricsUtil.treeifyMetricsUpdaters(plan) 会把 Spark 物理计划转换为一棵指标更新器树。这就是 Spark 侧的 treefy 过程。生成的树描述了哪些 Spark 算子应当接收原生指标,以及应当按照怎样的子节点顺序遍历它们。

关键场景如下:

  • HashJoinLikeExecTransformer:创建一个 join 更新器节点,其子节点按 (buildPlan, streamedPlan) 顺序排列。
  • SortMergeJoinExecTransformer:创建一个 join 更新器节点,其子节点按 (bufferedPlan, streamedPlan) 顺序排列。
  • TransformSupport 搭配 MetricsUpdater.None:跳过当前节点,只对其子节点执行 treefy。当某个 Spark 节点仅为规划出所需的计划形态而存在、本身不应接收原生指标时,就会使用这种方式。
  • 其他 TransformSupport:创建一个更新器节点,并对 children.reverse 执行 treefy。
  • 非 transform 类型的 Spark 节点:会变为 MetricsUpdater.Terminate,从而终止该分支上的原生指标传递。

子节点顺序的反转是有意为之。原生指标之后会从扁平化列表的末尾开始消费,因此更新器树必须与 Substrait 规划及原生 orderedNodeIds 所产生的顺序保持镜像对应。

从概念上讲:

SparkPlan
  -> treeifyMetricsUpdaters
  -> MetricsUpdaterTree(updater, children)

该树不包含指标数值,只包含将原生指标回放到 Spark 算子上所需的更新器拓扑。

第 4 步:消费并映射指标套件

在 Scala 映射代码中,一个 OperatorMetrics 对象代表一个原生指标套件。一个套件通常对应一个 Velox 算子统计项。某些 Spark 算子消费一个套件,另一些则消费多个套件,因为该 Spark 算子会展开为多个 Substrait/Velox 算子。

最初分配给一个 Spark 算子的套件数量为:

relMap(operatorIdx).size()

updateTransformerMetricsInternal 执行算子级映射。对于每个更新节点,它会:

  1. 读取当前 operatorIdx 的 Substrait rel id。
  2. 为每个 rel id 消费一个原生 OperatorMetrics 套件。
  3. 应用针对该算子的特定处理。
  4. 以递减的索引递归进入子更新节点。

对于普通的单元算子,消费到的套件会被合并并传递给:

u.updateNativeMetrics(mergedOperatorMetrics)

合并行为是围绕 Velox 管线结构设计的:

  • 输入侧计数器取自最后一个被消费的 suite。
  • 输出侧与写入计数器取自第一个被消费的 suite。
  • CPU、墙钟时间、溢写、内存分配以及大多数自定义计数器进行累加。
  • loadLazyVectorTime 附加到最终展平的 suite,并跨被消费的 suite 累加。
  • 峰值内存取最大值。

这使得 Spark 算子即使由多个原生 rel 实现,也能得到一行连贯的指标。

对齐方式可以概括为:

native orderedNodeIds treefy
  -> flattened OperatorMetrics
  -> Spark MetricsUpdaterTree traversal
  -> Spark SQLMetric updates

算子专属映射

有些算子并不遵循简单的"消耗关系计数、合并、更新"规则。

Join

Join 更新器会消耗 relMap 分配的指标套件,然后额外再消耗一个套件,用于 build/probe 两侧的指标。该更新器还会接收来自规划阶段的 join 参数,从而能够把 Velox 的 join 侧数值映射到正确的 Spark SQL 指标上。

HashJoinLikeExecTransformer 和 SortMergeJoinExecTransformer 在 Spark 侧进行 treefy 时也会使用自定义的树子节点排序,因为 build/buffered 侧与 streamed 侧必须与原生遍历顺序保持一致。

Union

UnionMetricsUpdater 会额外消耗一个套件,并根据合并后的原生数值更新 union 专属指标。

Hash Aggregate

HashAggregateMetricsUpdater 使用规划阶段记录的聚合参数。原生套件仍然来自 relMap,但更新器需要这些参数才能确定聚合指标到 Spark 指标的映射方式。

Limit Over Sort

Velox 可能把 Limit over Sort 实现为 TopN 风格的原生算子。在这种情况下,原生指标套件归属于 sort 更新器。limit 更新器既不更新指标,也不消耗套件,因此下游索引保持对齐。

端到端示例

对于一个简单的转换后计划:

Project
  Filter
    Scan

native treefy 生成:

orderedNodeIds:
  [scan node, filter node, project node]

Scala 将该 JSON 解析为:

nativeMetrics:
  0 -> scan suite
  1 -> filter suite
  2 -> project suite

Spark 侧 treefy 构建:

ProjectUpdater
  FilterUpdater
    ScanUpdater

假设规划阶段记录了:

operatorToRelsMap:
  0 -> [scan rel]
  1 -> [filter rel]
  2 -> [project rel]

更新器从末尾开始:

operatorIdx = 2, metricsIdx = 2

它先更新 project,然后是 filter,最后是 scan。每一步都会消费当前算子对应的套件,并使两个索引各自递减。更复杂的算子沿用相同的遍历方式,但可能会消费多个套件。

添加或调试映射

在添加指标或调试错误值时,请按照与运行时相同的路径排查:

  1. 检查原生 orderedNodeIds 与 omittedNodeIds,确认是否展平了错误的 Velox 节点,或者遗漏了某个融合节点。
  2. 检查 parseNativeOperatorMetrics 是否按预期的数量和顺序产出 OperatorMetrics 套件。
  3. 检查 VeloxMetricsApi 中的 Spark 指标键。
  4. 查看目标 MetricsUpdater,确认 OperatorMetrics 套件是如何写入 Spark SQL 指标的。
  5. 检查该算子消费的套件数量是否正常,或是否在 updateTransformerMetricsInternal 中有特殊处理逻辑。
  6. 如果被更新的是连接或多子节点算子的错误一侧,请检查 Spark 侧 treefy 的子节点排序。

最有用的不变式是:

parsed native metric count == Metrics.numMetrics

如果上述情况成立,但数值被归到了错误的 Spark 算子上,请检查原生 orderedNodeIds、MetricsUpdaterTree 的结构、operatorToRelsMap,以及任何算子特有的额外 suite 消耗。

评论

登录后参与评论

正在加载评论…