When would you schedule a Spark Structured Streaming query with Trigger.AvailableNow instead of running it continuously?
answer
- streaming bookkeeping, batch-shaped cost
- the cluster does not have to be awake
- it drains the backlog, then stops
- several batches, not one giant one
- it replaced an older run-once trigger
basics
~20 sUse Trigger.AvailableNow when the latency budget is minutes to hours and cluster cost dominates: it processes everything currently available in several rate-limited micro-batches, then stops, keeping checkpointed offsets and exactly-once semantics while the cluster runs only briefly.
solid answer
~50 s`Trigger.AvailableNow` turns a streaming query into a bounded run: it consumes all data available at start, in **multiple** micro-batches that honour rate limits like `maxFilesPerTrigger` and `maxOffsetsPerTrigger`, then terminates. You keep everything the streaming engine gives you — the offsets write-ahead log, exactly-once semantics against the sink, stateful operators, incremental reads — but pay for a cluster that runs ten minutes an hour instead of continuously. Choose it when the SLA is measured in minutes or hours, when arrival is bursty, when idle cluster cost dominates, or when you want a natural maintenance window. Stay always-on when latency must be seconds, when state is large enough that reloading it each run dominates the run, or when downstream expects a live feed. It replaced `Trigger.Once`, which processed everything in a single batch and ignored those rate limits, risking an out-of-memory failure on a large backlog.
code
python · 10 lines# always-on: fixed cadence, cluster runs continuously
(stream.writeStream
.trigger(processingTime="30 seconds")
.option("checkpointLocation", "/ckpt/events").start())
# scheduled: drain the backlog in rate-limited batches, then stop
(stream.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", "/ckpt/events").start()
.awaitTermination())go deeper
Know that the trigger controls when micro-batches run, that the default keeps running continuously, and that Trigger.AvailableNow drains what is available and then stops the query.
Explain that AvailableNow uses multiple rate-limited micro-batches rather than one, and that offsets, state and watermark all persist in the checkpoint so the next run resumes rather than restarting.
Show you weigh state reload cost and event-time latency against cluster spend, and that you check the consumer's freshness requirement before proposing a schedule. Be ready to explain why Trigger.Once was replaced.
Own the platform position: one incremental, exactly-once expression of every pipeline with freshness as a scheduling parameter, versus maintaining parallel batch and streaming code. Name the obligations it creates — checkpoint lifecycle, source retention, shuffle-partition sizing fixed at launch.
## The four triggers A Structured Streaming query's trigger decides when micro-batches are planned, and the choice is more consequential than it looks because it determines the *shape of the compute bill* as much as the latency. - **Default (no trigger specified).** Spark starts the next micro-batch as soon as the previous one finishes. Latency is as low as micro-batching allows; the cluster is busy continuously. - **`Trigger.ProcessingTime("30 seconds")`.** A batch is started on a fixed interval. If a batch overruns the interval the next one starts immediately after; if it finishes early Spark waits. This is the usual always-on setting because it makes batch size predictable. - **`Trigger.AvailableNow()`.** Process everything available at start, across as many micro-batches as the rate limits imply, then stop the query. Added in Spark 3.3. - **`Trigger.Continuous("1 second")`.** Experimental long-running task model with millisecond latency, at-least-once only, and support for a narrow set of map-like operations. Rarely the right answer in production. `Trigger.Once()` — process everything in exactly one batch and stop — was deprecated in Spark 3.4 in favour of `AvailableNow`. The reason is the failure mode: a single batch ignores `maxFilesPerTrigger` and `maxOffsetsPerTrigger`, so a query that has been paused for a week tries to plan one enormous batch and either takes hours or dies. `AvailableNow` fixes exactly that by keeping the rate limits and simply looping until the backlog at start time is exhausted. ## What scheduling actually buys you The pitch for `AvailableNow` is that you get streaming's *bookkeeping* with batch's *cost profile*. The bookkeeping is the valuable half. The query still writes its offset range to `offsets/N` before processing and its commit after, so a run that dies halfway is replayed over identical input on the next run. It still reads only what arrived since last time, with no `WHERE ingest_date = ...` predicate to get wrong and no manual high-water-mark table to maintain. Stateful operators still work, with state carried in the checkpoint from run to run. Exactly-once against an idempotent sink still holds. You are not giving any of that up by scheduling it — you are only choosing when the compute exists. The cost profile is the other half. An always-on query holds a cluster twenty-four hours a day whether or not data is arriving. Many pipelines feed dashboards refreshed hourly, or downstream tables consumed by a morning report. For those, running for ten minutes each hour is the same result at a fraction of the spend, and it comes with an operational bonus: there is a moment every hour when nothing is running, which makes deploys, config changes and library upgrades ordinary instead of delicate. ## When to stay always-on **Latency.** If a consumer needs results within seconds — fraud scoring, alerting, an operational feed — a scheduled run cannot deliver it, and no trigger interval on a scheduled job closes the gap. **Large state.** A scheduled run reloads the state store from the checkpoint at start. If the query holds tens of gigabytes of session or dedup state, that reload can dominate the run, and running every fifteen minutes means paying it four times an hour. Past a certain state size the always-on query is genuinely cheaper. **Watermark behaviour.** Watermarks advance from event times in data read, and they persist across runs in the checkpoint, so scheduled runs do not lose watermark progress. But append-mode output is still gated by the watermark, and a run that reads a small backlog may not advance it far enough to close windows — meaning output appears in bursts tied to arrival, not to your schedule. Reason about the event-time latency, not the wall-clock cadence. **Downstream expectations.** A topic that other teams consume as a live feed is a contract. Turning it into an hourly batch is a cross-team change, not a cost optimisation. ## The judgment, framed for a lead The honest way to make this call is to separate two questions that are usually conflated: *how fresh does the consumer need the data* and *how do we want to express the computation*. Structured Streaming with `AvailableNow` decouples them — it lets incremental, exactly-once, stateful processing be the default expression for everything, with freshness as a scheduling parameter you can turn up for the handful of pipelines that need it and down for the majority that do not. That is a better platform position than maintaining two codebases, one batch and one streaming, for the same logic and discovering they have drifted. The cost of that position is that every pipeline carries a checkpoint whose lifecycle you must manage — source retention long enough to cover the gap between runs, a documented procedure for rebuilding a checkpoint, and shuffle-partition sizing fixed at launch for stateful queries. Those are real obligations, and they are the ones worth naming when someone proposes making everything a scheduled streaming job. ## Migration note If you inherit jobs using `Trigger.Once`, moving them to `AvailableNow` is usually a one-line change that improves their worst case: the same semantics, but a backlog is drained in bounded batches instead of one unbounded one.
- Why was Trigger.Once deprecated in favour of Trigger.AvailableNow?`Trigger.Once` processes the entire available backlog in a single micro-batch and ignores rate limits such as `maxFilesPerTrigger` and `maxOffsetsPerTrigger`. After an outage that means one enormous batch that runs for hours or fails on memory. `AvailableNow` drains the same backlog across multiple rate-limited batches and then stops, giving identical semantics with a bounded worst case. Deprecated in Spark 3.4.
- What does a scheduled query pay on every run that an always-on query pays once?Cluster startup and, for stateful queries, reloading the state store from the checkpoint. With small state that is negligible; with tens of gigabytes of session or dedup state it can dominate a short run, and running every fifteen minutes means paying it four times an hour. That crossover is the main reason to keep a large stateful pipeline always-on.
- Does scheduling a query lose watermark progress between runs?No — the watermark is checkpointed along with offsets and state, so the next run resumes from where it left off. What can still surprise you is append-mode output: a run that reads only a small backlog may not advance the watermark past a window's end, so results appear in bursts driven by data arrival rather than by your schedule.
- What obligations does standardising on scheduled streaming create?Every pipeline now owns a checkpoint. Source retention must exceed the gap between runs or a resumed query cannot find its offsets. You need a documented procedure for rebuilding a checkpoint and reprocessing safely into idempotent sinks. And for stateful queries, `spark.sql.shuffle.partitions` is fixed at first launch, so it must be sized before the pipeline ships.
saying these in an interview costs you the question
- Says Trigger.AvailableNow processes everything in one micro-batch
- Thinks scheduling a query gives up exactly-once semantics
- Assumes the watermark restarts from scratch on every scheduled run
- Treats an always-on query as always cheaper because it avoids startup
- Recommends Trigger.Continuous for general production workloads