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

  1. Default (Unspecified): Spark runs micro-batches back-to-back. As soon as one micro-batch finishes, the next one starts.
  2. processingTime: Spark runs a micro-batch at specified intervals (e.g. 10 seconds), waiting if no new data has arrived.
  3. once=True or availableNow=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.
  4. 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.

Related Topics