Courseiva

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.

33 questions9 easy12 medium12 hard

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.

1

Which component in Structured Streaming is responsible for providing fault tolerance and ensuring data is processed exactly once?

Easy
2

A 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?

Medium
3

A 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.)

Medium
4

A 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?

Easy
5

You 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?

Medium
6

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)

Hard
7

A 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?

Hard
8

An 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?

Medium
9

Which of the following describes the function of the checkpoint location in a Structured Streaming query?

Easy
10

You 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?

Medium
11

You 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.)

Medium
12

You 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?

Medium
13

A 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?

Easy
14

A 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?

Hard
15

You 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?

Medium
16

Which state management strategy should you implement if you notice your streaming aggregation query is failing due to excessive memory consumption on the executor nodes?

Hard
17

You 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.)

Medium
18

A 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?

Hard
19

A streaming job using 'mapGroupsWithState' is failing due to excessive memory usage. Which strategy is most effective for mitigating this?

Hard
20

What is the primary difference between a micro-batch streaming query and a continuous processing query in Spark?

Easy
21

You 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?

Easy
22

You 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?

Easy
23

A 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?

Easy
24

A 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?

Medium
25

You 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?

Hard
26

A 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?

Hard
27

A 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?

Hard
28

A 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?

Medium
29

You 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?

Hard
30

A 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?

Easy
31

Which THREE of the following sources support streaming read operations in Spark Structured Streaming?

Medium
32

You 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?

Hard
33

You 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?

Hard

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.
databricks-spark-developer-associate DATABRICKS-SPARK-DEVELOPER-ASSOCIATE structured streaming Practice Questions