Courseiva

CCNA Structured Streaming Questions

33 questions · Structured Streaming · All types, answers revealed

1
MCQeasy

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

A.The State Store
B.Checkpointing
C.The Write Ahead Log
D.The Driver
AnswerB

Checkpointing records the metadata and offsets of the micro-batches in a persistent store. This allows the streaming query to recover from failures and resume processing exactly where it left off, which is the cornerstone of providing the exactly-once fault-tolerance guarantees expected in enterprise data pipelines.

Why this answer

Checkpointing is the mechanism that stores the query's metadata and progress in durable storage (like DBFS or S3). By recording the offset of the processed data, Spark can recover from failures and restart from the exact point where it left off. This is fundamental to ensuring that streaming applications are reliable and maintain exactly-once processing guarantees across restarts or cluster crashes.

Exam trap

Candidates often confuse checkpointing with logging. They think checkpointing is just for debugging output, whereas it is actually the critical mechanism for state recovery and exactly-once processing guarantees.

2
MCQmedium

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?

A.The checkpoint location stores the last processed offset and query metadata, enabling the query to resume exactly where it left off after a restart.
B.The checkpoint location caches the entire streamed DataFrame in memory, allowing faster recovery after a driver restart.
C.The checkpoint location stores a copy of the Kafka topic's data, enabling replay from the checkpoint instead of Kafka.
D.The checkpoint location automatically deduplicates records in the Delta table by comparing primary keys.
AnswerA

The checkpoint location persists the streaming query's progress, including offsets and state, so after a failure the query can resume from the last committed offset and avoid reprocessing or data loss. This is fundamental to Structured Streaming's exactly-once semantics when used with a replayable source and idempotent sink.

Why this answer

The checkpoint location is critical for fault tolerance in Structured Streaming. It records the progress of the query, including offsets and state, so that after a failure the query can resume from the last committed offset. This, combined with a replayable source and an idempotent sink, enables exactly-once processing.

The other options misattribute caching, deduplication, or data storage to the checkpoint.

Exam trap

The trap here is assuming that the checkpoint location stores actual data or performs deduplication, when it only stores metadata and state.

3
Multi-Selectmedium

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

Select 2 answers
A.Set the output mode to `complete`.
B.Ensure the sink is idempotent or transactional, such as Delta Lake.
C.Configure a checkpoint location for the query.
D.Use the `foreachBatch` sink to write to the Delta table.
E.Use `trigger(processingTime='0 seconds')` for continuous processing.
AnswersB, C

Exactly-once semantics require that the sink can handle replays without duplicating data. Delta Lake provides transactional writes and idempotent merges, which allow the streaming query to retry a micro-batch without creating duplicates. In this scenario, writing to a Delta table as the sink ensures that even if a batch is reprocessed after failure, the result remains consistent, thus satisfying exactly-once.

Why this answer

Exactly-once in Structured Streaming relies on two pillars: a reliable checkpoint to track progress and an idempotent or transactional sink to handle replays. A checkpoint location stores offsets and state, enabling recovery without data loss. A sink like Delta Lake ensures that re-executed batches do not duplicate output.

Together, they provide end-to-end exactly-once. Other options affect output mode, custom sink logic, or trigger frequency, none of which guarantee exactly-once.

Exam trap

The trap here is thinking that a specific output mode or trigger setting alone can guarantee exactly-once, when checkpointing and sink idempotency are the real requirements.

4
MCQeasy

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?

A.No output is produced until the query is stopped.
B.Only new rows added since the last micro-batch are written to the console.
C.Only rows that have been updated since the last micro-batch are written to the console.
D.The entire result table is rewritten to the console after every micro-batch.
AnswerB

In append mode, only new rows that have been added to the result table since the last trigger are output. For a non-aggregated streaming query, this means each new record is emitted once. This matches the typical debugging use case where you want to see incoming data as it arrives. Therefore this statement correctly describes append mode behavior.

Why this answer

Append mode in Structured Streaming outputs only new rows that are added to the result table since the last micro-batch. For a simple file source without aggregations, each incoming record is emitted once. This is ideal for debugging because you see data as it arrives.

Complete mode rewrites the whole table, and update mode emits only updated rows, neither of which matches the described behavior.

Exam trap

The trap here is mixing up append mode with complete or update mode, especially when aggregations are not involved.

5
MCQmedium

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?

