intermediate
Monitoring Streaming Jobs
10 min read
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).