Every Structured Streaming query has a trigger, whether you set one or not. The trigger decides when the engine looks for new input and runs the next unit of work, and that single setting determines latency, cost, output file counts, how a query behaves after a backlog, and whether it runs forever or stops when it has caught up. Teams often copy a trigger from an example and never revisit it, then wonder why their table has millions of small files or why a nightly job silently processed only part of its input.

This article explains how the micro-batch loop works, what each trigger type actually does according to the Apache Spark documentation, how triggers interact with source rate limits and watermarks, how to choose an interval with arithmetic rather than guesswork, how to monitor whether a query keeps up, and the failure modes of each choice. It covers the default, fixed-interval, one-time, available-now and continuous triggers, plus the Real-time trigger added in Spark 4.1.

The micro-batch loop

Fixed 10-second interval: what happens when a batch overruns0s10s20s30s40s50s60sbatch 0batch 1 (overruns)b2batch 3idle until the 10s boundarystarts at once: boundary missedno new data: no batch1. Planread latest offsets, apply rate limits2. Write offset logthen run the batch job3. Commit logstate and sink doneThe trigger only decides WHEN step 1 starts. Steps 1 to 3 are the same for every micro-batch trigger.
Top: a fixed 10-second trigger waits after a fast batch, starts immediately after an overrun, and skips intervals with no new data. Bottom: the steps every micro-batch performs.

In micro-batch execution the driver repeats a loop. It asks each source for its latest available offsets, applies any per-trigger rate limits, and writes the chosen offset range to the offset log in the checkpoint directory before running anything. It then runs an ordinary Spark job over exactly that range, updates state stores and writes to the sink. Finally it records the batch in the commit log. On restart, an offset entry without a matching commit means the batch is re-run with the same range, which is how Structured Streaming gets exactly-once processing with replayable sources and idempotent sinks.

The trigger controls only when the next plan step begins. It does not change what a batch contains beyond the offsets available at that moment, and it does not make a slow batch faster. Keeping that separation in mind prevents the most common misunderstanding: a shorter interval does not reduce latency if each batch already takes longer than the interval.

The trigger types

TriggerPythonScala / JavaBehaviour
Unspecified (default)no callno callNext micro-batch starts as soon as the previous one completes.
Fixed intervaltrigger(processingTime='10 seconds')Trigger.ProcessingTime("10 seconds")Starts batches on interval boundaries; an overrun starts the next batch immediately; no new data, no batch.
One-time (deprecated since 3.4)trigger(once=True)Trigger.Once()One batch over everything available, then stop.
Available-now (since 3.3)trigger(availableNow=True)Trigger.AvailableNow()Processes everything available at start, in possibly many batches that respect rate limits, then stops.
Continuous (experimental, since 2.3)trigger(continuous='1 second')Trigger.Continuous("1 second")Long-running tasks, map-like queries only, at-least-once; the interval is a checkpoint interval.
Real-time (since 4.1)see belowTrigger.RealTime("5 minutes")Long-running tasks with exactly-once processing; the duration is a checkpoint interval, not a latency target.

The rest of this article takes them in the order most teams need them: always-on micro-batch queries, scheduled incremental jobs, and low-latency modes.

Default and fixed-interval triggers

With no trigger, the engine runs batch after batch with no gap, which minimises latency but maximises the number of batches, and therefore the number of commits, small output files and calls to the source. A fixed interval caps that rate. The documented rules are precise: if a batch finishes early, the engine waits for the next boundary; if a batch overruns, the next one starts as soon as the previous completes rather than waiting for a later boundary; and if no new data is available, no batch is started. Missed boundaries are not queued up and replayed, so an overrun does not create a burst of catch-up batches.

query = (events.writeStream
    .format("delta")
    .option("checkpointLocation", "s3://lake/_chk/events_silver")
    .trigger(processingTime="1 minute")
    .toTable("silver.events"))

One exception to the no-data rule matters for stateful queries. When a watermark has moved far enough that windows can be closed or state evicted, Spark can run a batch with no new input so that results are emitted without waiting for more data. This is governed by spark.sql.streaming.noDataMicroBatches.enabled, which defaults to true. If an aggregation only emits when new data arrives, check that this setting has not been turned off.

Rate limits decide batch size

