What you will be able to do
- Create a streaming DataFrame with spark.readStream and explain why defining it does not start processing
- Build a streaming write with writeStream, a checkpoint location, and start() or toTable()
- Choose between append, update and complete output modes and predict what each emits on a trigger
Key concept
Streaming DataFrame — A streaming DataFrame is an ordinary DataFrame whose source never ends. You write the same transformations you would write for a batch job, and the engine runs them incrementally as new data arrives.
1.Creating a streaming DataFrame with spark.readStream
Every streaming query starts at spark.readStream, which returns a DataStreamReader. It works like the batch spark.read: you set a format, an optional schema, and option values, then call load(path). There are also shortcuts such as json(path), parquet(path), csv(path) and table(tableName), which loads a streaming Delta table. What comes back is a DataFrame. In Scala the same reader gives you a streaming Dataset: the Databricks docs use one name, sdf, for "a streaming DataFrame/Dataset generated with sparkSession.readStream". The sources describe no Dataset-specific behaviour beyond that, so treat the two as created and written the same way.
The engine treats the source as unbounded. Because of that, some transformations that need to sort an infinite number of rows are not supported. Most row-level operations work as they do in batch, such as select, adding columns, and enrichment with SQL functions. The example below uses Auto Loader (cloudFiles) to create a streaming DataFrame from a directory of JSON files.
raw_df = (spark.readStream
.format("cloudFiles")
.option("cloudFiles.format", "json")
.option("cloudFiles.schemaLocation", checkpoint_path)
.load(file_path)
)So what starts the stream? In a notebook, calling display() on a streaming DataFrame starts a streaming job, which helps while you are exploring. Databricks is clear that this is not the production pattern: the action that triggers a stream should normally be a write to a sink. That write is the subject of the next section.
Checkpoint 1 of 5· Check yourself
Which action do the Databricks docs say should normally trigger a production Structured Streaming query?
display() does start a job, but the docs say most use cases should start the stream by writing to a sink. load() only defines the source.
“For most Structured Streaming use cases, the action that triggers a stream should be writing data to a sink.”Source: docs.databricks.com
2.Writing the stream: writeStream, checkpoints and start()
df.writeStream returns a DataStreamWriter, the streaming counterpart of df.write. You configure the sink in a chain of calls, and nothing runs until you call one of the terminal methods: start() or toTable(). Each returns a StreamingQuery object that you can monitor and stop.
| Method | What it does |
|---|---|
| outputMode(outputMode) | Specifies how data of a streaming DataFrame is written to the sink: append, complete, or update |
| format(source) | Specifies the output data source format |
| option(key, value) | Adds an output option, such as checkpointLocation or path |
| trigger(**kwargs) | Sets the trigger for the streaming query execution |
| foreachBatch(func) | Sends each micro-batch's output to the given function |
| start(path) | Starts the query and returns a StreamingQuery object |
| toTable(tableName) | Starts the query, continually writing results to the given table (table() is an alias) |
Checkpoint 2 of 5· Fill the gap
This sample reads a rate stream and writes it to the console. Which accessor completes it?
import time
df = spark.readStream.format("rate").load()
df = df.selectExpr("value % 3 as v")
q = df. ? .format("console").start()
time.sleep(3)
q.stop()A streaming DataFrame is written through df.writeStream, which returns a DataStreamWriter. df.write is the batch writer and does not work on a streaming DataFrame.
Source: docs.databricks.comA production writer always sets checkpointLocation. The checkpoint gives the stream its identity: it records which data has been processed and any state the query holds, so a restarted query continues where it stopped. Every writer needs its own location, and you must set it before the query runs. The tutorial example below writes to a file path and uses the availableNow trigger, which processes all records not yet processed and then shuts down.
transformed_df.writeStream
.trigger(availableNow=True)
.option("checkpointLocation", checkpoint_path)
.option("path", target_path)
.start()No. The checkpoint is what identifies a stream, and the docs say each writer must have a unique checkpoint location. If two queries share one, each query's progress tracking corrupts the other's.
Checkpoint 3 of 5· Check yourself
What does calling start() on a configured DataStreamWriter return?
start() and toTable() both start execution and return a StreamingQuery, which you can later stop or wait on.
“Starts the execution of the streaming query and returns a StreamingQuery object.”Source: docs.databricks.com
3.Output modes: append, update and complete
outputMode() decides which result rows a query sends to the sink on each trigger. It only changes anything when the query is stateful, for example a streaming aggregation, where a result row such as a window's total can change from one trigger to the next. For stateless queries all three modes behave the same: on each trigger, they emit the source rows that trigger processed. If you set nothing, the query runs in append mode.
| Mode | What it emits on each trigger | Note |
|---|---|---|
| append (default) | Only rows that will not change in future triggers | Stateful operators use the watermark to decide when a row is final |
| update | All rows that changed during this trigger | A row may be emitted again if it changes later |
| complete | Every result row the operator has ever produced | Only works with streaming aggregations; can perform poorly as data scales |
In the same first trigger, update mode emits both windows because both were just modified, and complete mode emits every row in the state. Next, a $20 record at 3:20pm arrives. The watermark moves to 3:05pm, which is past the end of the [2pm, 3pm] window, so that window can no longer change. Now append mode emits [2pm, 3pm] $25 exactly once. Update mode emits only [3pm, 4pm], whose total went from $30 to $50, and complete mode again emits everything. The trade-off follows from this. Append makes you wait at least the watermark delay, but it writes each row once, which suits a downstream service that acts on every write. Update keeps the sink fresh, but it writes once per trigger per changed aggregate, which can be expensive if the sink charges per write.
Checkpoint 4 of 5· Match them up
Match each output mode to what it writes for a stateful aggregation
Tap a term, then the definition that fits it.
The docs summarise the stateful behaviour of each mode in exactly these terms.
“In update mode, write records that have changed since the previous trigger.”Source: docs.databricks.com
Checkpoint 5 of 5· Exam question
A data engineer builds a streaming DataFrame from JSON files in cloud storage. The code below throws an `AnalysisException` at start time because the file source has no schema to work with. Which single change fixes the error while keeping the same file source? ```python raw_stream = ( spark.readStream .format("json") .load("/mnt/raw/events") ) ```
Correct answer: A — Add `.schema(event_schema)` before `.load()`, supplying an explicit `StructType` that matches the JSON files.
- A. File-based streaming sources require an explicit schema by default so that the query behaves consistently across restarts, and supplying a `StructType` up front satisfies that requirement. This is the documented fix for schema errors on file sources such as JSON.
- B. There is no per-source `inferSchema` streaming option that enables inference on `readStream`; schema inference for streaming file sources is controlled separately and is disabled by default. Setting this option does not resolve the missing-schema error.
- C. Switching to the text format changes what the source produces entirely, turning each line into a single raw string column instead of parsed JSON fields, so the pipeline would lose its structured columns rather than fixing the schema error.
- D. The `multiLine` option only controls how JSON records spanning multiple lines are parsed once a schema is known; it does not cause Spark to infer a schema for a streaming file source, so the error remains.
- E. `spark.read.json(...)` produces a static, non-streaming DataFrame, and `.toDF()` does not turn a static DataFrame into a streaming source, so this approach cannot be used to build a streaming query at all.
Sources5
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Calling spark.readStream...load() starts reading data from the source.Why is that wrong?
load() only defines a streaming DataFrame. Data is read once an action starts the stream, usually writeStream...start() or toTable().
Covered in Creating a streaming DataFrame with spark.readStream
2.Complete mode can be used on any streaming query to rewrite the full result each trigger.Why is that wrong?
Complete mode is valid only for streaming aggregations. A plain select/filter stream cannot use it.
Covered in Output modes: append, update and complete
3.Switching a stateless projection stream from append to update changes which rows reach the sink.Why is that wrong?
Output mode affects only stateful operators. For stateless streams every mode emits the rows processed in that trigger.
Covered in Output modes: append, update and complete
Sources
Every claim above is drawn from one of these pages, quoted as it was written on the date shown.
- 1.
“Interface used to load a streaming DataFrame from external storage systems (for example, file systems and key-value stores).”
↩︎ Creating a streaming DataFrame with spark.readStream - 2.
“sdf represents a streaming DataFrame/Dataset generated with sparkSession.readStream.”
↩︎ Creating a streaming DataFrame with spark.readStream“You must specify the checkpointLocation option before you run a streaming query”
↩︎ Writing the stream: writeStream, checkpoints and start() - 3.
“Structured Streaming treats data sources as unbounded or infinite datasets.”
↩︎ Creating a streaming DataFrame with spark.readStream“Calling display() on a streaming DataFrame starts a streaming job.”
↩︎ Creating a streaming DataFrame with spark.readStream“Always make sure you specify a unique checkpoint location for each streaming writer you configure.”
↩︎ Writing the stream: writeStream, checkpoints and start()“You must trigger an action on the data before the stream begins.”
↩︎ Exam trap 1“Like other read operations on Databricks, configuring a streaming read does not actually load data.”
↩︎ Prediction“For most Structured Streaming use cases, the action that triggers a stream should be writing data to a sink.”
↩︎ Checkpoint - 4.
“Interface used to write a streaming DataFrame to external storage systems”
↩︎ Writing the stream: writeStream, checkpoints and start()“Starts the execution of the streaming query and returns a StreamingQuery object.”
↩︎ Checkpoint - 5.
“By default, streaming queries run in append mode.”
↩︎ Output modes: append, update and complete“In update mode, operators emit all rows that changed during the trigger, even if the emitted record might change in a subsequent trigger.”
↩︎ Output modes: append, update and complete“Append mode forces stateful operators to emit results only after stateful results are finalized, which is at least as long as your watermark delay.”
↩︎ Output modes: append, update and complete“Update mode results in one write per trigger per aggregate value.”
↩︎ Output modes: append, update and complete“Complete mode only works with streaming aggregations.”
↩︎ Exam trap 2“For stateless streaming, all output modes behave the same.”
↩︎ Exam trap 3“The streaming aggregation operator does not emit anything downstream.”
↩︎ Prediction“In update mode, write records that have changed since the previous trigger.”
↩︎ Checkpoint
Also cited
“Structured Streaming lets you express computation on streaming data in the same way you express a batch computation on static data.”
↩︎ Key concept