A.Configure the Kafka source with `startingOffsets` set to `earliest`.
B.Enable idempotent writes by setting the Delta table property `delta.enableChangeDataFeed` to true.
C.Use `foreachBatch` to manually deduplicate records based on a unique key.
D.Set the `checkpointLocation` option to a reliable storage location.
AnswerD

The checkpoint location stores the progress information of the streaming query, including which offsets have been processed. On restart, Spark uses this to resume from where it left off, ensuring each record is processed exactly once. Combined with idempotent sinks like Delta Lake, this provides end-to-end exactly-once guarantees. Therefore, setting a checkpoint location is essential.

Why this answer

Exactly-once processing in Structured Streaming relies on checkpointing to record progress and idempotent sinks to avoid duplicates. The checkpoint location stores offset information, allowing the query to resume without reprocessing. Delta Lake supports idempotent writes when used with checkpoints.

Other options like Change Data Feed or manual deduplication do not provide the necessary fault tolerance.

Exam trap

The trap here is assuming that any Delta table property or deduplication logic automatically ensures exactly-once semantics, when actually checkpointing is the core mechanism.

6
Multi-Selecthard

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)

Select 2 answers
A.Configure a valid checkpoint location using the .option("checkpointLocation", path) method.
B.Enable Delta Lake format as the sink to leverage ACID transactions and transactional metadata.
C.Set the output mode of the streaming query to complete mode.
D.Use an idempotent or transactional sink combined with proper source offset management.
E.Increase the driver memory allocation to at least 64GB to store all Kafka offset metadata.
AnswersA, D

A checkpoint location records Kafka offsets and commit metadata durably, letting Structured Streaming resume exactly where it stopped after a restart. Without it, the query cannot recover progress, so exactly-once delivery and automatic restart recovery are impossible.

Why this answer

Achieving end-to-end exactly-once guarantees in Spark Structured Streaming requires idempotent or transactional sinks combined with persistent checkpoint directories. The checkpoint mechanism saves the exact state and offsets, allowing the streaming query to resume seamlessly after an interruption without duplicating data processing.

Exam trap

Candidates often forget that exactly-once semantics requires BOTH a transactional/idempotent sink AND a persistent checkpoint location, selecting only one of these mandatory components.

7
MCQhard

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?

A.Configure `spark.sql.streaming.checkpointLocation` to a reliable cloud storage path, use `withWatermark("eventTime", "30 minutes")`, and set the output mode to `update`.
B.Set the source option `maxFilesPerTrigger` to a high value, define a watermark of 30 minutes on the event-time column, and rely on the default checkpoint location provided by the Spark session.
C.Set a checkpoint location on a durable file system, apply `withWatermark("eventTime", "30 minutes")`, and use the `append` output mode with a Delta table sink.
D.Use `withWatermark("eventTime", "30 minutes")`, set the output mode to `complete`, and configure a checkpoint location on DBFS.
AnswerC

A durable checkpoint location preserves query progress and state across restarts, satisfying the no-reprocessing requirement. The 30-minute watermark allows the engine to accept and incorporate late events up to 30 minutes past the window boundary. The `append` output mode emits a window's final result only after the watermark passes the window end, ensuring each window is emitted exactly once to the Delta sink.

Why this answer

The pipeline needs durable state recovery and late-data tolerance without re-emitting finalized windows. A checkpoint location on reliable storage ensures the query resumes from where it left off. A 30-minute watermark defines how long late events are accepted.

The append output mode emits each window only once, after the watermark passes the window end, which matches the no-reprocessing requirement when writing to Delta.

Exam trap

The trap here is assuming that any watermark combined with any output mode will prevent reprocessing, when only append mode guarantees a window is emitted once after finalization.

8
MCQmedium

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?

A.Late-arriving data arriving after the watermark is automatically updated in-place within the existing Delta files using ACID merge operations.
B.Late-arriving data is buffered in an unlimited memory state store until the streaming query is manually stopped and restarted by an operator.
C.Late-arriving data that falls behind the specified watermark threshold is dropped and never written to the Delta table.
D.Late-arriving data forces the Delta Lake table to automatically switch its output mode to complete mode for that specific micro-batch.
AnswerC

Watermarks define how long the engine waits for late data. Any event whose event-time falls behind the current watermark is considered too late and is dropped to prevent unbounded state accumulation in streaming aggregations and joins.

Why this answer

Streaming queries using append mode require that new rows are entirely independent of previously processed outputs, meaning they are simply appended as new files. Late data arriving after the watermark threshold is dropped entirely by Spark and never written to the Delta table, preventing unbounded state growth and maintaining strict correctness guarantees.