A trigger says when to plan; source options say how much to take. Without limits, the first batch after a long outage tries to read the entire backlog at once, which can take hours, overflow executor memory and delay every downstream consumer. The relevant options are:

  • Kafka: maxOffsetsPerTrigger caps the offsets per batch, split proportionally across partitions. minOffsetsPerTrigger delays a batch until enough data accumulates, bounded by maxTriggerDelay (default 15 minutes), which is useful for reducing tiny batches on quiet topics.
  • Files: maxFilesPerTrigger or maxBytesPerTrigger, but not both. A batch always reads at least one file so a single oversized file cannot stall the stream.

Rate limits turn a backlog into a predictable series of bounded batches. The cost is recovery time: if the cap is 600,000 offsets per batch and batches run every minute, a backlog of 36 million offsets takes about an hour to drain, so size the cap from the recovery time you can accept.

Scheduled jobs: one-time versus available-now

Many streaming pipelines do not need to run continuously. Ingesting into a lakehouse every hour from a cluster that starts, catches up and shuts down is far cheaper than an always-on cluster. That is what the one-time and available-now triggers are for: the query uses the same checkpoint as a continuous run, processes what has arrived since last time, and then stops.

One-time runs a single batch over everything available and ignores rate limits, so a large backlog becomes one enormous batch. Available-now processes the same data in as many batches as the source's rate limits imply, guarantees that all data available at start time is processed (including uncommitted batches from a previous failed run), and advances the watermark per batch, running a final no-data batch so stateful output is emitted before the query stops. Spark deprecated one-time in 3.4 in favour of available-now.

# Scheduled hourly by Airflow, cron or a job scheduler.
q = (spark.readStream.format("kafka")
        .option("kafka.bootstrap.servers", BROKERS)
        .option("subscribe", "orders")
        .option("maxOffsetsPerTrigger", 2_000_000)
        .load()
        .selectExpr("CAST(value AS STRING) AS json", "timestamp")
        .writeStream
        .format("delta")
        .option("checkpointLocation", "s3://lake/_chk/orders_bronze")
        .trigger(availableNow=True)
        .toTable("bronze.orders"))
q.awaitTermination()          # returns once the backlog at start time is processed
if q.exception():
    raise q.exception()       # fail the scheduler task so it retries

Available-now is active only when every source in the query supports it. Since Spark 4.0, if any source does not, Spark falls back to running a single batch, as one-time did, so a query that seems to ignore its rate limit may contain such a source. Kafka and file sources support it; check the documentation of any other source.

Low latency: continuous and Real-time

Micro-batch latency has a floor of roughly 100 milliseconds according to the Spark documentation, because every batch pays for planning, task scheduling and a commit. Two modes avoid that by running long-lived tasks instead of a new job per batch.

Continuous processing, introduced in Spark 2.3, is still experimental. It supports only map-like operations (projections and filters, no aggregations), Kafka and rate sources, and Kafka, memory and console sinks, and it gives at-least-once guarantees. It needs at least one core per source partition, and a failed task stops the query, which must be restarted manually from the checkpoint.

Real-time mode, added in Spark 4.1, is the documented successor for new low-latency work. It is enabled with a Real-time trigger, Trigger.RealTime(...) in Scala and Java, and the current documentation also shows a realTime= keyword for Python; check that your PySpark release accepts it. Tasks stay alive for the whole batch and process records as they arrive, so the duration (default 5 minutes, minimum 5 seconds by default) is mainly a checkpoint interval, not a latency target. The query must use update output mode and a checkpoint location. Spark 4.1 supports stateless queries; the development documentation describes stateful operators such as aggregation and deduplication arriving in later releases, so check the guide for your version before relying on them. Processing is exactly-once, but delivery depends on the sink, and the built-in Kafka sink is at-least-once.

// Scala, Spark 4.1+: stateless enrichment from Kafka to Kafka.
import org.apache.spark.sql.streaming.Trigger

spark.readStream.format("kafka")
  .option("kafka.bootstrap.servers", brokers)
  .option("subscribe", "payments")
  .load()
  .selectExpr("CAST(key AS STRING) AS key", "CAST(value AS STRING) AS value")
  .writeStream.format("kafka")
  .option("kafka.bootstrap.servers", brokers)
  .option("topic", "payments_scored")
  .option("checkpointLocation", "/chk/payments_scored")
  .outputMode("update")
  .trigger(Trigger.RealTime("1 minute"))
  .start()

Worked example: choosing an interval

Suppose a query reads a 12-partition Kafka topic at 30,000 events per second and writes a Delta table, and the business wants data queryable within two minutes. Measure first: run with the default trigger and read batch durations from the progress events. Say a batch of 10 seconds of data takes 4 seconds and a batch of 60 seconds of data takes 9, since fixed overhead dominates small batches.

