Spark Thrift Server监控:通过Listener获取查询行数及关联作业问题
Spark Thrift Server 监控采集问题解答
问题1:能否从SparkListenerSQLExecutionEnd事件中获取正确的输出行数
可以,你当前拿到的指标不准是实现逻辑错误:你直接访问的queryExecution.executedPlan是物理执行计划的静态模板,没有填充运行时实际累加的指标值。
- 正确实现方式:直接读取
SparkListenerSQLExecutionEnd事件自带的聚合后指标,事件提供的executionMetrics字段已经封装了该SQL执行全生命周期所有聚合完成的运行指标,无需自行从执行计划中爬取。 - 调整后代码示例:
def onExecutionEnd(event: SparkListenerSQLExecutionEnd): Unit = { logger.info("-----onExecutionEnd") logger.info(s"---------executionId: ${event.executionId}") logger.info(s"---------endTime: ${event.time}") // 直接使用事件携带的执行指标 val executionMetrics = event.executionMetrics // 先打印全量指标确认对应关系,不同Spark版本metric名称略有差异 executionMetrics.foreach { case (metricId, metricValue) => logger.info(s"Metric name: ${metricValue.name.getOrElse("")}, value: ${metricValue.value}") } // 读取行数、输出行数取值示例(以Spark 3.x为例) val numReadRows = executionMetrics.values.find(_.name.contains("number of rows read")).map(_.value).getOrElse(0L) val numOutputRows = executionMetrics.values.find(_.name.contains("number of output rows")).map(_.value).getOrElse(0L) logger.info(s"读取行数:$numReadRows,返回行数:$numOutputRows") }
- 补充说明:如果你需要的是返回给Thrift客户端的实际行数,部分版本STS会将最终结果拉取到Driver端再返回,最终输出行数需要匹配根执行节点的输出指标,可先打印全量
executionMetrics的名称和值确认对应关系。
问题2:是否可以关联SQL执行ID与对应的stage或job
可以,关联链路通过事件携带的属性传递:
- 所有SQL触发的Job,在
SparkListenerJobStart事件的properties字段中都会携带spark.sql.execution.id参数,参数值就是对应的SQL执行ID,可通过该参数建立executionId和jobId的关联。 - 拿到
executionId和jobId的关联关系后,再通过SparkListenerStageCompleted事件的stageInfo.jobIds字段,将stage和对应job关联,即可完成executionId -> jobId -> stageId的全链路关联。
补充问题:Docker环境生效生产环境不生效的排查方向
该问题基本都是环境差异导致,优先排查以下几点:
- 版本差异:确认Docker环境和生产环境的Spark版本是否一致,低版本Spark(2.3及以下)没有将executionId直接封装到
SparkListenerJobStart的公共字段,必须从properties中读取,如果你是直接调用jobStart.executionId字段在低版本会返回空。 - 配置差异:检查生产环境STS是否开启了
spark.sql.execution.id.propagate配置(默认开启,部分安全加固版本可能会关闭),该参数控制SQL执行ID是否会传递到Job的属性中。 - Listener注册时机:确认自定义Listener是在SparkContext初始化完成前注册的,如果是运行时动态注册的Listener,可能会漏掉部分启动阶段的事件属性传递。
内容的提问来源于stack exchange,提问作者SCouto
相关产品推荐
相关产品推荐

