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

