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
相关产品推荐
相关产品推荐

