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