Exam trap

Candidates often assume that late-arriving data is automatically updated in the destination table or buffered indefinitely, forgetting that watermarks strictly drop data that falls behind the threshold in append mode.

9
MCQeasy

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

A.It stores the physical data ingested by the stream.
B.It keeps track of the offsets for fault tolerance.
C.It manages the Spark UI logs for performance monitoring.
D.It serves as a cache for shuffle operations.
AnswerB

The checkpoint location records the progress of the streaming query, specifically the offsets of the data that has been processed. This mechanism enables Spark to resume from the exact point of failure, ensuring fault tolerance and preventing data loss or duplication when the streaming application is restarted after an interruption.

Why this answer

Checkpoints are the backbone of fault tolerance in Structured Streaming. They store the query's progress, including offsets and metadata, in reliable storage (like DBFS or S3). This allows the query to recover from failures and resume exactly where it left off.

Without checkpoints, any crash would cause the loss of processed offset information, forcing the system to reprocess all data from the beginning, which is inefficient.

Exam trap

Candidates often confuse checkpointing with general caching or logging. They fail to realize its primary purpose is tracking state and offsets to enable fault-tolerant recovery in streaming pipelines.

10
MCQmedium

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?

A.Use a UDF to parse the JSON string.
B.Use the from_json function with a defined schema.
C.Convert the JSON to a Map and extract keys.
D.Use the split function to break the JSON string.
AnswerB

Using from_json with a predefined schema is the standard and most performant approach in Structured Streaming. It allows Spark to parse the JSON content efficiently while enforcing schema constraints, ensuring data quality and type safety, which is essential for downstream analytical workloads and reliable streaming pipelines in production.

Why this answer

When dealing with JSON in streaming, using schema enforcement and the `from_json` function is the recommended practice. It allows you to transform the raw JSON string column into a structured format with clear types. This approach is highly efficient because it leverages Spark's Catalyst optimizer to process the nested fields, which is far superior to manual string parsing or complex regex operations.

Exam trap

Candidates often attempt to parse JSON using manual split/substring operations or custom Python UDFs. These are inefficient and ignore Spark’s built-in, optimized schema-based parsing capabilities.

11
Multi-Selectmedium

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

Select 2 answers
A.Enable `spark.sql.streaming.forceDeleteTempCheckpointLocation` to clean up temporary files.
B.Configure the Kafka source with `startingOffsets` set to `earliest`.
C.Set the `maxOffsetsPerTrigger` option to limit the number of records per batch.
D.Use a checkpoint location to store progress information.
E.Use a Delta table as the sink with idempotent writes.
AnswersD, E

A checkpoint location is essential for fault tolerance and exactly-once semantics. It stores the progress of the streaming query, including offsets processed and state information. On restart, the query resumes from where it left off, ensuring no data is lost or duplicated. Without a checkpoint, the job cannot recover correctly.

Why this answer

Exactly-once semantics in Structured Streaming require a reliable checkpoint location to track progress and an idempotent sink like Delta Lake that can handle replays without duplicating data. Checkpointing ensures the query can recover from failures, while Delta's transactional writes guarantee that each batch is applied exactly once.

Exam trap

The trap here is thinking that Kafka offset settings or rate limits provide exactly-once semantics, when actually checkpointing and an idempotent sink are the key requirements.

12
MCQmedium

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?

A.Use window(timestamp, '10 minutes') without a watermark.
B.Use window(timestamp, '10 minutes') with withWatermark('timestamp', '10 minutes').
C.Use window(timestamp, '10 minutes') with withWatermark('timestamp', '2 minutes').
D.Use a trigger interval of 2 minutes with no watermark.
AnswerC

This configuration perfectly matches the requirement by grouping events into 10-minute buckets while allowing a 2-minute buffer for late data. The watermark correctly signals to Spark that state for windows older than 2 minutes from the maximum event time can be safely cleared, optimizing memory usage and ensuring production stability.

Why this answer

Event-time windowing with watermarks is essential for handling late-arriving data in stream processing. By defining a window duration of 10 minutes and a watermark delay of 2 minutes, Spark maintains state for late events while discarding data older than the watermark threshold. This ensures the output remains accurate even when network latency or ingestion bottlenecks occur, preventing unbounded state growth in the processing engine.

Exam trap

Candidates often forget to pair the window function with a watermark, which causes Spark to throw an analysis error or accumulate unbounded state.

13
MCQeasy

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?

