What you will be able to do
- Explain how Structured Streaming runs a batch-style query incrementally over an unbounded input
- Tell stateless queries apart from stateful ones, and say when output mode matters
- Pick the trigger mode for a given latency and cost need, and spot the deprecated and unsupported ones
Key concept
Incremental execution over an unbounded input — You write a streaming query with the same DataFrame operations you would use on a static table. The engine treats the source as a table that never ends, and on each trigger it processes only the new data and updates the result.
1.One API for batch and streaming, run incrementally
Structured Streaming does not give you a separate streaming language. The documentation says it "lets you express computation on streaming data in the same way you express a batch computation on static data." You build a streaming DataFrame with spark.readStream instead of spark.read, and after that you use the familiar transformations: select, filter, SQL functions, and even MLflow models loaded as UDFs. What changes is how the query runs. The engine does not compute the answer once. It runs the same logical plan again and again over each new slice of input and updates the result each time.
raw_df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
)The answer is none. As with batch reads, a streaming read is lazy. Transformations on raw_df only add instructions to the plan, and nothing starts until an action begins the stream. Calling display() in a notebook starts a streaming job, but in production the action that starts a stream should normally be a write to a sink.
The engine also changes which operations make sense. A streaming source has no last row, and the docs say Structured Streaming "treats data sources as unbounded or infinite datasets." In other words, the engine treats the source as a table that never ends, and each trigger processes only the rows that have arrived since the last one. Because of that, some transformations are not supported on streams: they would require sorting an infinite number of items. You can't sort something that never finishes arriving.
Checkpoint 1 of 4· Exam question
In Spark Structured Streaming's programming model, how does the engine conceptually treat a streaming data source?
Correct answer: A — As a table that is continuously appended to, where each new batch of data is treated as new rows added to an unbounded input table
- A. This is correct: Structured Streaming models a data stream as an unbounded input table that grows with each arriving batch, and queries are expressed as if run against that table, which lets the same DataFrame/Dataset API work for both batch and streaming.
- B. This is incorrect because the engine does not overwrite prior data in place; new data is appended to the conceptual input table rather than replacing what came before it.
- C. This is incorrect because Structured Streaming's execution model is not built on merging independent RDD checkpoints at query completion; it incrementally updates a result table on each trigger.
- D. This is incorrect because the engine tracks incremental offsets from the source rather than reloading the entire dataset from disk on every trigger.
Checkpoint 2 of 4· Check yourself
A colleague defines a streaming DataFrame with spark.readStream, adds a select with current_timestamp(), and then sees no activity on the cluster. What explains this?
Defining a streaming read and its transformations only builds the plan. The stream starts when an action runs, and in production that action is a write to a sink.
“Like other read operations on Databricks, configuring a streaming read does not actually load data.”Source: docs.databricks.com
2.Stateless and stateful queries
Once the engine runs your query in repeated passes, the next question is whether a pass needs to remember anything from earlier passes. The docs split queries into two kinds: "Stateless queries process rows without retaining state. Stateful queries maintain intermediate state for aggregations, joins, and deduplication."
A projection like the enrichment select above is stateless. Each row can be transformed and written without looking at any other row. A running count per customer is stateful. The engine has to keep each customer's current count between passes so that it can add new rows to it. This difference decides how much the engine has to persist and recover, and it decides whether output mode has any effect.
An output mode determines "which records the query's operators emit during each trigger." There are three: append (the default), update, and complete. For a stateless operator the answer is always the same: the records it emits during a trigger are the source records processed during that trigger. So for stateless streaming, all output modes behave the same, and only stateful streams that contain aggregations need an output mode configured. Joins only support append mode, and output mode does not impact deduplication.
Stateful queries also bring in event time and watermarks. Watermarks control how long Structured Streaming waits for late-arriving data in stateful operations. In append mode, stateful operators use the watermark to decide when a row can no longer change and is safe to emit. The rest of the output-mode detail is covered in the lesson on output modes and sinks.
Checkpoint 3 of 4· Check yourself
Which streaming query has its emitted records changed by the output mode setting?
Only the aggregation keeps state, so only its output depends on the mode. The other three are stateless, and every output mode emits the same records for them.
“Only stateful streams containing aggregations require an output mode configuration.”Source: docs.databricks.com
3.Micro-batches and the trigger that schedules them
Each incremental pass is a micro-batch. The trigger decides when the engine checks the source and starts the next one. In the docs' words, "Trigger intervals control how frequently Structured Streaming checks for new data." The trigger is your main control for balancing latency against cost. Checking more often lowers latency but adds overhead. With cloud storage sources, frequent checks can also produce unexpected storage API costs.
| Trigger mode | Syntax (Python) | Behaviour and fit |
|---|---|---|
| Unspecified (Default) | N/A | Equivalent to processingTime trigger with 0 ms intervals, so the next micro-batch starts as soon as the previous one finishes; the docs describe it as general-purpose streaming with 3-5 second latency |
| Processing Time | .trigger(processingTime='10 seconds') | Fixed interval micro-batches; reduces overhead by not checking for data too often |
| Available Now | .trigger(availableNow=True) | Scheduled incremental batch: processes everything available when triggered, then stops |
| Real-time mode | .trigger(realTime='5 minutes') | Sub-second operational workloads; '5 minutes' is the length of a micro-batch, chosen to minimize per-batch overhead such as query compilation. The batch length is not the latency: end-to-end latency is under 1 second at the tail |
| Continuous | .trigger(continuous='1 second') | Not supported on Databricks; experimental in Spark OSS. Use real-time mode instead |
Several of these rows show up as exam distractors. Time-based triggers are what Structured Streaming calls "fixed interval micro-batches." Trigger.Once is older than AvailableNow, and in Databricks Runtime 11.3 LTS and above it is deprecated: use Trigger.AvailableNow for all incremental batch processing workloads. AvailableNow takes all available records as an incremental batch, and you can size it with options such as maxBytesPerTrigger. Spark's Continuous Processing trigger "has been classified as experimental since Spark 2.3," and Databricks neither supports nor recommends it. On serverless compute, only Trigger.AvailableNow() and Trigger.Once() are supported.
The tutorial write below uses availableNow so that the stream processes every record not yet processed, then shuts down. It is safe to run in a notebook without leaving a stream running.
transformed_df.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", checkpoint_path)
.option("path", target_path)
.start()Checkpoint 4 of 4· Check yourself
A team on Databricks Runtime 14.3 LTS wants a nightly job that processes all new files since the last run and then stops. Which trigger should they use?
AvailableNow is the incremental batch trigger. Once is deprecated from 11.3 LTS, a 24-hour processingTime trigger keeps a cluster running all day, and Continuous is not supported.
“In Databricks Runtime 11.3 LTS and above, Trigger.Once is deprecated. Use Trigger.AvailableNow for all incremental batch processing workloads.”Source: docs.databricks.com
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Trigger.Once is still the recommended way to run a streaming query as a scheduled incremental batch.Why is that wrong?
Trigger.Once is deprecated in Databricks Runtime 11.3 LTS and above. AvailableNow replaces it and also lets you control batch size.
Covered in Micro-batches and the trigger that schedules them
2.The Continuous trigger is the supported low-latency option on Databricks.Why is that wrong?
Continuous Processing is still experimental in Spark OSS, and Databricks does not support it. Real-time mode is the low-latency option.
Covered in Micro-batches and the trigger that schedules them
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Structured Streaming lets you express computation on streaming data in the same way you express a batch computation on static data.”
↩︎ One API for batch and streaming, run incrementally“Stateless queries process rows without retaining state. Stateful queries maintain intermediate state for aggregations, joins, and deduplication.”
↩︎ Stateless and stateful queries“Control how long Structured Streaming waits for late-arriving data in stateful operations.”
↩︎ Stateless and stateful queries“The Structured Streaming engine performs the computation incrementally and continuously updates the result as streaming data arrives.”
↩︎ Key concept - 2.
“some transformations are not supported in Structured Streaming workloads because they would require sorting an infinite number of items.”
↩︎ One API for batch and streaming, run incrementally“For most Structured Streaming use cases, the action that triggers a stream should be writing data to a sink.”
↩︎ One API for batch and streaming, run incrementally“instructs Structured Streaming to process all previously unprocessed records from the source dataset and then shut down”
↩︎ Micro-batches and the trigger that schedules them“Like other read operations on Databricks, configuring a streaming read does not actually load data.”
↩︎ Checkpoint - 3.
“The records a stateless operator emits during a trigger are always the source records processed during that trigger.”
↩︎ Stateless and stateful queries“Joins only support the append output mode, and output mode doesn't impact deduplication.”
↩︎ Stateless and stateful queries“Stateful operators use the watermark to determine when this happens.”
↩︎ Stateless and stateful queries“Only stateful streams containing aggregations require an output mode configuration.”
↩︎ Checkpoint - 4.
“Trigger intervals control how frequently Structured Streaming checks for new data.”
↩︎ Micro-batches and the trigger that schedules them“On serverless compute, only Trigger.AvailableNow() and Trigger.Once() are supported.”
↩︎ Micro-batches and the trigger that schedules them“Use 5 minutes to minimize per-batch overhead such as query compilation.”
↩︎ Micro-batches and the trigger that schedules them“In Databricks Runtime 11.3 LTS and above, Trigger.Once is deprecated.”
↩︎ Exam trap 1“This mode has been classified as experimental since Spark 2.3. Databricks doesn't support or recommend this mode.”
↩︎ Exam trap 2“Equivalent to processingTime trigger with 0 ms intervals.”
↩︎ Prediction“In Databricks Runtime 11.3 LTS and above, Trigger.Once is deprecated. Use Trigger.AvailableNow for all incremental batch processing workloads.”
↩︎ Checkpoint