You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark SQL多SQL场景下如何获取作业状态及当前处理SQL?

跟踪Spark SQL执行状态与当前运行SQL的解决方案

我之前也碰到过类似的问题——用Spark Listener的Job事件确实很难直接关联到对应的SQL,毕竟一个SQL可能拆成多个Job,而且Job事件里没有直接带SQL文本。下面给你几个实用的方案,分版本适配,你可以根据自己的Spark版本选:

方案一:Spark 3.x+ 专属——用SparkListener的Query级事件

Spark 3.0之后新增了专门针对SQL查询的事件(QueryStartEvent、QueryEndEvent等),这些事件直接绑定到每个spark.sql()调用,能直接拿到原始SQL文本和执行状态,简直是为这个场景量身定做的!

步骤1:自定义SparkListener

import org.apache.spark.scheduler._
import org.apache.spark.sql.execution.QueryExecutionStatus

class SQLTrackingListener extends SparkListener {
  // 监听查询开始事件
  override def onQueryStart(event: QueryStartEvent): Unit = {
    println(s"[查询启动] ID: ${event.queryId}, SQL: ${event.sqlText}")
  }

  // 监听查询结束事件
  override def onQueryEnd(event: QueryEndEvent): Unit = {
    val status = event.status match {
      case QueryExecutionStatus.SUCCESS => "✅ 已完成"
      case QueryExecutionStatus.FAILED => "❌ 失败"
      case QueryExecutionStatus.CANCELED => "⏹️ 已取消"
    }
    println(s"[查询结束] ID: ${event.queryId}, SQL: ${event.sqlText}, 状态: $status")
  }
}

步骤2:注册监听器并执行SQL

import org.apache.spark.sql.SparkSession

object SQLTracker {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("SQLExecutionTracker")
      .master("local[*]") // 生产环境去掉这个
      .getOrCreate()

    // 注册自定义监听器
    spark.sparkContext.addSparkListener(new SQLTrackingListener())

    // 你的SQL列表
    val sqls = List(
      "SELECT * FROM test_table WHERE id < 100",
      "INSERT INTO result_table SELECT name, count(*) FROM test_table GROUP BY name"
    )

    // 注意:spark.sql()是懒加载,必须触发Action才会执行!
    sqls.foreach { sql =>
      spark.sql(sql).collect() // collect()只是示例,你可以换成write()等实际业务需要的Action
    }

    spark.stop()
  }
}

这个方案最省心,不需要额外的属性绑定,直接从事件里拿SQL和状态,还能配合onQueryProgressEvent跟踪查询的实时进度。

方案二:兼容Spark 2.x+——QueryExecutionListener + 自定义属性

如果你的Spark版本还停留在2.x,那就用QueryExecutionListener(专门监听SQL执行生命周期的监听器),再配合Spark配置属性来绑定当前执行的SQL。

步骤1:自定义QueryExecutionListener

import org.apache.spark.sql.{QueryExecution, SparkSession}
import org.apache.spark.sql.util.QueryExecutionListener

class SQLExecutionListener extends QueryExecutionListener {
  override def onSuccess(queryExecution: QueryExecution, durationNs: Long): Unit = {
    // 从配置中获取当前执行的SQL
    val currentSql = queryExecution.sparkSession.conf.get("custom.current.sql", "未知SQL")
    println(s"[SQL执行成功] 耗时: ${durationNs / 1000000}ms, SQL: $currentSql")
  }

  override def onFailure(queryExecution: QueryExecution, exception: Exception): Unit = {
    val currentSql = queryExecution.sparkSession.conf.get("custom.current.sql", "未知SQL")
    println(s"[SQL执行失败] 错误: ${exception.getMessage}, SQL: $currentSql")
  }
}

步骤2:注册监听器并执行SQL

import org.apache.spark.sql.SparkSession

object SQLTracker2x {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder()
      .appName("SQLExecutionTracker2x")
      .master("local[*]")
      .getOrCreate()

    // 注册监听器
    spark.listenerManager.register(new SQLExecutionListener())

    val sqls = List(/* 你的SQL列表 */)

    sqls.foreach { sql =>
      try {
        // 设置当前SQL到Spark临时配置
        spark.conf.set("custom.current.sql", sql)
        println(s"[开始执行] SQL: $sql")
        // 触发Action执行SQL
        spark.sql(sql).collect()
      } catch {
        case e: Exception => 
          // 这里也可以直接捕获异常,标记失败
          println(s"[执行异常] SQL: $sql, 错误: ${e.getMessage}")
      } finally {
        // 清理临时配置
        spark.conf.unset("custom.current.sql")
      }
    }

    spark.stop()
  }
}

这个方案的核心是用Spark的临时配置传递当前SQL,监听器通过配置拿到对应关系,兼容所有Spark 2.x及以上版本。

方案三:最基础的串行执行跟踪

如果你的SQL是串行执行的(就像你代码里的foreach),其实可以不用监听器,直接在循环里用try-catch捕获状态,代码最简单:

sqls.foreach { sql =>
  println(s"SQL 状态: 运行中 | $sql")
  try {
    spark.sql(sql).collect() // 触发执行
    println(s"SQL 状态: 已完成 | $sql")
  } catch {
    case e: Exception =>
      println(s"SQL 状态: 失败 | $sql, 错误信息: ${e.getMessage}")
  }
}

这个方案适合简单场景,不需要额外的监听器代码,但缺点是如果改成并行执行(比如foreachPar),就需要处理线程安全问题。

关键注意点

  • 懒加载陷阱:spark.sql(sql)只是生成逻辑计划,不会实际执行!必须调用collect()、write()、count()这类Action操作,才会触发SQL执行,产生对应的监听器事件。
  • 并行执行的线程安全:如果是并行执行SQL,要改用ThreadLocal存储当前SQL,或者用线程安全的ConcurrentHashMap关联执行标识和SQL。
  • 集群模式输出:监听器的打印日志会输出在Driver节点,如果需要把状态同步到外部系统(比如监控平台、数据库),可以在监听器的方法里添加对应的写入逻辑。

内容的提问来源于stack exchange,提问作者AI Joes

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 03:52:14