A.complete
B.snapshot
C.append
D.update
AnswerD

Update mode emits only the rows that were updated since the previous trigger, which is exactly the requirement. For windowed aggregations, as new events arrive and counts change, the affected window rows are re-emitted each trigger with their updated values. This provides near-real-time visibility into evolving aggregates without reprinting the entire result set, making it suitable for a debugging console sink where you want to observe changes.

Why this answer

Update mode is designed to emit only rows whose aggregate values changed since the last trigger. This matches the debugging need to see evolving window counts without reprinting the full result each time. Append withholds results until windows finalize, complete reprints everything, and snapshot is not a supported mode, so update is the correct choice.

Exam trap

The trap here is conflating update mode with complete mode: complete mode also shows changes, but it re-emits every row each trigger rather than only the rows that changed.

14
MCQhard

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?

A.The watermark ensures exactly-once processing by deduplicating records that arrive within the 15-minute window.
B.The watermark allows the engine to drop state for windows that are older than the watermark, and late data within 15 minutes is still processed and can update the aggregation.
C.The watermark forces the output mode to complete, so the entire aggregation result is recomputed on every trigger.
D.The watermark causes the aggregation state to be dropped immediately after 15 minutes, and any late data beyond that is ignored.
AnswerB

The watermark specifies how long the engine waits for late data. State for a window is kept until the watermark (max event time seen minus 15 minutes) passes the window's end time. Data arriving within the 15-minute delay is still processed and can update the aggregation. After the watermark passes, state is dropped and further late data is ignored. This is the intended behavior for handling late data.

Why this answer

A watermark in Structured Streaming defines a threshold for how late data can arrive. It allows the engine to drop old state once the watermark passes the window end time, preventing unbounded state growth. Late data within the watermark delay is still processed and can update aggregations.

This balances correctness and resource usage, enabling efficient handling of late-arriving events in streaming aggregations.

Exam trap

The trap here is thinking that a watermark immediately drops state or deduplicates data, when it actually only defines a delay threshold for late data and state cleanup based on event time.

15
MCQmedium

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?

A.An output mode of `complete`
B.A watermark on the event time column
C.A checkpoint location
D.A trigger interval of `processingTime='1 second'`
AnswerC

A checkpoint location is required for any Structured Streaming query to track progress and maintain state. It stores metadata about which offsets have been processed, as well as aggregation state. Without it, the query cannot recover from failures and would either fail to start or lose data. Configuring `checkpointLocation` ensures fault tolerance and exactly-once processing when combined with a replayable source and idempotent sink.

Why this answer

To recover from failures and ensure exactly-once processing, a Structured Streaming query must have a checkpoint location. This location stores the progress information and state, allowing the query to resume from where it left off. Watermarks, trigger intervals, and output modes are unrelated to basic fault tolerance.

Therefore, the checkpoint location is the essential configuration.

Exam trap

The trap here is thinking that a watermark or a specific trigger is required for fault tolerance, when in fact the checkpoint location is the only mandatory configuration for recovery.

16
MCQhard

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?

A.Increase the checkpoint interval.
B.Implement or tighten the watermark duration.
C.Disable checkpointing to free memory.
D.Use a larger instance type for all nodes.
AnswerB

Tightening the watermark duration instructs Spark to discard state for old data sooner. This directly reduces the memory footprint of the state store, as the engine no longer needs to track windows that have already 'expired' according to the watermark, which is the most effective way to address memory pressure.

Why this answer

Aggregations in streaming inherently require state management. If memory is exhausted, it is often due to the state store growing too large because of late data or a lack of proper watermark cleanup. Implementing a stricter watermark and periodically cleaning the state using stateful operators like 'dropDuplicates' or correctly configured windowing is necessary.

This ensures the cluster stays within its memory limits during long-term operation.

Exam trap

Candidates often try to increase executor memory or change the shuffle partition count, ignoring the root cause: the state store is growing indefinitely due to missing or loose watermarks.

17
Multi-Selectmedium

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

Select 2 answers
A.Configure a unique checkpoint location for the query.
B.Enable the option latestFirst on the file source.
C.Set the query to use trigger(once=True) so it stops after each run.
D.Set the file source option maxFilesPerTrigger to a high value to process all files at once.
E.Use a Delta table sink so the write is idempotent and transactional.
AnswersA, E

