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

Spark流查询使用awaitTermination后如何获取进度?

问题:Spark Structured Streaming中使用awaitTermination时lastProgress返回null的解决方法

我是Spark新手,正在学习Spark应用监控相关内容。我想了解Spark应用在指定触发时间内处理的记录数以及查询进度。我知道lastProgress可以提供这些指标,但当我结合awaitTermination使用时,它总是返回null。

我的代码如下:

val q4s = spark.readStream 
.format("kafka") 
.option("kafka.bootstrap.servers", brokers) 
.option("subscribe", topic) 
.option("startingOffsets", "earliest") 
.load() 
.writeStream 
.outputMode("append") 
.option("checkpointLocation", checkpoint_loc) 
.trigger(Trigger.ProcessingTime("10 seconds")) 
.format("console") 
.start() 
println("Query Id: "+ q4s.id.toString()) 
println("QUERY PROGRESS.........") 
println(q4s.lastProgress); 
q4s.awaitTermination();

输出结果:

Query Id: efd6bc15-f10c-4938-a1aa-c81fdb2b33e3 
QUERY PROGRESS......... 
null

请问如何在使用awaitTermination时获取查询进度,或者无需使用该方法如何让查询持续运行?提前感谢。


回答

哈哈,这个问题我刚学Structured Streaming的时候也踩过坑!你现在碰到的lastProgress返回null的原因很简单:你调用lastProgress的时机太早了。start()方法是异步启动流查询的,你刚调用完start()就立刻打印lastProgress,这时候流查询还没完成第一次触发处理,自然没有任何进度数据可以返回。

给你两个实用的解决方案,按需选择:

方案一:用后台线程定期拉取进度,同时保留awaitTermination

这种方式既能让主线程通过awaitTermination()维持流查询的运行(避免程序直接退出),又能在后台持续获取进度信息。代码示例如下:

val q4s = spark.readStream 
.format("kafka") 
.option("kafka.bootstrap.servers", brokers) 
.option("subscribe", topic) 
.option("startingOffsets", "earliest") 
.load() 
.writeStream 
.outputMode("append") 
.option("checkpointLocation", checkpoint_loc) 
.trigger(Trigger.ProcessingTime("10 seconds")) 
.format("console") 
.start() 

// 启动一个后台线程专门打印进度
new Thread(() => {
  // 先等查询激活
  while (!q4s.isActive) {
    Thread.sleep(1000)
  }
  // 只要查询还在运行,就定期拉取进度
  while (q4s.isActive) {
    val progress = q4s.lastProgress
    if (progress != null) {
      println("\n=== 当前查询进度 ===")
      println(s"本次触发处理的记录数: ${progress.numInputRows}")
      println(s"触发时间: ${progress.timestamp}")
      // 还可以打印更多细节,比如处理延迟、批次耗时等
    }
    // 和你的触发间隔保持一致,避免频繁查询
    Thread.sleep(10000)
  }
}).start()

println("Query Id: "+ q4s.id.toString())
// 主线程阻塞,保持查询运行
q4s.awaitTermination()

方案二:不用awaitTermination,用循环维持查询运行

如果不想依赖awaitTermination(),可以用一个循环持续检查查询的活跃状态,同时在循环里获取进度:

val q4s = spark.readStream 
.format("kafka") 
.option("kafka.bootstrap.servers", brokers) 
.option("subscribe", topic) 
.option("startingOffsets", "earliest") 
.load() 
.writeStream 
.outputMode("append") 
.option("checkpointLocation", checkpoint_loc) 
.trigger(Trigger.ProcessingTime("10 seconds")) 
.format("console") 
.start() 

println("Query Id: "+ q4s.id.toString())

// 循环检查查询状态,同时拉取进度
while (q4s.isActive) {
  val progress = q4s.lastProgress
  if (progress != null) {
    println("\n=== 当前查询进度 ===")
    println(s"本次触发处理的记录数: ${progress.numInputRows}")
    println(s"触发时间: ${progress.timestamp}")
  }
  // 每隔10秒检查一次,和触发间隔匹配
  Thread.sleep(10000)
}

额外小提示:

  • 如果需要收集所有批次的历史进度,而不只是最近一次,推荐使用QueryListener来监听进度事件,这是Spark官方推荐的监控方式,能更全面地捕获查询的生命周期事件。
  • 别忘了确认你的Kafka主题里有数据哦!如果主题是空的,即使触发了处理,numInputRows会是0,但lastProgress不会再是null了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.11 08:50:33