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

Related Topics