The checkpoint location stores progress information, including which files or offsets have been processed and the state of stateful operations. On restart, the engine reads this metadata to resume exactly where it left off, preventing reprocessing of already-handled files. Without a checkpoint, the query cannot recover its position and will either fail to restart or reprocess from the beginning, so this is essential for resilient incremental processing.

Why this answer

Resilient incremental file ingestion requires two things: durable progress tracking and a sink that can commit batches atomically. A unique checkpoint location records which files and offsets have been processed, and a Delta sink provides transactional, idempotent writes that prevent duplicate effects when batches are retried. Throughput and ordering options do not affect recovery correctness.

Exam trap

The trap here is focusing on performance-oriented file source options such as maxFilesPerTrigger or latestFirst, which tune behavior but do not persist progress or guarantee idempotent writes.

18
MCQhard

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?

A.The record is held in state until the next watermark update and then emitted.
B.The record causes the query to fail with a watermark violation error.
C.The record is dropped because its event time is older than the watermark.
D.The record is processed and appended to the Delta table immediately.
AnswerC

In append mode, records with event times older than the current watermark are considered too late and are dropped. The watermark is at 10:10, and the record's event time is 10:00, which is 10 minutes behind, exceeding the 5-minute delay. Therefore the record is discarded and will not appear in the output. This is the correct behavior for this scenario.

Why this answer

With a 5-minute watermark, any record whose event time is more than 5 minutes behind the current watermark is considered late. In append mode, such records are dropped because the system assumes all data for that time has already been emitted. The record at 10:00 arrives when the watermark is 10:10, so it is 10 minutes late and is discarded.

This matches the expected watermark behavior.

Exam trap

The trap here is assuming that late data is buffered or triggers an error, rather than being silently dropped in append mode once the watermark has passed.

19
MCQhard

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

A.Increasing the number of partitions.
B.Reducing the trigger interval.
C.Implementing state timeouts.
D.Switching to Append mode.
AnswerC

State timeouts enable the engine to remove inactive state entries after a defined period. By using processing or event-time timeouts, you ensure that memory is reclaimed for keys that are no longer active, which is the most effective way to prevent OutOfMemory errors in stateful streaming.

Why this answer

The 'mapGroupsWithState' function maintains state for keys until they are explicitly timed out or removed. If keys are not removed, the state grows indefinitely, leading to memory issues. Using 'GroupStateTimeout' (either processing or event-time) allows the engine to automatically expire state after a specific duration, effectively bounding memory usage.

This is a critical pattern for managing state in long-running streaming applications on Databricks.

Exam trap

Candidates often assume Spark automatically cleans up intermediate streaming state or rely solely on increasing cluster memory, missing that mapGroupsWithState requires explicit timeouts to purge stale keys and prevent infinite memory growth.

20
MCQeasy

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

A.Continuous processing supports all SQL operations.
B.Continuous processing offers lower latency.
C.Micro-batch processing provides lower latency.
D.Continuous processing is the default mode.
AnswerB

Continuous processing is specifically engineered to provide sub-millisecond latency by avoiding the batch scheduling overhead present in micro-batch processing. By running tasks continuously on the executors, it eliminates the start-stop cycle, making it ideal for extremely latency-sensitive applications that require instantaneous response times to streaming data inputs.

Why this answer

Micro-batch processing handles data in small discrete batches, providing high throughput and fault tolerance with second-level latency. Continuous processing, by contrast, runs a task continuously on each executor, enabling sub-millisecond latency. Choosing between them involves a trade-off between strict latency requirements and the operational complexity or feature limitations inherent in the continuous processing model, which does not support all operations yet.

Exam trap

Candidates often confuse the terminology, wrongly assuming that continuous processing is the default mode or that it provides higher throughput rather than just lower latency at higher resource costs.

21
MCQeasy

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?

A.Query the `spark.streaming` log file on the driver node.
B.Run `spark.sql("SHOW STREAMING QUERIES")`.
C.Check the Delta table's transaction log for streaming metrics.
D.Use the `StreamingQuery.lastProgress` method on the query object.
AnswerD

`StreamingQuery.lastProgress` returns a JSON object containing detailed metrics for the most recent micro-batch, including input rows per second, processing rate, batch duration, and more. It is a direct API call on the streaming query object and provides the most immediate and structured access to the query's status. This method is commonly used to programmatically monitor streaming queries.

Why this answer

The `StreamingQuery.lastProgress` method provides a structured snapshot of the most recent micro-batch's metrics, including input rate and batch duration. It is the direct API for monitoring. Log files are not structured, there is no SQL command for showing streaming queries, and Delta logs do not contain streaming metrics.

