intermediate

Monitoring Streaming Jobs

10 min readLast updated: 2026-07-09

Overview

Learn how to monitor running Spark Structured Streaming queries, track throughput metrics, and listen to query progress events.

What You Will Learn

In this lesson, you will learn:
  • Query Status: Checking current stream activity.
  • Query Progress: Inspecting input rates and processing throughput metrics.
  • StreamingQueryListener: Registering custom callback listeners.

Detailed Concept Explanation

To monitor streaming queries, Spark provides status and progress metrics directly on the query object.

Monitoring Methods

  • query.status: Returns what the query is doing at that exact moment (e.g. "message": "Processing new data").
  • query.lastProgress: Returns a detailed JSON report of the last executed micro-batch, containing:
    • inputRowsPerSecond: The rate at which data is arriving from the source.
    • processedRowsPerSecond: The rate at which Spark is processing the data.
    • durationMs: A breakdown of time spent on tasks (trigger, walCommit, stateStore, etc.).

Code Examples

Python (PySpark) Implementation

python
from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("MonitoringStreaming").getOrCreate()

rate_stream = spark.readStream.format("rate").load()
query = rate_stream.writeStream.format("console").start()

# Print current query state
print("Query Status Details:", query.status)

# Print last progress metrics
print("Last Progress Metrics:", query.lastProgress)

query.awaitTermination(5)
query.stop()

Expected Output

text
Query Status Details: {'message': 'Initializing sources', 'isDataAvailable': False, 'isTriggerActive': False}
Last Progress Metrics: None

Execution Plan Diagram (Python & Scala)

Execution Plan Diagram
SparkSession.builder
readStream
writeStream.format(console)
start()
print(query.status)
print(query.lastProgress)

Scala Implementation

scala
import org.apache.spark.sql.SparkSession

val spark = SparkSession.builder().appName("MonitoringScala").getOrCreate()

val rateStream = spark.readStream.format("rate").load()
val query = rateStream.writeStream.format("console").start()

println(s"Query Status Details: ${query.status}")
println(s"Last Progress Metrics: ${query.lastProgress}")

query.awaitTermination(5000)
query.stop()

Best Practices

  • Use StreamingQueryListener: In production, register a custom StreamingQueryListener on the SparkSession to capture query metrics and send them to your monitoring dashboard (like Prometheus or Datadog).

Related Topics