With a 10-second interval, worst-case freshness is about the interval plus the batch time, roughly 14 seconds, but the query commits 8,640 batches a day. If each batch writes 12 files (one per task), that is 103,680 files a day before compaction. With a 60-second interval, freshness is about 69 seconds, comfortably inside two minutes, and the table receives 17,280 files a day, six times fewer, with lower checkpoint and listing cost. The 60-second trigger wins. Then set maxOffsetsPerTrigger to about three minutes of traffic, 5.4 million offsets, so a backlog drains in bounded batches.

The stability rule is simple: average batch duration must stay well below the interval. If it creeps toward the interval, the query is about to fall behind, and the right fix is more executors or less work per record, not a shorter trigger.

Monitoring whether a query keeps up

Every query reports a progress object after each batch. Watch the trigger duration, the input and processing rates, and the source lag:

from pyspark.sql.streaming import StreamingQueryListener

class TriggerHealth(StreamingQueryListener):
    def __init__(self, interval_ms):
        self.interval_ms = interval_ms
    def onQueryStarted(self, event): pass
    def onQueryIdle(self, event): pass
    def onQueryTerminated(self, event): pass
    def onQueryProgress(self, event):
        p = event.progress
        took = p.durationMs.get("triggerExecution", 0)
        if took > 0.8 * self.interval_ms:
            log.warning("batch %s took %s ms of a %s ms trigger", p.batchId, took, self.interval_ms)
        if p.processedRowsPerSecond and p.inputRowsPerSecond > p.processedRowsPerSecond:
            log.warning("falling behind: in %.0f/s, processed %.0f/s",
                        p.inputRowsPerSecond, p.processedRowsPerSecond)

spark.streams.addListener(TriggerHealth(interval_ms=60_000))

Alert on sustained conditions, not single batches: one slow batch after a compaction is normal, ten in a row is not. For Kafka, also export consumer lag in offsets, because a query can process at its rate limit indefinitely while lag grows.

Failure modes

  • Small-file explosion. A default or seconds-level trigger writing to a lake table creates files faster than compaction removes them. Lengthen the interval or schedule compaction.
  • The giant first batch. A query restarted after a long outage with no rate limit tries to process days of data in one batch and runs out of memory. Always set a per-trigger limit on production sources.
  • One-time against a backlog. A scheduled job still using one-time processes a week of data in one batch. Migrate to available-now.
  • Silent fallback. Available-now degrades to a single batch when any source lacks support. Check every source in multi-source queries.
  • Results that wait for data. Windowed aggregations stop emitting on a quiet stream when no-data batches are disabled. Keep the default.
  • Assuming continuous is exactly-once. It is at-least-once; make the sink idempotent.
  • Treating the Real-time duration as latency. A short duration only adds commit pauses and raises tail latency. Pick it for recovery granularity instead.

Trade-offs and related reading

Shorter intervals buy freshness with overhead: more commits, more files and more load on sources and metadata services. Longer intervals buy efficiency with staleness. Available-now trades always-on freshness for large cost savings on hourly or daily pipelines. Continuous and Real-time trade breadth of supported operations for millisecond latency. Micro-batch checkpoints can be restarted with a different micro-batch trigger, so you can move between fixed-interval and available-now as requirements change without rebuilding state.

Related pages: Structured Streaming fundamentals, the Kafka source for offset handling, watermarks for when stateful results are emitted, foreachBatch for per-batch sink logic, and streaming sinks for delivery guarantees.

What to do next

  1. List every streaming query and record its trigger; flag any that rely on the default without a reason.
  2. Measure batch duration from progress events and set intervals from freshness targets and file counts.
  3. Add maxOffsetsPerTrigger, maxFilesPerTrigger or maxBytesPerTrigger to every production source.
  4. Replace one-time triggers with available-now and check that every source supports it.
  5. Install a progress listener that alerts when batch time nears the interval or input outpaces processing.
  6. Use continuous only for map-like Kafka pipelines, and evaluate Real-time mode on Spark 4.1+ for new low-latency work.
Key takeaway: A trigger decides only when the next unit of work starts. Use fixed intervals sized from freshness targets and file counts for always-on queries, available-now with rate limits for scheduled jobs, and Real-time mode on Spark 4.1+ rather than experimental continuous processing for millisecond latency. Watch batch duration against the interval and fix slowness with resources, not shorter triggers.