Thus, `lastProgress` is the correct choice.

Exam trap

The trap here is assuming that there is a SQL command to show streaming queries or that Delta logs contain runtime metrics, when in fact these are accessed via the streaming query API.

22
MCQeasy

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?

A.Update with watermark
B.Update
C.Append
D.Complete
AnswerC

Append mode outputs only new rows that have been added since the last trigger. Since there are no aggregations, each micro-batch contains only the new data from Kafka. This matches the requirement to see only newly arrived data. Therefore, Append mode is correct.

Why this answer

In Structured Streaming, Append mode is the default and outputs only new rows added since the last trigger. For a stateless query like reading from Kafka and writing to console, each micro-batch contains only new data. Update and Complete modes are designed for aggregations, where they output updated or full results, respectively.

Exam trap

The trap here is confusing Update mode with Append mode, especially when no aggregations are involved, leading to unnecessary complexity.

23
MCQeasy

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?

A.`.option("ignoreChanges", "false")`
B.No option is needed; the Delta source reads the initial snapshot by default on the first run when no checkpoint exists.
C.`.option("includeExistingData", "true")`
D.`.option("startingVersion", "0")`
AnswerB

When a Structured Streaming query reads from a Delta table and starts without a checkpoint, the Delta source processes the entire current snapshot of the table as the first micro-batch, then follows with new changes. This is the default behavior, so the developer does not need to set any special option to read existing data on the first run.

Why this answer

The Delta Lake streaming source is designed to read the full table snapshot as the first micro-batch when a query starts without a checkpoint, and then to continue with subsequent changes. No option is required to include existing data; the behavior is the default. Options like `startingVersion` and `ignoreChanges` alter other aspects of the read, not the initial snapshot inclusion.

Exam trap

The trap here is assuming an option like `includeExistingData` exists, when the Delta source already reads existing data by default on the first run.

24
MCQmedium

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?

A.The late events are buffered and processed once the watermark advances.
B.The late events are dropped and not included in the output.
C.The late events cause the query to fail with a watermark violation error.
D.The late events are written to a separate dead-letter queue automatically.
AnswerB

In Structured Streaming, when a watermark is defined on an event-time column, any event whose event-time timestamp is older than the current watermark (max event time seen minus the watermark delay) is considered too late and is dropped from the aggregation. This is the documented behavior: watermarking allows the engine to discard late data to bound state. Since the job uses append mode with a watermark, late events beyond the threshold are excluded from the result set, which is the correct outcome here.

Why this answer

When a watermark is set on an event-time column, Structured Streaming uses it to determine when data is too late to be included in an aggregation. Events with timestamps older than the watermark are dropped. This behavior is intentional to bound state and ensure timely results.

The other options describe mechanisms that are not part of the watermarking feature, such as buffering, erroring, or automatic dead-lettering.

Exam trap

The trap here is assuming that late data is either reprocessed or causes an error, rather than being silently dropped as the watermark dictates.

25
MCQhard

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?

A.Set the checkpoint location to a new path and restart the query so the old state is discarded.
B.Increase the spark.sql.streaming.stateStore.maintenanceInterval to force more frequent cleanup of old keys.
C.Add withWatermark on the event-time column and use dropDuplicatesWithinWatermark on the key columns.
D.Repartition the stream by the composite key before calling dropDuplicates to spread state across more partitions.
AnswerC

dropDuplicatesWithinWatermark combined with withWatermark limits how long the engine retains keys for deduplication. Once the watermark passes a key's event time beyond the configured delay, the key can be evicted from state, bounding memory usage. This directly addresses unbounded state growth while still deduplicating events that arrive within the allowed lateness window, which matches the requirement to remove duplicates within a reasonable period.

Why this answer

Bounding state requires a rule that tells the engine when a key will no longer be needed. A watermark on event time provides that rule, and dropDuplicatesWithinWatermark uses it to evict keys older than the watermark. This keeps deduplication correct within the allowed lateness while preventing state from growing without limit.

Configuration tweaks and repartitioning do not change the fundamental retention policy.

Exam trap

The trap here is believing that tuning state store internals or partitioning can bound state, when only a watermark-based retention rule lets the engine decide a key can be safely forgotten.

26
MCQhard

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?

A.The rate source does not support `outputMode("append")` and therefore emits no rows.
B.The Delta sink requires `checkpointLocation` to be set, and without it the query silently discards all output.
C.The Delta table was created without specifying a schema, causing all writes to be rejected.
D.The query includes a windowed aggregation with a watermark, and in append mode results are emitted only after the watermark passes the window end, so no rows are emitted until then.
AnswerD

