Courseiva

Databricks-Spark-Assoc · topic practice

Structured Streaming practice questions

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.

Courseiva uses original exam-style practice questions designed for learning and revision. The goal is to understand the concepts, recognise exam patterns, and improve through explanations — not memorise copied exam dumps.

Editorial oversight:Johnson Ajibi· MSc IT Security, IEEE Senior Member
20 questionsDomain: Structured Streaming

What the exam tests

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

Practice set

Structured Streaming questions

20 questions · select your answer, then reveal the explanation

You are developing a Structured Streaming job in Databricks that reads JSON data from an Auto Loader source, transforms the schema, and writes the output to a Delta Lake table. During execution, downstream consumers complain that the data contains duplicates due to upstream retries. Which operation should you apply to the DataFrame to ensure exactly-once processing semantics before writing to the Delta table?

A data engineer configures a Structured Streaming job to read from an Apache Kafka source and writes the incoming data directly to Delta Lake using the append output mode. During testing, the job throws a AnalysisException stating that streaming queries with aggregation and output mode update require a watermark. Which action resolves this issue correctly?

A streaming pipeline reads from Delta Lake using `spark.readStream.table("events")` and performs stateful aggregation using `groupBy("userId").count()`. The engineering team notices that the shuffle partitions default to 200, which causes excessive task overhead and latency. Which configuration setting should be applied to optimize the shuffle partition count specifically for this Structured Streaming query?

A data engineer is building a Structured Streaming pipeline that reads from an Apache Kafka source and writes the incoming data directly to Delta Lake. The pipeline must execute multiple streaming aggregates across event time, but downstream operational reporting requires complete updates to be written out for every micro-batch. Which output mode should be configured to meet this requirement?

A data engineer writes a Structured Streaming query that reads files from an AWS S3 bucket directory and applies transformations before writing to a Delta table. The source files are continuously dropped into the folder by an upstream ingestion process. However, the engineer notices that newly added files are completely ignored by the stream. What is the most likely cause of this behavior?

You are developing a Structured Streaming job in Databricks that reads from an Apache Kafka source and writes the output to Delta Lake using the append mode. The streaming query fails with a schema evolution error because an upstream producer added a new nested field to the incoming JSON payload. Which approach allows you to seamlessly ingest this evolving schema without failing the running stream?

A data engineer is building a Structured Streaming pipeline that reads from an Apache Kafka source and writes the output in append mode to Delta Lake. The pipeline fails during execution due to an upstream data corruption issue, and the team needs to restart it from a specific historical offset without losing stateful aggregations. Which approach should the engineer use to resume processing correctly?

When reading from a file source in Structured Streaming, which configuration setting must be explicitly enabled if you want the streaming query to process existing files in the directory before streaming new ones?

Which TWO of the following are valid output modes in Structured Streaming that support aggregation operations?

Refer to the exhibit. You are running a streaming query on Databricks. You need to update the trigger to ensure the query processes data as soon as it becomes available, rather than waiting for a 1-minute interval. What should you modify in the configuration?

Exhibit

{
  "queryName": "sensor_stream",
  "checkpointLocation": "/mnt/logs/checkpoint",
  "outputMode": "complete",
  "trigger": "ProcessingTime('1 minute')"
}

You are performing a stream-stream join between two dataframes, `orders` and `clicks`. Both streams have watermarks defined. What happens if the join condition does not include a time range constraint?

Which THREE of the following are true about the 'foreachBatch' sink in Structured Streaming?

What happens if a streaming source produces data with a schema that differs from the one initially defined in the schema inference phase?

Which TWO of the following are valid ways to monitor a running Structured Streaming query on Databricks?

You are processing a streaming dataset with a watermark of 10 minutes. A late record arrives with an event time of 12:00, while the current watermark is 12:15. What happens to this record?

Which TWO statements are true regarding the use of 'foreachBatch' in Structured Streaming?

Refer to the exhibit. You are configuring a streaming job. Based on the provided configuration, what is the impact of the 'checkpointLocation' setting?

Exhibit

{
  "source": "kafka",
  "startingOffsets": "earliest",
  "checkpointLocation": "/mnt/logs/checkpoint",
  "trigger": "ProcessingTime(5 seconds)"
}

A data engineer is using `foreachBatch` in a Structured Streaming job to write to multiple sinks, including a Delta table and an external database. The engineer needs to ensure the operation is idempotent and handles reprocessing correctly. Which TWO statements about `foreachBatch` are correct in this context? (Choose two.)

A developer is testing a Structured Streaming query that reads from a rate source and writes to the console. The developer wants the query to process all available data and then stop automatically. Which trigger type should be used to achieve this behavior?

You are building a Structured Streaming pipeline in Databricks that reads from a Delta table using `spark.readStream.table("transactions")`. The pipeline performs a `groupBy("accountId").count()` aggregation and writes to another Delta table. You need the streaming query to emit updated counts for all accounts as new data arrives, even if an account has no new transactions in the current micro-batch. Which output mode should you use?

Free account

Track your progress over time

Create a free account to save your results and see which topics improve across sessions.

Focused Structured Streaming sessions

Start a Structured Streaming only practice session

Every question in these sessions is drawn from the Structured Streaming domain — nothing else.

Related practice questions

Related Databricks-Spark-Assoc topic practice pages

Move into related areas when this topic feels solid.

Frequently asked questions

What does the Databricks-Spark-Assoc exam test 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.
How should I use these practice questions?
Select your answer before revealing the explanation. Then read why each option is right or wrong — this active recall approach builds retention far faster than re-reading notes.
Can I practise just Structured Streaming questions in a focused session?
Yes — the session launcher on this page draws every question from the Structured Streaming domain. Use a 10-question session first to gauge your baseline, then move to 20 or 30 once the weak spots are clear.
Where can I practise other Databricks-Spark-Assoc topics?
Use the topic links above to move to related areas, or go back to the Databricks-Spark-Assoc question bank to see all topics.
Are these real exam questions or dumps?
These are original practice questions written to test the same concepts the Databricks-Spark-Assoc exam covers. They are not copied from any real exam or dump site.