A candidate must configure streaming reads, apply event-time windows with watermarks, and manage stateful operators correctly. The single most important thing: pair watermarks with stateful operations and use a checkpoint plus an idempotent sink to achieve exactly-once processing.
Start practicing
Structured Streaming — choose a session length
Free · No account required
Domain overview
Structured Streaming on Databricks covers reading from Delta, Kafka, and file sources, applying event-time windowing with watermarks, and managing stateful operations like dropDuplicates and foreachBatch. Questions test source support, late-data handling, exactly-once semantics via checkpoints, and diagnosing state-store or watermark failures in production pipelines.
Exam objectives
Streaming reads from Delta tables, Kafka topics, and file/rate sources using readStream
Event-time tumbling windows with watermark to bound late-arriving sensor data
Stateful deduplication via dropDuplicates on composite keys and state store limits
Exactly-once delivery using checkpointLocation with foreachBatch and idempotent sinks
Assuming dropDuplicates state grows unbounded; without a watermark the state store eventually fails on long-running jobs
Confusing processing-time windows with event-time windows, so late data is dropped instead of admitted within the watermark
Believing checkpointing alone guarantees exactly-once; the sink must also be idempotent or transactional like Delta
Click any question to see the full explanation and answer options, or start a focused practice session above.
An analytics engineer is designing a Structured Streaming pipeline that reads from an Apache Kafka topic and writes JSON-formatted output to cloud object storage. The pipeline must maintain exactly-once processing semantics and support automatic recovery from cluster restarts. Which TWO actions are mandatory to achieve these requirements? (Choose 2)
2You are processing a streaming dataset of sensor readings. You need to calculate the average temperature every 10 minutes, allowing data to arrive up to 2 minutes late. Which windowing approach correctly handles this requirement in Structured Streaming?
3Which of the following describes the function of the checkpoint location in a Structured Streaming query?
4You are processing JSON data from a stream. You need to extract nested fields from the JSON structure. Which function is the most efficient and standard way to handle this in Structured Streaming?
5What is the primary difference between a micro-batch streaming query and a continuous processing query in Spark?
6Which state management strategy should you implement if you notice your streaming aggregation query is failing due to excessive memory consumption on the executor nodes?
7Which component in Structured Streaming is responsible for providing fault tolerance and ensuring data is processed exactly once?
8A streaming job using 'mapGroupsWithState' is failing due to excessive memory usage. Which strategy is most effective for mitigating this?
9Which THREE of the following sources support streaming read operations in Spark Structured Streaming?
10An engineer is developing a Structured Streaming job that reads from an Apache Kafka source and writes the output continuously toDelta Lake using outputMode("append"). The stream occasionally experiences late-arriving data. Which downstream behavior can the engineer expect regarding the Delta Lake table?
11A developer is building a Structured Streaming pipeline that reads from a Kafka topic and writes to a Delta table. The pipeline must tolerate occasional downstream failures and reprocess data without duplicates. The developer sets a checkpoint location and uses the default output mode. Which statement correctly describes how the checkpoint location contributes to fault tolerance in this scenario?
12A developer writes a Structured Streaming query that reads from a Kafka topic with `spark.readStream.format("kafka")` and then calls `.writeStream.format("console").start()`. The query runs, but after a few minutes the driver logs show that the query is only processing newly arriving offsets and older messages in the topic are never read. What is the most likely cause?
13You are building a Structured Streaming job that reads from a Kafka source and writes to a Delta table. You need to ensure that the job can recover from failures and process data exactly once. Which TWO of the following are required to achieve exactly-once semantics? (Choose two.)
14You are running a Structured Streaming query on Databricks that reads from a Kafka topic and writes to a Delta table. The query uses `option("maxOffsetsPerTrigger", 10000)` to limit the number of records per micro-batch. During a peak, the Kafka topic accumulates a large backlog. You notice that the query is processing data but the backlog is not decreasing. What is the most likely cause?
15A streaming DataFrame `df` has a watermark defined on `eventTime` with a delay of 5 minutes. The query uses `withWatermark("eventTime", "5 minutes")` and writes to a Delta table in `append` output mode. A record with event time 10:00 arrives when the watermark is at 10:10. What happens to this record?
16You are developing a Structured Streaming job that reads from a Delta table and writes to another Delta table. You need to ensure that the streaming query can recover from failures and continue processing without data loss or duplication. Which of the following must be configured?
17A developer is building a Structured Streaming job that reads from a Kafka topic and writes to a Delta table. The job must handle late data up to 15 minutes and ensure that aggregations are updated correctly. The developer adds a watermark of 15 minutes on the event time column. What is the effect of this watermark on the aggregation state and output?
18A developer is building a Structured Streaming job that reads from a Delta table as a stream and writes to another Delta table. The job must support exactly-once processing and allow the output to be updated incrementally. Which two options are required to achieve exactly-once semantics? (Choose two.)
19A developer wants to start a Structured Streaming query that reads from a Delta table and writes to another Delta table, and needs the query to process all existing data in the source table on its first run. Which option should be set on the read stream?
20You are performing a stream-stream join between two streaming DataFrames, `orders` and `payments`, both with watermarks defined on their event time columns. The join condition is `orders.orderId == payments.orderId` and it is an inner join. What happens to state in this join?
21A developer is building a Structured Streaming job in PySpark that reads from a Kafka topic and writes to a Delta Lake table. The job uses `outputMode("append")` and a 10-minute watermark on the event-time column `event_time`. A batch of late data arrives with events whose `event_time` is older than the watermark. What happens to these late events?
22You are monitoring a Structured Streaming query in Databricks and want to see the current status, including the number of input rows per second and the batch duration. Which of the following is the most direct way to access this information?
23A developer wants to start a Structured Streaming query that reads from a Kafka topic and writes to the console for debugging. The developer uses `writeStream.format("console").start()`. What is the default trigger for this query?
24You have a Structured Streaming job that reads from a Kafka topic and writes to a Delta table. You need to ensure that the job processes each record exactly once, even after failures. Which of the following should you configure?
25You are building a Structured Streaming pipeline that reads from a Delta table source and applies a stateful deduplication using dropDuplicates on a composite key. After several hours, the job fails with an error indicating that the state store has grown too large. You need to bound the state size while still removing duplicate events that arrive within a reasonable window. Which approach should you take?
26A developer is writing a Structured Streaming query that reads from a JSON file source and writes to the console for debugging. The query uses `outputMode("append")`. Which statement describes the output behavior?
27A data engineer wants to run a Structured Streaming query that reads from a Kafka topic and writes aggregated counts to a console sink for debugging. The query uses a grouping aggregation on a tumbling event-time window. Which output mode must be used so that only rows that changed since the last trigger are emitted?
28A streaming query reads from a rate source and writes to a Delta table using `outputMode("append")`. The developer observes that the query processes data continuously but the Delta table remains empty. Which condition explains why no rows are written?
29You are designing a Structured Streaming job that must read from a file source and write to a Delta table. You want the job to be resilient to failures and to continue processing only new files after a restart. Which two actions should you take? (Choose two.)
30You are writing a Structured Streaming query that reads from a Kafka topic and outputs to the console. You want to see only the newly arrived data in each micro-batch, without aggregations. Which output mode should you use?
31You are troubleshooting a Structured Streaming job that reads from Kafka and writes to a Delta table using foreachBatch. The job occasionally processes the same Kafka offsets twice after a task retry, causing duplicate rows in the Delta table. You want to ensure that each micro-batch's output is applied exactly once. Which change should you make inside the foreachBatch function?
32A developer is using Structured Streaming with a Kafka source and wants to ensure that each message is processed exactly once, even in the event of failures. The developer has set a checkpoint location and is using `foreachBatch` to write to an external database. Which additional step is necessary to achieve exactly-once semantics?
33A Databricks Structured Streaming job reads JSON files from a cloud storage directory and writes aggregated results to a Delta table. The source directory receives new files continuously and files are never modified after being written. The pipeline must tolerate late-arriving event data by up to 30 minutes and must not reprocess already-emitted windows when the query is restarted. Which combination of configurations is required to meet these requirements?
A candidate must configure streaming reads, apply event-time windows with watermarks, and manage stateful operators correctly. The single most important thing: pair watermarks with stateful operations and use a checkpoint plus an idempotent sink to achieve exactly-once processing.
The Courseiva Databricks-Spark-Assoc question bank contains 33 questions in the Structured Streaming domain. Click any question to see the full explanation and answer breakdown.
Start with a 10-question focused session to identify your baseline accuracy in this domain. Read every explanation — even for questions you answer correctly — to understand the reasoning. Once you score consistently above 80%, move to a 20–30 question session to confirm depth before moving to the next domain.
Yes — the session launcher on this page draws questions exclusively from the Structured Streaming domain. Choose 10, 20, 30, or 50 questions for a focused session, or click individual questions to review them one by one.
Save your results, see per-domain analytics, and get readiness scores — free, for every certification.
Sign Up FreeFree forever · Every certification included