In append mode, windowed aggregations emit a window's result only once the watermark has advanced past the window's end time, guaranteeing no further updates. If the watermark is large or the data stream is short, the watermark may not have passed any window end yet, so no rows are written. This explains why the query processes data but the Delta table remains empty.

Why this answer

With append output mode and a windowed aggregation, Structured Streaming waits until the watermark passes the end of a window before emitting that window's results. If the watermark has not advanced sufficiently, no rows are emitted, so the target Delta table stays empty even though the query is running and processing data. This is expected behavior, not a failure.

Exam trap

The trap here is assuming append mode always writes rows immediately, when for windowed aggregations it delays emission until the watermark finalizes each window.

27
MCQhard

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?

A.Set the Kafka source option `failOnDataLoss` to true.
B.Configure the external database to use a transaction isolation level of SERIALIZABLE.
C.Enable Kafka's idempotent producer by setting `enable.idempotence=true`.
D.Use the batch ID from `foreachBatch` to make writes to the external database idempotent.
AnswerD

To achieve exactly-once semantics when writing to an external database via `foreachBatch`, the writes must be idempotent. The batch ID provided by `foreachBatch` is unique for each micro-batch and can be used to deduplicate writes. For example, you can store the batch ID in the database and check it before applying the batch, or use it as part of an upsert key. Combined with checkpointing, which ensures that the stream resumes from the correct batch after a failure, this provides exactly-once semantics. Without idempotent writes, a failure could cause a batch to be reprocessed, leading to duplicates.

Why this answer

Exactly-once semantics in Structured Streaming with an external sink like a database require both checkpointing and idempotent writes. The batch ID from `foreachBatch` enables idempotency by allowing you to deduplicate writes. Checkpointing ensures the stream can recover from failures without reprocessing completed batches.

The other options do not address the need for idempotent writes to the external system.

Exam trap

The trap here is assuming that Kafka producer settings or database isolation levels alone can guarantee exactly-once, when the critical piece is idempotent writes using the batch ID.

28
MCQmedium

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?

A.Kafka consumers always begin reading from the committed consumer-group offset, and Structured Streaming reuses that group offset as its start position.
B.The console sink silently discards any record whose Kafka `offset` value is lower than the current micro-batch ID.
C.The topic's retention policy removed all historical segments before the query started, so only new records remained available to read.
D.The Kafka source defaults to `startingOffsets` = "latest", so the first micro-batch begins at the tail of each partition and earlier records are skipped.
AnswerD

The Kafka Structured Streaming source uses `startingOffsets` to decide where the very first query starts when no checkpoint exists. Its default value is "latest", which means the query begins at the newest offset per partition. Older records already present in the topic are therefore never delivered to the first micro-batch, exactly matching the observed behavior.

Why this answer

When a Kafka-backed streaming query starts without an existing checkpoint, the source resolves its first offsets from `startingOffsets`. Because the default is "latest", the initial micro-batch aligns to the end of each partition, so any records produced before the query started are ignored. Setting `startingOffsets` to "earliest" or a JSON offset map is required to consume the backlog.

Exam trap

The trap here is assuming the Kafka source always reads the full topic backlog by default, when in fact its default start position is the latest offset.

29
MCQhard

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?

A.State is maintained only for the right stream, as it is the stream being joined to.
B.State is not maintained; the join is stateless and relies on the current micro-batch.
C.State is maintained only for the left stream, as it is the driving stream.
D.State is maintained for both streams and is cleaned up based on the watermarks.
AnswerD

In a stream-stream join, the engine maintains state for both sides to match records. Watermarks are used to determine when state can be evicted. For an inner join, state for a record is kept until the watermark passes the record's event time plus the watermark delay, ensuring that late matches can still occur within the allowed lateness.

Why this answer

In a stream-stream inner join, the engine maintains state for both streams to match records across micro-batches. Watermarks define when state can be evicted: once the watermark passes the event time of a record plus the allowed lateness, its state is removed. This ensures that late-arriving matches within the watermark window are still processed.

Exam trap

The trap here is assuming that one stream is the driver and only its state is kept, but stream-stream joins require state on both sides.

30
MCQeasy

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?

