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

如何为SparkSession注册StreamingListener?StreamingQueryListener指标获取

结构化流中通过StreamingQueryListener实现微批指标采集

在结构化流应用中,只能通过sparkSession.streams.addListener()注册StreamingQueryListener,无法直接关联原有的StreamingListener(基于StreamingContext的监听逻辑不适用于SparkSession驱动的结构化流)。以下是将原微批指标采集逻辑(调度延迟、处理时间、记录数)迁移到StreamingQueryListener的实现方案:

实现代码

import io.prometheus.client.{CollectorRegistry, Gauge}
import io.prometheus.client.exporter.PushGateway
import org.apache.spark.sql.streaming.{StreamingQueryListener, QueryProgressEvent, QueryTerminatedEvent, QueryStartedEvent}

class MicroBatchStatsQueryListener(val pushGateway: PushGateway, jobName: String) extends StreamingQueryListener {

  private val processingTimeGauge = Gauge.build()
    .name("processingTimeGauge")
    .help("Time it took to process this microbatch (seconds)")
    .register(CollectorRegistry.defaultRegistry)

  private val schedulingDelayGauge = Gauge.build()
    .name("schedulingDelayGauge")
    .help("Scheduling delay of this microbatch (seconds)")
    .register(CollectorRegistry.defaultRegistry)

  private val numRecordsGauge = Gauge.build()
    .name("numRecordsGauge")
    .help("Number of records received in this microbatch")
    .register(CollectorRegistry.defaultRegistry)

  override def onQueryStarted(event: QueryStartedEvent): Unit = {
    // 可选:处理查询启动事件,无需逻辑可留空
  }

  override def onQueryProgress(event: QueryProgressEvent): Unit = {
    val progress = event.progress
    
    // 处理时间:毫秒转秒,无值时默认0
    val processingTime = progress.processingDelay.getOrElse(0L) / 1000.0
    // 调度延迟:毫秒转秒,无值时默认0
    val schedulingDelay = progress.schedulingDelay.getOrElse(0L) / 1000.0
    // 当前微批的输入记录数
    val numRecords = progress.numInputRows

    // 更新Prometheus指标
    processingTimeGauge.set(processingTime)
    schedulingDelayGauge.set(schedulingDelay)
    numRecordsGauge.set(numRecords)

    // 推送指标到PushGateway
    pushGateway.push(CollectorRegistry.defaultRegistry, jobName)
  }

  override def onQueryTerminated(event: QueryTerminatedEvent): Unit = {
    // 可选:处理查询终止事件,无需逻辑可留空
  }
}

关键指标获取说明

  • 调度延迟:从QueryProgress.schedulingDelay获取,该字段表示微批从进入调度队列到开始执行的间隔(单位毫秒),转成秒后使用。
  • 处理时间:从QueryProgress.processingDelay获取,该字段表示微批实际处理数据消耗的时间(单位毫秒),转成秒后使用。
  • 记录数:直接从QueryProgress.numInputRows获取,代表当前微批接收到的输入记录总数。

注册Listener

创建实例后通过SparkSession完成注册:

val spark = SparkSession.builder().appName("StructuredStreamingApp").getOrCreate()
val pushGateway = new PushGateway("pushgateway-host:port")
val listener = new MicroBatchStatsQueryListener(pushGateway, "your-job-name")
spark.streams.addListener(listener)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 15:53:09