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
StreamingQueryListeneron the SparkSession to capture query metrics and send them to your monitoring dashboard (like Prometheus or Datadog).