A.Default trigger, which processes data in micro-batches as soon as the previous batch completes, without a fixed interval.
B.Continuous trigger with a 1-second checkpoint interval, which provides low-latency processing.
C.Once trigger, which processes all available data in a single batch and then stops.
D.ProcessingTime trigger with an interval of 0 seconds, which processes data as fast as possible.
AnswerA

When no trigger is specified, Structured Streaming uses the default trigger, which processes each micro-batch as soon as the previous one finishes. This provides the lowest latency without a fixed interval. It is the standard behavior for streaming queries that need continuous processing. The developer does not need to set a trigger for this default behavior.

Why this answer

The default trigger in Structured Streaming processes data in micro-batches as soon as the previous batch completes, without a fixed interval. This is the behavior when no trigger is explicitly set. It provides continuous processing with low latency, making it suitable for most streaming use cases.

Developers can override this by specifying a processing time, once, or continuous trigger.

Exam trap

The trap here is confusing the default trigger with a processing time trigger of 0 seconds or a once trigger, when the default is actually an unspecified trigger that runs micro-batches back-to-back.

31
Multi-Selectmedium

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

Select 3 answers
A.Delta Lake tables
B.Apache Kafka
C.File system directories (e.g., Parquet)
D.Standard SQL Databases (JDBC)
E.In-memory lists
AnswersA, B, C

Delta Lake is a first-class streaming source in Spark. It uses the transaction log to track data changes, allowing it to efficiently read new rows as they are added to the table. This is the recommended approach for building robust streaming pipelines on the Databricks Lakehouse platform.

Why this answer

Structured Streaming requires data sources to provide a way to track offsets. Delta Lake, Kafka, and File-based sources (like Parquet or JSON) are natively supported because they provide the necessary metadata for Spark to determine what data has been read and what is new. Understanding these sources is essential for designing robust streaming architectures that integrate seamlessly with the Databricks ecosystem.

Exam trap

Test-takers frequently select traditional batch-only formats or generic SQL sources that lack native transaction logs or offset tracking capabilities, incorrectly assuming every Spark SQL source supports streaming reads.

32
MCQhard

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?

A.Use the batchId argument to make the Delta write idempotent, for example by merging on a batch identifier or using replaceWhere with the batchId.
B.Increase the number of shuffle partitions so each batch writes to more files.
C.Call df.cache() on the batch DataFrame before writing to Delta.
D.Set the Delta table property delta.appendOnly to true.
AnswerA

foreachBatch provides the batch identifier, which is monotonically increasing and stable across retries of the same batch. By using it to make the write idempotent, such as merging on a batch column or replacing a partition keyed by batchId, a retried batch overwrites rather than duplicates its prior output. This delivers exactly-once effects at the sink even when the batch function is re-executed.

Why this answer

Exactly-once sink semantics with foreachBatch require the batch function itself to be idempotent, because the function may be invoked more than once for the same batch. The provided batchId is stable across retries, so using it as a deduplication key or partition selector lets a retried write overwrite its previous output. Caching, append-only properties, and partition tuning do not provide this guarantee.

Exam trap

The trap here is assuming that Structured Streaming automatically guarantees exactly-once for arbitrary foreachBatch code, when in fact the custom function must be made idempotent by the developer.

33
MCQhard

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?

A.The Kafka source is not configured with `startingOffsets` set to `latest`.
B.The trigger interval is too short, causing the query to start new micro-batches before the previous one finishes.
C.The `maxOffsetsPerTrigger` value is lower than the rate at which new data is arriving.
D.The `maxOffsetsPerTrigger` option is ignored when writing to Delta Lake.
AnswerC

If `maxOffsetsPerTrigger` is set to 10000 and the Kafka topic receives more than 10000 new records per trigger interval, the query will only process 10000 records per micro-batch. The backlog will grow because the processing rate is capped below the arrival rate. To reduce the backlog, you would need to increase `maxOffsetsPerTrigger` or increase the trigger frequency, provided the cluster can handle the load.

Why this answer

The backlog is not decreasing because the query is limited to processing 10000 records per micro-batch, which is likely less than the incoming rate. This cap is set by `maxOffsetsPerTrigger`. To reduce the backlog, you must increase this limit or the processing capacity.

The other options do not address the rate mismatch: the option is not ignored, triggers do not overlap, and starting offsets only affect initial reads.

Exam trap

The trap here is assuming that `maxOffsetsPerTrigger` is a soft limit or that it only applies to the initial load, when in fact it strictly caps the number of records processed per micro-batch.

Ready to test yourself?

Try a timed practice session using only Structured Streaming questions.