intermediate
Trigger Types
8 min readLast updated: 2026-07-09
Overview
Learn how to control when Spark processes streaming micro-batches using default, ProcessingTime, Once, and Continuous triggers.
What You Will Learn
In this lesson, you will learn:
- Processing Time Trigger: Setting interval-based micro-batches.
- Once / AvailableNow Trigger: Running streams as batch-style workloads.
- Continuous Processing: Sub-millisecond latency execution.
Detailed Concept Explanation
A Trigger defines the timing of streaming query executions.
Trigger Options
- Default (Unspecified): Spark runs micro-batches back-to-back. As soon as one micro-batch finishes, the next one starts.
processingTime: Spark runs a micro-batch at specified intervals (e.g.10 seconds), waiting if no new data has arrived.once=TrueoravailableNow=True: Spark processes all available data in the source and then stops the query automatically. This is useful for running daily streaming queries as cheap batch jobs.continuous: A low-latency engine mode that processes events continuously rather than in micro-batches, achieving sub-millisecond latencies. (Only supports simple maps/filters).
Code Examples
Python (PySpark) Implementation
python
from pyspark.sql import SparkSession
spark = SparkSession.builder.appName("TriggerTypes").getOrCreate()
rate_stream = spark.readStream.format("rate").load()
# Run micro-batches every 10 seconds
query = rate_stream.writeStream \
.format("console") \
.trigger(processingTime="10 seconds") \
.start()
query.awaitTermination(5)
query.stop()
Expected Output
text
+-------------------+-----+
| timestamp|value|
+-------------------+-----+
|2026-07-09 19:30:00| 1|
+-------------------+-----+
Execution Plan Diagram (Python & Scala)
Execution Plan Diagram
SparkSession.builder
readStream
writeStream.trigger(processingTime=10s)
start()
Scala Implementation
scala
import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.streaming.Trigger
val spark = SparkSession.builder().appName("TriggerTypesScala").getOrCreate()
val rateStream = spark.readStream.format("rate").load()
val query = rateStream.writeStream
.format("console")
.trigger(Trigger.ProcessingTime("10 seconds"))
.start()
query.awaitTermination(5000)
query.stop()
Expected Output
text
+-------------------+-----+
| timestamp|value|
+-------------------+-----+
|2026-07-09 19:30:00| 1|
+-------------------+-----+
Common Mistakes
- Confusing Continuous Processing with Micro-Batching: Expecting continuous processing (
Trigger.Continuous) to support complex aggregations or window queries. It only supports simple row-by-row mapping transformations.