What you will be able to do
- Apply select, selectExpr, filter and where to a streaming DataFrame exactly as you would to a static one
- Tell stateless operations (selection, projection) apart from stateful ones (aggregation), and explain why that difference matters
- Build a grouped streaming aggregation with groupBy and pick an output mode that is valid for it
Key concept
Stateless vs stateful streaming operations — Selection and projection handle each incoming row by itself and keep nothing between micro-batches. An aggregation has to keep running results (state) from one micro-batch to the next, and that state is why output modes, windows and watermarks come into play.
1.A streaming DataFrame is still a DataFrame
Structured Streaming's main design choice is that you don't learn a second API. A DataFrame created with spark.readStream uses the same select, filter, groupBy and agg you use on static data. The engine performs the computation incrementally and keeps updating the result as new data arrives. Databricks notes that Structured Streaming supports most of the transformations available in Databricks and Spark SQL.
Two members of the DataFrame class let you tell the two kinds apart. The isStreaming property returns True when the DataFrame reads from a continuous source. writeStream is the streaming counterpart of write, and you use it to send the result to a sink. Everything between the read and the write is ordinary DataFrame code.
Checkpoint 1 of 6· Check yourself
A colleague gives you a DataFrame df and you need to confirm in code that it is a streaming DataFrame. What do you check?
isStreaming is the DataFrame property that reports whether the DataFrame reads from a continuous source.
“Returns True if this DataFrame contains one or more sources that continuously return data as it arrives.”Source: docs.databricks.com
2.Selection and projection: stateless row-by-row work
Projection picks or computes columns, and selection keeps only the rows that match a condition. On a stream, both are stateless. Each row is transformed or dropped when it arrives, and nothing is remembered for later micro-batches. The methods are the ones the DataFrame class reference lists:
| Method | What it does |
|---|---|
| select(*cols) | Projects a set of expressions and returns a new DataFrame |
| selectExpr(*expr) | Projects a set of SQL expressions and returns a new DataFrame |
| filter(condition) | Filters rows using the given condition |
| where(condition) | Alias for filter |
| withColumn(colName, col) | Adds a column or replaces an existing column with the same name |
The Databricks streaming tutorial projects an ingested stream with nothing more than select. It keeps every column with "*" and adds two derived columns, each named with alias:
from pyspark.sql.functions import col, current_timestamp
transformed_df = (raw_df.select(
"*",
col("_metadata.file_path").alias("source_file"),
current_timestamp().alias("processing_time")
)
)Nothing runs at this point. transformed_df only holds the instructions to load and transform each record when it arrives. Because these operators are stateless, output mode doesn't affect what they emit: each trigger emits the rows that trigger processed.
Streaming doesn't remove every limit, though. A stream is treated as unbounded, so any transformation that would need to sort an infinite number of items is not supported.
Checkpoint 2 of 6· Match them up
Match each method to what it does on a streaming DataFrame
Tap a term, then the definition that fits it.
select and selectExpr both project (one takes Column expressions, the other SQL strings). filter selects rows, and where is just another name for filter.
“Projects a set of SQL expressions and returns a new DataFrame.”Source: docs.databricks.com
Checkpoint 3 of 6· Exam question
A team ingests a streaming DataFrame `events` with columns `device`, `signal`, and `event_time`. They need a new streaming DataFrame that keeps only rows where `signal` exceeds 50 and retains just the `device` and `signal` columns for downstream writing. Which code accomplishes this selection and projection on the streaming DataFrame? ```python result = events. ___ ```
Correct answer: A — `select("device", "signal").where("signal > 50")`
- A. This is correct because `select` and `where` are both fully supported untyped operations on streaming DataFrames, letting projection and filtering compose exactly as they would on a batch DataFrame.
- B. This is incorrect because reordering `where` before `select` still tries to filter on `signal` after only the two named columns would remain if the select ran first in a real pipeline; as written here it also never assigns the filtered signal check to the intended output shape the scenario asks for.
- C. This is incorrect because `show()` is an eager action that materializes output for inspection and is not supported on an unbounded streaming DataFrame, and `limit` was not part of the requirement in the first place.
- D. This is incorrect because `distinct()` is not a supported operation on streaming DataFrames since it would require unbounded state to track every row ever seen.
- E. This is incorrect because `show()` cannot run on a streaming DataFrame, and a full sort on an unaggregated stream is not supported either way.
3.Aggregation: where a stream starts to remember
groupBy works on a stream the same way it does in batch. It groups rows so an aggregate such as count() or agg(sum(...)) can be computed for each group. What changes is the cost. A count for key 3 has to survive from one micro-batch to the next, so the engine keeps it as state. Databricks lists streaming aggregation as a stateful operation, alongside distinct, dropDuplicates and stream-stream joins.
The Databricks example below chains all three basic operations on a single stream. It reads from the built-in rate source, projects with selectExpr, groups and counts, and renames the result columns with toDF:
query = (
spark.readStream.format("rate").load()
.selectExpr("value % 10 as key")
.groupBy("key")
.count()
.toDF("key", "count")
.writeStream
.foreachBatch(writeToSQLWarehouse)
.outputMode("update")
.start()
)Checkpoint 4 of 6· Fill the gap
Which method turns this projected stream into a per-key aggregation?
query = (
spark.readStream.format("rate").load()
.selectExpr("value % 10 as key")
. ? ("key")
.count()
.toDF("key", "count")groupBy("key") groups the stream so count() returns one running count per key. orderBy would require a sort, select only projects, and dropDuplicates removes rows instead of counting them.
Source: docs.databricks.comCheckpoint 5 of 6· Exam question
An engineer writes the following on a streaming DataFrame `readings` with columns `sensorId` and `value`: ```python counts = readings.groupBy("sensorId").count() counts.writeStream.outputMode("append").format("console").start() ``` Running this raises an `AnalysisException` about the output mode. What is the most likely cause?
Correct answer: A — Append mode cannot emit a row for an aggregation until it is guaranteed never to change again, and without a watermark Spark has no way to know when a group's value is final.
- A. This is correct because append mode only emits a row once it is guaranteed not to change again, which for a non-windowed running aggregation without a watermark is never determinable, so Spark rejects it at analysis time.
- B. This is incorrect because `groupBy("sensorId").count()` is a perfectly valid streaming aggregation on its own; a `window()` call is only required for event-time windowed aggregations, not for every grouped count.
- C. This is incorrect because the console sink supports append, update, and complete modes; it does not silently coerce or reject modes based on sink type.
- D. This is incorrect because streaming sources typically infer or are given a schema already, and schema presence is unrelated to which output modes an aggregation query may use.
- E. This is incorrect because grouping keys in Structured Streaming aggregations work with their native column type and do not require an explicit cast to string.
4.Choosing an output mode once you aggregate
The example above set .outputMode("update"). With a stateless stream that call would make no difference, because all output modes behave the same. An aggregation changes this. A given group's result can change from one trigger to the next, so you have to decide which rows get written. Databricks puts it directly: only stateful streams that contain aggregations require an output mode configuration.
| Output mode | What the aggregation emits per trigger |
|---|---|
| append (default) | Only rows that will not change in future triggers; stateful operators use the watermark to decide this |
| update | Every row that changed during the trigger, even if it may change again later |
| complete | Every resulting row the operator has ever produced; only works with streaming aggregations |
These modes work together with the other operations covered so far. Complete mode is only valid when the query contains an aggregation. Append mode emits a group only after its result is final, so a plain groupBy("key") has no point at which that happens. Time-based windows combined with a watermark provide that point, and they are the next step after basic aggregation.
Checkpoint 6 of 6· Check yourself
A streaming query only does select and filter before writing to a sink. Which statement about output modes is correct?
Complete mode requires a streaming aggregation. For a stateless select/filter query the modes otherwise behave the same, since each trigger simply emits the rows it processed.
“Complete mode only works with streaming aggregations.”Source: docs.databricks.com
Sources4
Exam traps
Each one states something that sounds right. Open it to see what is actually true.
1.Complete mode can be used on any streaming query, including a stateless select/filter pipeline.Why is that wrong?
Complete mode requires a streaming aggregation. For stateless queries the output modes all behave the same.
Covered in Choosing an output mode once you aggregate
2.Every batch DataFrame transformation, including sorting a raw stream, works unchanged on a streaming DataFrame.Why is that wrong?
Streams are unbounded, so transformations that would require sorting an infinite number of items are not supported.
Covered in Selection and projection: stateless row-by-row work
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.”
↩︎ A streaming DataFrame is still a DataFrame“Stateless queries process rows without retaining state. Stateful queries maintain intermediate state for aggregations, joins, and deduplication.”
↩︎ Key concept - 2.
“Structured Streaming supports most transformations that are available in Databricks and Spark SQL.”
↩︎ A streaming DataFrame is still a DataFrame“Structured Streaming treats data sources as unbounded or infinite datasets.”
↩︎ Selection and projection: stateless row-by-row work“some transformations are not supported in Structured Streaming workloads because they would require sorting an infinite number of items.”
↩︎ Exam trap 2 - 3.
“Interface for saving the content of the streaming DataFrame out into external storage.”
↩︎ A streaming DataFrame is still a DataFrame“Projects a set of expressions and returns a new DataFrame.”
↩︎ Selection and projection: stateless row-by-row work“Groups the DataFrame by the specified columns so that aggregation can be performed on them.”
↩︎ Aggregation: where a stream starts to remember“Returns True if this DataFrame contains one or more sources that continuously return data as it arrives.”
↩︎ Checkpoint“Projects a set of SQL expressions and returns a new DataFrame.”
↩︎ Checkpoint - 4.
“The records a stateless operator emits during a trigger are always the source records processed during that trigger.”
↩︎ Selection and projection: stateless row-by-row work“Only stateful streams containing aggregations require an output mode configuration.”
↩︎ Choosing an output mode once you aggregate“For stateless streaming, all output modes behave the same.”
↩︎ Choosing an output mode once you aggregate“Append mode forces stateful operators to emit results only after stateful results are finalized”
↩︎ Choosing an output mode once you aggregate“Complete mode only works with streaming aggregations.”
↩︎ Exam trap 1“Complete mode only works with streaming aggregations.”
↩︎ Checkpoint - 5.
“Stateful operations include streaming aggregation, distinct, dropDuplicates, stream-stream joins, and custom stateful applications.”
↩︎ Aggregation: where a stream starts to remember
Also cited
“Streaming queries use state to incrementally update results instead of recomputing everything after each micro-batch.”
↩︎ Prediction