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

Spark中StreamQueryListener的onQueryProgress()未执行问题求助

问题排查与解决方案

核心问题1:Listener输出无法在Notebook单元格显示

你的代码中使用println输出Listener事件内容,但在Databricks Notebook环境下,StreamingQueryListener的回调方法运行在Spark内部异步线程中,println的内容不会被捕获到Notebook单元格输出里,而是直接输出到Spark Driver的日志中。

解决方法:改用Spark日志系统记录

将println替换为Spark日志工具,输出会写入Driver日志,可在集群日志页面查看:

%scala
import org.apache.spark.sql.streaming._
import org.apache.log4j.Logger

// 初始化日志实例
val logger = Logger.getLogger(getClass.getName)

val streamingCountsListener = new StreamingQueryListener() {
  override def onQueryStarted(queryStarted: StreamingQueryListener.QueryStartedEvent): Unit = {
      logger.info("query started")
  }
  override def onQueryTerminated(queryTerminated: StreamingQueryListener.QueryTerminatedEvent): Unit = {
        logger.info("query made stopped")
      
  }
  override def onQueryProgress(queryProgress: StreamingQueryListener.QueryProgressEvent): Unit = {
    queryProgress.progress.sources.foreach(src => {
      logger.info(s"批次输入行数: ${src.numInputRows}")
    })
  }
}

// 注册Listener
spark.streams.addListener(streamingCountsListener)

// 读取Delta流
var x = spark.readStream.format("delta")
  .option("ignoreChanges", "true")
  .table(s"qastg.student")

// 启动流查询,设置触发间隔确保定期检查
val query = x.writeStream
  .format("console")
  .outputMode("append")
  .trigger(Trigger.ProcessingTime("5 seconds")) // 新增触发间隔,无数据时也会触发检查
  .start()

query.awaitTermination()

查看Driver日志步骤:

  1. 点击Databricks界面右上角集群图标,进入集群详情页
  2. 切换到「Logs」标签页,查看Driver节点的stdout或stderr日志即可看到输出内容

核心问题2:流查询未触发批次处理

如果Delta表(或Kafka主题)没有新数据写入,流查询不会触发批次处理,也就不会调用onQueryProgress方法。

解决方法:

  • 手动向qastg.studentDelta表写入测试数据,触发流批次处理
  • 添加trigger(Trigger.ProcessingTime("5 seconds"))配置,让流每隔5秒检查一次数据,即使无新数据也会触发进度检查(此时numInputRows会显示为0)

额外验证:Listener注册时机

你的代码先注册Listener再创建流查询,这个顺序是正确的。若在流查询启动后才注册Listener,将无法捕获后续事件,这点你已做对,可排除该问题。

内容的提问来源于stack exchange,提问作者Muqtada Hussain

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.25 08:47:24