Databricks-Spark-Assoc · domain
Structured Streaming
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.
Focused practice
Practice Structured Streaming questions
Scored sessions drawing only from this domain — pick a length below.
Start 20-question practice test →What this domain covers
What to know about Structured Streaming
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.
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
Watch out for
Common Structured Streaming exam traps
- ▸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
Question index
All Structured Streaming questions (33)
Click any question to see the full explanation, or start a practice session above.
Which component in Structured Streaming is responsible for providing fault tolerance and ensuring data is processed exactly once?
Easy2A 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?
Medium3A 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.)
Medium4A 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?
Easy5You 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?
Medium6An 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)
Hard7A 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?
Hard8An 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?
Medium9Which of the following describes the function of the checkpoint location in a Structured Streaming query?
Easy10You 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?
Medium11You 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.)
Medium12You 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?
Medium13A 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?
Easy14A 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?
Hard15You 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?
Medium16Which state management strategy should you implement if you notice your streaming aggregation query is failing due to excessive memory consumption on the executor nodes?
Hard17You 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.)
Medium18A 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?
Hard19A streaming job using 'mapGroupsWithState' is failing due to excessive memory usage. Which strategy is most effective for mitigating this?
Hard20What is the primary difference between a micro-batch streaming query and a continuous processing query in Spark?
Easy21You 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?
Easy22You 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?
Easy23A 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?
Easy24A 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?
Medium25You 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?
Hard26A 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?
Hard27A 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?
Hard28A 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?
Medium29You 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?
Hard30A 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?
Easy31Which THREE of the following sources support streaming read operations in Spark Structured Streaming?
Medium32You 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?
Hard33You 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?
HardOther domains
All Databricks-Spark-Assoc exam domains
Frequently asked questions
- What does the Structured Streaming domain cover on the Databricks-Spark-Assoc exam?
- 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.
- How many questions are in this domain?
- This page lists all 33 Structured Streaming questions in the Databricks-Spark-Assoc question bank. The actual exam draws from this domain proportionally to its weighting in the official exam blueprint.
- What is the best way to practise this domain?
- Start with a short focused session (10 questions) to identify gaps, then work through explanations. Repeat with a longer session once the weak areas feel solid.
- Can I practise only Structured Streaming questions?
- Yes — the session launcher on this page filters questions to this domain only. Choose any session length for inline explanations and scoring.