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日志步骤:
- 点击Databricks界面右上角集群图标,进入集群详情页
- 切换到「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
相关产品推荐
相关产品推荐

