Courseiva

Databricks Certified Associate Developer for Apache Spark (Databricks-Spark-Assoc) — Questions 226–295

295 questions total · 4pages · All types, answers revealed

Page 3

Page 4 of 4

226
MCQmedium

A data engineer has a Pandas-on-Spark DataFrame `psdf` with a column `event_time` stored as string. They run `psdf['event_time'] = pd.to_datetime(psdf['event_time'])` where `pd` is the Pandas API on Spark module. What is the most likely outcome?

A.The operation triggers immediate collection of all data to the driver to perform the conversion locally.
B.The operation succeeds and returns a new Pandas-on-Spark Series with datetime64[ns] dtype, executed lazily.
C.The operation raises a TypeError because Pandas-on-Spark does not support datetime conversion on string columns.
D.The operation succeeds only if the DataFrame has a single partition, otherwise it fails with an AnalysisException.
AnswerB

Pandas API on Spark implements `to_datetime` and returns a Series backed by Spark. Assignment to an existing column updates the DataFrame lazily; the conversion is applied per partition when an action triggers computation, and the resulting dtype is datetime64[ns] as exposed by the pandas-compatible API.

Why this answer

The Pandas API on Spark provides `to_datetime` that operates in a distributed manner, returning a Series with datetime64[ns] dtype. Assigning it back to a column updates the DataFrame lazily, and no driver collection occurs. This aligns with the goal of scaling pandas-like code on Spark without changing semantics.

Exam trap

The trap here is assuming that pandas API on Spark operations like `to_datetime` force local execution or fail on distributed data, when in fact they are implemented as Spark transformations.

227
Multi-Selectmedium

A Spark application is running on a cluster with 5 executors. The driver program creates a broadcast variable that is used in a transformation. Which two components are directly involved in distributing and using the broadcast variable? (Choose two.)

Select 2 answers
A.The DAG Scheduler ensures that broadcast variables are only used in narrow transformations.
B.The Cluster Manager replicates the broadcast variable across all nodes in the cluster.
C.The driver serializes the broadcast variable and sends it to each executor via a BitTorrent-like protocol.
D.The Task Scheduler assigns the broadcast variable to each task individually.
E.Each executor caches the broadcast variable in memory and makes it available to all tasks within that executor.
AnswersC, E

The driver is responsible for serializing the broadcast variable and initiating its distribution. Spark uses an efficient broadcast mechanism, often BitTorrent-like, to disseminate the variable to executors without overwhelming the driver. This ensures that each executor receives the variable once and can share it with other executors if needed.

Why this answer

Broadcast variables are distributed by the driver, which serializes and sends them to executors using an efficient broadcast protocol. Each executor then caches the variable and makes it available to all its tasks. This two-step process minimizes network traffic and ensures that the variable is easily accessible during task execution.

Exam trap

The trap here is assuming that the Cluster Manager or Task Scheduler plays a role in broadcast variable distribution, when in fact it is handled entirely by the driver and executors.

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

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

230
MCQmedium

A data engineer is building a feature pipeline and needs to add a monotonically increasing integer column `row_num` to a DataFrame `df` that assigns consecutive numbers to rows within each partition of a specified ordering, similar to a window function. Which approach uses the DataFrame API to compute this value?

A.df.withColumn("row_num", monotonically_increasing_id())
B.df.withColumn("row_num", row_number().over(Window.orderBy("some_col")))
C.df.withColumn("row_num", lit(1).cast("int"))
D.df.rdd.zipWithIndex().toDF()
AnswerB

row_number() is a window function that assigns consecutive integers starting at one according to the window's ordering, which matches the requested behavior. Used with over(Window.orderBy(...)), it produces the running sequence the engineer needs. Note that a global orderBy in the window moves all rows to a single partition, so the engineer should be aware of that cost, but the semantics are exactly correct.

Why this answer

Consecutive numbering based on a defined ordering is the job of the row_number() window function applied with over(Window.orderBy(...)). It yields one-based consecutive integers according to the window specification, which is what the feature pipeline needs. monotonically_increasing_id() leaves gaps and ignores ordering, a literal constant does not increment, and the RDD zipWithIndex approach collapses to one partition and bypasses DataFrame optimizations.

Exam trap

The trap here is treating monotonically_increasing_id() as a drop-in row counter, when its values are only increasing and unique rather than consecutive, and it cannot honor an ordering clause.

231
MCQmedium

Refer to the exhibit. Traceback (most recent call last): File "app.py", line 12, in <module> df = spark.read.table("default.sales") File "/opt/spark/python/pyspark/sql/session.py", line 314, in table return DataFrame(self._client.execute_plan(parser.parse_table(name)))) File "/opt/spark/python/pyspark/sql/connect/client/core.py", line 112, in execute_plan(y+"sessionID"), grpc.RpcError: StatusCode.UNAVAILABLE An engineer attempts to run a PySpark script using Spark Connect but encounters the traceback shown above. What is the most likely root cause of this execution failure?

A.The target Delta table default.sales contains corrupted parquet files that cause schema resolution errors on the remote server.
B.The Spark Connect server is unreachable, turned off, or listening on a different network port than specified in the client connection string.
C.The PySpark client library version installed locally is newer than the server-side Spark runtime version, causing protocol serialization mismatches.
D.The user running app.py lacks Hive metastore permissions to read the default database catalog on the Databricks workspace.
AnswerB

A StatusCode.UNAVAILABLE status code is raised by the gRPC client library when it fails to connect to the remote server endpoint. This confirms a network connectivity barrier, incorrect host specification, or an inactive Spark Connect background service on the cluster.

Why this answer

The StatusCode.UNAVAILABLE gRPC error indicates that the client application cannot establish or maintain a network connection with the Spark Connect server endpoint. This typically happens when the cluster is stopped, the port is blocked by a firewall, or the connection string URL is incorrect. Verifying network accessibility and cluster status is an essential first troubleshooting step for Spark Connect deployments.

Exam trap

Test-takers often assume the traceback points to a syntax error or a missing database table, ignoring the gRPC status code indicating a network connectivity or server availability issue.

232
MCQmedium

A developer is using the Pandas API on Spark to process a large dataset. They need to apply a custom Python function to each value in a column. They consider using `psdf['col'].apply(custom_func)`. What should they be aware of regarding performance?

A.`apply` with a Python function is executed row-by-row in Python, which can be slow and may cause out-of-memory errors if the data is not partitioned well.
B.`apply` with a Python function is not supported in Pandas API on Spark and will raise a NotImplementedError.
C.`apply` with a Python function is automatically parallelized across all cores of the driver node, ensuring optimal performance.
D.`apply` with a Python function is executed in a vectorized manner using Arrow, so it is as fast as built-in Pandas functions.
AnswerA

This is correct because Pandas API on Spark's `apply` method applies the function to each element individually, which involves Python serialization and overhead. This can be significantly slower than using vectorized operations or Spark SQL functions. Additionally, if partitions are large, the row-by-row processing can lead to memory issues, as the entire partition is loaded into memory for the operation.

Why this answer

Using a Python function with `apply` in Pandas API on Spark processes data row-by-row in Python, which is slower than vectorized operations and can lead to memory issues if partitions are large. It is important to use built-in Pandas functions or Spark SQL functions when possible for better performance.

Exam trap

The trap here is assuming that Pandas API on Spark automatically optimizes arbitrary Python functions with Arrow for vectorized execution, when in fact such functions are executed row-by-row in Python.

233
MCQhard

You are working with a Pandas-on-Spark DataFrame `psdf` that has a default index generated by Spark. You need to perform a join with another Pandas-on-Spark DataFrame `other` that also has a default index. After the join, you notice that the resulting DataFrame has a new index and the original indices are lost. Which of the following best explains this behavior?

A.The default index is lost because the join operation uses the index as the join key by default, and since both DataFrames have the same default index, it causes a conflict.
B.The default index is not preserved across joins because it is not a true pandas index; it is a synthetic index generated per partition, and joins require shuffling which discards the original index.
C.The default index is preserved across joins only if you set the Spark configuration `spark.pandas.join.index` to true.
D.The default index is lost because the join operation converts both DataFrames to Spark DataFrames, performs the join, and then converts back to Pandas-on-Spark, which resets the index.
AnswerB

This is correct. In Pandas API on Spark, when no explicit index is set, a default index is created using `distributed-sequence` or `distributed` index types. These indices are not stable across operations that require shuffling, such as joins. The join operation shuffles data based on join keys, and the default index is not carried over, resulting in a new default index for the output.

Why this answer

The default index in Pandas API on Spark is a synthetic index that is not preserved across operations that shuffle data, such as joins. When you perform a join, the data is redistributed across partitions, and the original default index is discarded. The resulting DataFrame gets a new default index.

To maintain a stable index, you should set an explicit index using `set_index` before the join.

Exam trap

The trap here is assuming that the default index behaves like a pandas index and is preserved across joins, but it is not stable across shuffles.

234
MCQmedium

A developer runs a PySpark job on Databricks that reads a large Delta table, filters on a timestamp column, and writes results to another Delta table. The job takes 45 minutes, but the Spark UI shows that 90% of task time is spent reading from the source table. The developer wants to reduce the read time. Which action should the developer take?

A.Increase the number of shuffle partitions by setting spark.sql.shuffle.partitions to a higher value.
B.Enable Delta Lake data skipping by ensuring the timestamp column is in the table's partitioning or Z-ORDER BY columns.
C.Cache the source DataFrame in memory using .cache() before applying the filter.
D.Repartition the source DataFrame by the timestamp column before filtering.
AnswerB

Delta Lake data skipping uses file-level statistics to skip reading files that do not contain relevant data. If the timestamp column is a partition column or has been optimized with Z-ORDER BY, the query engine can prune files based on the filter predicate, drastically reducing I/O and read time. This directly addresses the observed bottleneck.

Why this answer

The Spark UI indicates that the bottleneck is reading from the source Delta table. Delta Lake data skipping leverages file statistics to avoid reading irrelevant files when a filter predicate is applied. Ensuring the timestamp column is a partition column or has Z-ORDER BY applied allows the engine to prune files effectively, reducing I/O and overall job time.

Exam trap

The trap here is assuming that caching or repartitioning will speed up a slow read, when the real solution is to reduce the amount of data read via data skipping.

235
MCQmedium

A data engineer submits a Spark application using spark-submit in client deploy mode from an edge node. The application reads a large Parquet dataset, performs a groupBy aggregation, and writes the result to a Delta table. The engineer notices that the Driver process runs on the edge node and remains alive throughout the application's lifetime. Which statement best describes the role of the Driver in this scenario?

A.The Driver executes the actual data processing tasks and stores intermediate shuffle data on local disk.
B.The Driver is responsible for storing the final output data and serving it to downstream consumers.
C.The Driver schedules tasks, maintains the DAG, and coordinates with the cluster manager to allocate executors.
D.The Driver acts as a passive monitor that only collects metrics and logs, while the cluster manager handles all scheduling.
AnswerC

In client deploy mode, the Driver runs on the submitting machine (edge node) and is responsible for converting the user program into a DAG, splitting it into stages, scheduling tasks on executors, and negotiating resources with the cluster manager. It also tracks task status and aggregates results. This matches the scenario where the Driver remains alive on the edge node.

Why this answer

The Driver in client deploy mode runs on the submitting host and orchestrates the application: it builds the DAG, schedules stages and tasks, and communicates with the cluster manager to acquire executors. It does not execute data processing tasks or store data. Therefore, the statement that it schedules tasks and coordinates resource allocation is correct.

Exam trap

The trap here is assuming that because the Driver runs on the edge node, it also performs data processing or storage, when in fact it only coordinates.

236
MCQmedium

When working with Delta Lake tables in Databricks, which command should you use to optimize the physical layout of files to improve query performance?

A.COMPACT TABLE table_name
B.OPTIMIZE table_name
C.REORGANIZE TABLE table_name
D.VACUUM table_name
AnswerB

The OPTIMIZE command is the standard Delta Lake operation used to coalesce small files into larger, more performant files. It is an essential maintenance task for Databricks environments to ensure that storage layouts remain optimized for analytical queries, which significantly reduces the time spent on I/O operations.

Why this answer

Delta Lake provides the `OPTIMIZE` command to compact small files into larger ones, which is vital for maintaining performance as data grows. Frequent small writes can lead to file proliferation, degrading read speeds. Running `OPTIMIZE` regularly helps maintain efficient file sizes, enabling faster query execution by reducing metadata overhead and maximizing the benefits of data skipping through better file-level statistics.

Exam trap

Candidates often select 'VACUUM' or 'COMPACT' instead of 'OPTIMIZE'. They confuse the command for removing old files (VACUUM) with the command for improving query performance by compacting small files.

237
MCQmedium

A developer notices a Spark job is failing with an OutOfMemoryError during a join operation on two large tables. The join key is highly skewed, causing one task to process significantly more data than others. Which technique should be applied to resolve this skew?

A.Increase the spark.driver.memory configuration.
B.Implement salting on the join key to distribute the skewed keys.
C.Enable AQE and increase the spark.sql.shuffle.partitions value.
D.Convert the larger table into a broadcast variable.
AnswerB

Salting distributes rows with the same join key across multiple partitions by appending a random integer. This ensures that the massive volume of data associated with a single key is processed in parallel by different executors, effectively eliminating the bottleneck that triggers OutOfMemoryError in skewed join scenarios.

Why this answer

Salting involves adding a random prefix or suffix to the join key to redistribute the data across multiple partitions. This prevents a single executor from handling a disproportionate amount of data. This approach is a standard industry pattern for resolving skew-related OOM errors, as it breaks the hot key into smaller, manageable chunks that can be processed in parallel across the cluster without bottlenecking at a single node.

Exam trap

Candidates frequently suggest repartitioning by the skewed key itself. This is ineffective because it simply moves the same skewed data into a single partition, failing to alleviate the bottleneck on that specific task.

238
MCQeasy

You are writing a Databricks notebook and want to use the Pandas API on Spark. Which import statement should you use to access the Pandas API on Spark?

A.from databricks import pandas as ps
B.import pyspark.pandas as ps
C.from pyspark.sql import pandas as ps
D.import pandas as ps
AnswerB

This is correct. The Pandas API on Spark is available in the `pyspark.pandas` module. Importing it as `ps` is a common convention. This provides access to functions like `ps.DataFrame`, `ps.read_csv`, etc., which mimic the pandas API but operate on Spark DataFrames for scalability.

Why this answer

To use the Pandas API on Spark, you must import it from `pyspark.pandas`. This module provides a pandas-like interface on top of Spark, allowing you to scale your pandas code. The conventional alias is `ps`.

Other imports either refer to the standard pandas library or non-existent modules.

Exam trap

The trap here is confusing the standard pandas library with the Pandas API on Spark, which requires a different import.

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

240
MCQmedium

Your Spark application is experiencing severe data skew while performing a join between a large fact table and a small dimension table. Which technique should you apply to optimize performance?

A.Increase the number of partitions using repartition() on the large table.
B.Enable AQE and increase the spark.sql.shuffle.partitions configuration.
C.Force a broadcast join for the small dimension table.
D.Apply a salting technique to the join keys of the small table.
AnswerC

Broadcasting the smaller DataFrame eliminates the need for a sort-merge join, which is where skewed data causes bottlenecks. By copying the small table to every executor, you perform a map-side join, effectively avoiding the shuffle of the large table and preventing skewed keys from overloading specific worker nodes.

Why this answer

Broadcast hash joins are the most effective solution for skew caused by a large table join when one side is small enough to fit in memory. By broadcasting the small table to all executor nodes, Spark avoids the expensive shuffle operation that causes data skew. This prevents specific partitions from becoming hotspots, which is critical for maintaining stable and performant ETL pipelines in production Databricks environments.

Exam trap

Candidates often suggest salting or repartitioning as the first step, ignoring that broadcasting is the most efficient and direct way to handle skew when a small table is involved.

241
MCQeasy

A data engineer has a PySpark DataFrame `orders` with a string column `order_ts` formatted as `yyyy-MM-dd HH:mm:ss`. They need a new column `order_date` containing only the date portion as a `date` type, so downstream code can filter by day. Which approach is correct?

A.`orders.withColumn('order_date', from_unixtime(unix_timestamp('order_ts', 'yyyy-MM-dd HH:mm:ss'), 'yyyy-MM-dd'))`
B.`orders.withColumn('order_date', date_trunc('day', 'order_ts'))`
C.`orders.withColumn('order_date', substring('order_ts', 1, 10))`
D.`orders.withColumn('order_date', to_date('order_ts', 'yyyy-MM-dd HH:mm:ss'))`
AnswerD

This is correct because `to_date` parses the string using the supplied pattern and returns a `date` column, which is exactly what downstream day-level filtering needs. It preserves the original string column and adds a properly typed date, avoiding implicit conversions or extra casts.

Why this answer

Converting a formatted timestamp string to a `date` requires an explicit parse with a pattern. `to_date(col, pattern)` parses the string and returns a `date` column, which is the correct type for day-level filters and joins. String manipulation or `from_unixtime` yields strings, and `date_trunc` yields a timestamp, so none of them match the requirement as cleanly or as safely as `to_date`.

Exam trap

The trap here is treating any function that visually produces '2024-01-01' as equivalent, ignoring that column type drives downstream comparison behavior.

242
MCQmedium

A developer is building a Spark Connect client application in Python that runs on a local workstation and connects to a remote Databricks cluster. The application must construct a DataFrame from a list of Python dictionaries without requiring the data to be uploaded to cloud storage first. Which approach should the developer use?

A.Write the list of dictionaries to a local JSON file, then call spark.read.json() with the local file path.
B.Broadcast the list of dictionaries with spark.sparkContext.broadcast() and then call spark.createDataFrame() on the broadcast variable.
C.Use spark.createDataFrame(list_of_dicts) directly, because Spark Connect serializes local Python collections into Arrow batches and sends them to the server.
D.Use spark.sparkContext.parallelize(list_of_dicts) and convert the resulting RDD to a DataFrame with spark.createDataFrame().
AnswerC

Spark Connect supports creating DataFrames from local Python collections. The client serializes the list of dictionaries into Apache Arrow record batches and transmits them over gRPC to the server, which reconstructs the DataFrame. No intermediate cloud storage or manual upload step is required, making this the correct approach for the described scenario.

Why this answer

Spark Connect allows client-side Python collections to be converted into DataFrames through the standard createDataFrame API. The client serializes the data into Arrow format and streams it to the server over gRPC, so the data need not be staged in cloud storage. Approaches relying on SparkContext, RDD parallelization, or broadcasting fail because Spark Connect deliberately omits the low-level RDD and context APIs.

Exam trap

The trap here is assuming Spark Connect requires data to be staged in remote storage before a DataFrame can be built, when local Python collections are actually serialized and sent directly.

243
MCQhard

A data engineer has a DataFrame `events` with columns `user_id` and `ts`. They need to add a column `prev_ts` that holds the previous event timestamp for each user, ordered by `ts` ascending, without collapsing rows. Which operation accomplishes this?

A.events.groupBy('user_id').agg(max('ts').alias('prev_ts'))
B.Window.partitionBy('user_id').orderBy('ts') with lag('ts', 1) via withColumn
C.events.sortWithinPartitions('user_id', 'ts').withColumn('prev_ts', lead('ts'))
D.events.dropDuplicates(['user_id']).withColumnRenamed('ts', 'prev_ts')
AnswerB

A window partitioned by `user_id` and ordered by `ts`, combined with the `lag` function, returns the previous row's timestamp while preserving every original row. This is exactly the pattern for adding a prior-value column without collapsing data, and it executes as a distributed window operation under Catalyst.

Why this answer

Using a window partitioned by `user_id` and ordered by `ts` with the `lag` function returns the previous event's timestamp for each row while retaining all rows. The other choices either aggregate away detail, rely on partition-local sorting that cannot guarantee correct ordering, or drop rows entirely. Only the windowed `lag` approach yields a per-row previous timestamp at the original granularity.

Exam trap

The trap here is reaching for `sortWithinPartitions` or `lead` and assuming partition-local ordering plus the wrong offset function reproduces the previous-value semantics of a properly windowed `lag`.

244
MCQhard

A developer has a DataFrame `trades` with columns `trade_id`, `symbol`, and `price`. They want to add a column `prev_price` containing the price of the immediately preceding trade for the same symbol, ordered by `trade_id` ascending, without collapsing rows. Which transformation should they use?

A.trades.orderBy("trade_id").dropDuplicates(["symbol"])
B.trades.withColumn("prev_price", lag("price").over(Window.partitionBy("symbol").orderBy("trade_id")))
C.trades.groupBy("symbol").agg(last("price"))
D.trades.withColumn("prev_price", lead("price").over(Window.partitionBy("symbol").orderBy("trade_id")))
AnswerB

The lag window function returns the value of price from the row that is one position earlier within the window defined by partitionBy("symbol") and orderBy("trade_id"). Using withColumn keeps the original row count, and the result is exactly the previous trade price per symbol, matching the requirement precisely.

Why this answer

lag with a window partitioned by symbol and ordered by trade_id returns the prior row's price while preserving every row, which is precisely what adding prev_price requires. lead looks forward, aggregation collapses rows, and dropDuplicates discards data, so none of those alternatives satisfies the requirement.

Exam trap

The trap here is confusing lag with lead, or believing that any windowed aggregation will preserve row count when only analytic window functions like lag and lead do so.

245
Multi-Selecthard

Which TWO of the following statements accurately describe the role of the Spark Executor in a cluster deployment?

Select 2 answers
A.Executors are responsible for scheduling tasks across the worker nodes.
B.Executors execute the tasks dispatched by the Driver.
C.Executors perform the storage of data in memory or on disk.
D.Executors create the physical execution plan from user code.
E.Executors coordinate the cluster-wide resource allocation for the job.
AnswersB, C

Executors serve as the execution environment for tasks assigned by the Driver. Upon receiving a task, the executor deserializes the code and executes it against the data partitions stored on or fetched to the worker node, reporting the task status and metrics back to the Driver periodically.

Why this answer

Executors are the workhorses of the Spark architecture, responsible for executing tasks and storing data. Recognizing their duality—compute and storage—is critical for tuning Spark applications. If executors are improperly sized, they can lead to OOM errors or underutilization of cluster resources.

Understanding how executors manage task parallelism and block storage allows developers to optimize memory settings and partition counts effectively for high-performance data processing pipelines.

Exam trap

Candidates frequently attribute cluster coordination, DAG scheduling, and global metadata management to executors, forgetting that executors strictly execute tasks and cache data blocks.

246
MCQhard

A job joining a large fact table with a small dimension table runs out of memory on executors during the join. The dimension table is about 40 MB after filtering and the configured spark.sql.autoBroadcastJoinThreshold is 10 MB. The join key is highly skewed in the fact table. Which action is most appropriate?

A.Raise spark.sql.autoBroadcastJoinThreshold above the dimension table size so the small side is broadcast and no shuffle of the fact table occurs.
B.Repartition the fact table by the join key with a high partition count before the join to distribute the skew.
C.Increase spark.executor.memory so each executor can hold the skewed partitions of the fact table during the sort-merge join.
D.Salt the join key on both sides of the join so the skewed values are spread across many partitions.
AnswerA

Broadcasting the filtered 40 MB dimension table eliminates the shuffle of the large fact table and turns the join into a map-side operation. This avoids the memory pressure and shuffle associated with a sort-merge join on a skewed key. The threshold is a size guardrail, so raising it to cover the known small side is the intended control. This directly removes the source of the executor memory failure.

Why this answer

The dimension table is small enough to broadcast once the threshold is raised, which converts the join into a map-side operation and removes the shuffle of the large, skewed fact table. That eliminates the executor memory pressure caused by concentrating hot-key rows in a few reduce tasks. Salting and repartitioning address skew only in large-to-large joins and add cost here, while adding executor memory merely postpones the failure without changing the underlying plan.

Exam trap

The trap here is treating a large-to-small join as a skew problem to be salted, when broadcasting the small side removes the shuffle entirely.

247
MCQmedium

A developer is writing a PySpark job that must read a Parquet dataset, apply several transformations, and then write the result back to storage. They want to ensure the schema of the written data is inferred directly from the DataFrame rather than from any external definition, and they want to append to an existing Parquet directory. Which write configuration accomplishes this?

A.df.write.mode('append').format('json').save('/path/output')
B.df.write.mode('append').option('mergeSchema', 'true').parquet('/path/output')
C.df.write.mode('overwrite').format('parquet').save('/path/output')
D.df.write.mode('append').format('parquet').save('/path/output')
AnswerD

Writing with the DataFrame writer in `append` mode and the `parquet` format stores the DataFrame using its own schema, since Parquet embeds the schema in the file metadata. Appending adds new files to the existing directory without removing prior data. This directly satisfies both requirements: schema derived from the DataFrame and append semantics.

Why this answer

Appending Parquet output with the DataFrame writer uses the DataFrame's own schema because Parquet stores schema metadata in each file. The overwrite mode would destroy existing data, JSON output changes the storage format, and `mergeSchema` is a read-side option that does not affect how the DataFrame schema is written. The append-mode Parquet write is the only configuration that meets both conditions.

Exam trap

The trap here is assuming `mergeSchema` is a write option, when it is actually applied when reading Parquet files with differing schemas.

248
MCQmedium

A data engineer is working with a Spark SQL DataFrame in Databricks that has a column named event_time stored as a string in the format 'yyyy-MM-dd HH:mm:ss'. They need to filter rows where event_time falls within the last 7 days relative to the current timestamp. Which Spark SQL expression correctly achieves this?

A.SELECT * FROM events WHERE date_format(event_time, 'yyyy-MM-dd') >= date_sub(current_date(), 7)
B.SELECT * FROM events WHERE to_timestamp(event_time, 'yyyy-MM-dd HH:mm:ss') >= current_timestamp() - INTERVAL 7 DAYS
C.SELECT * FROM events WHERE event_time >= current_date() - 7
D.SELECT * FROM events WHERE unix_timestamp(event_time) >= unix_timestamp(current_timestamp()) - 604800
AnswerB

This expression uses to_timestamp with the correct format pattern to convert the string column into a timestamp, then compares it against current_timestamp() minus a 7-day interval. Spark SQL supports INTERVAL 7 DAYS syntax, and this correctly filters rows within the last week. It is the only option that both parses the string format correctly and uses a valid interval expression.

Why this answer

The correct expression uses to_timestamp to parse the string column with the explicit format, then compares it to current_timestamp() minus an interval of 7 days. This ensures accurate filtering based on the full timestamp, including time-of-day, and leverages Spark SQL's built-in interval arithmetic. Other options either ignore the time component, rely on implicit casting that may fail, or use less precise date-only comparisons.

Exam trap

The trap here is assuming that subtracting an integer from a date or timestamp will work as expected, when Spark SQL requires explicit interval syntax for timestamp arithmetic.

249
MCQmedium

What is the primary role of the 'Cluster Manager' in Spark?

A.It manages the Spark DAG and optimizes the query plan.
B.It handles the physical allocation of resources like CPU and memory.
C.It executes the tasks and stores intermediate shuffle data.
D.It monitors the progress of individual Spark tasks.
AnswerB

The Cluster Manager interacts with the underlying infrastructure to negotiate resource requests made by the Driver. It allocates containers on nodes where the Spark executors can run, ensuring that the Spark application has the requested compute power and memory capacity to execute its tasks according to the configuration.

Why this answer

The Cluster Manager, such as Kubernetes or YARN, acts as the resource broker for the Spark application. It is vital to understand that it does not manage the Spark execution logic (the Driver does that). Instead, it provides the 'raw materials'—the executor containers—that the Driver needs to run tasks.

Misunderstanding this can lead to incorrect assumptions about where failures occur: application logic failures happen in the Driver/Executors, while resource availability issues happen in the Manager.

Exam trap

Candidates mistakenly believe the Cluster Manager controls the job's internal execution logic, rather than acting solely as a resource provider for the Spark application.

250
MCQeasy

A developer submits a Spark application to a Databricks cluster using spark-submit with deploy mode set to cluster. During execution, one of the worker nodes hosting a task fails and is lost by the cluster manager. Which Spark component is responsible for rescheduling the failed task on another available executor?

A.The Task Scheduler
B.The Catalyst Optimizer
C.The Cluster Manager
D.The DAG Scheduler
AnswerA

The Task Scheduler in the Spark driver monitors each task within a stage, receives status updates from executors, and when a task fails or its executor is lost, it marks the task as failed and relaunches it on another available executor, up to spark.task.maxFailures. This is exactly the component that handles task-level fault tolerance in the scenario.

Why this answer

When an executor is lost, the driver's Task Scheduler detects the failed tasks and relaunches them on other executors, honoring spark.task.maxFailures before failing the stage. The DAG Scheduler only regenerates stages when shuffle map outputs are lost. The cluster manager merely allocates resources and reports executor loss, and Catalyst is a compile-time optimizer with no runtime task responsibility.

Exam trap

The trap here is assuming that because the DAG Scheduler builds the execution graph, it also owns per-task retries, when in fact task-level fault tolerance belongs to the Task Scheduler.

251
MCQeasy

A data analyst wants to connect a local Python script to a Databricks cluster using Spark Connect. The workspace URL and a personal access token are available. Which client-side step is required to create the remote Spark session?

A.Add the Databricks JDBC driver to the classpath and open a JDBC connection with the token as the password.
B.Call `SparkSession.builder.remote("sc://<workspace-url>:443/;token=<token>;use_ssl=true").getOrCreate()`.
C.Set the `SPARK_HOME` environment variable to the Databricks workspace URL and call `SparkSession.builder.getOrCreate()`.
D.Install and start a local Spark master with `start-master.sh`, then connect the script to that local master.
AnswerB

The `remote` method on the builder accepts a Spark Connect connection string. The `sc://` scheme with host, port, token, and SSL parameters establishes the gRPC connection to the Databricks workspace. This is the documented pattern for creating a remote session from a local Python environment without a local Spark installation.

Why this answer

Spark Connect clients create a remote session by passing a connection string to the builder's `remote` method or the `connect` function. The `sc://` URL encodes the workspace host, port 443, authentication token, and SSL flag, allowing the thin client to establish a gRPC channel to the Databricks cluster.

Exam trap

The trap here is confusing local Spark configuration variables like `SPARK_HOME` with the remote connection mechanism that Spark Connect actually requires.

252
MCQhard

A developer needs to optimize a Spark application that performs repetitive filtering and grouping on the same large DataFrame. Which feature should they implement to improve performance?

A.Increase the number of cores per executor.
B.Use the cache() or persist() method on the DataFrame.
C.Switch the file format from Parquet to CSV.
D.Enable speculative execution.
AnswerB

Caching pins the DataFrame in memory or disk, allowing subsequent actions to read the computed data directly rather than re-running the entire lineage. This significantly improves performance for iterative processing tasks, reducing latency and avoiding repeated data reads and transformations that occur in non-cached, re-evaluated Spark DataFrames.

Why this answer

Caching (or persisting) the DataFrame keeps it in memory or on disk for subsequent actions. This is essential for iterative algorithms or workloads where the same data is reused, as it avoids recomputing the lineage from scratch. Understanding the trade-offs between memory and disk persistence is fundamental for building performant Databricks pipelines that maximize resource reuse and minimize redundant computation time.

Exam trap

Candidates frequently confuse dataframe lineage optimization or broadcast hints with caching, missing that reused DataFrames need explicit persistence.

253
Multi-Selecthard

Which TWO of the following statements accurately describe the relationship between Spark Executors and memory management within a Databricks cluster?

Select 2 answers
A.Executors use fixed memory boundaries that cannot be adjusted during task execution.
B.The storage memory region is primarily used for caching RDDs and DataFrames.
C.Execution memory is reserved for intermediate shuffle and join calculations.
D.Executors are allowed to access the Driver's memory pool to store large datasets.
E.Memory management is handled entirely by the Cluster Manager, not the Spark process.
AnswersB, C

Storage memory is dedicated to keeping serialized or deserialized data in memory for rapid access. When a user explicitly calls cache() or persist() on a DataFrame, Spark stores these partitions in this region to avoid recomputing data from source files during subsequent iterations or multi-pass operations.

Why this answer

Executors manage memory through a unified memory manager, partitioning heap space between storage (caching) and execution (shuffles/joins). Understanding this architecture is vital because improper memory configuration leads to OOM errors or excessive spilling to disk. By balancing the memory pools dynamically, Spark avoids hard boundaries, allowing execution tasks to borrow space from storage when cache utilization is low, significantly improving overall job throughput.

Exam trap

Candidates often assume memory regions are static and strictly partitioned. They fail to realize that Spark uses a unified memory manager, allowing dynamic borrowing between storage and execution regions.

254
MCQhard

A Spark job writes a large DataFrame to a Delta table partitioned by date. The job is taking much longer than expected, and the Spark UI shows that many tasks are writing very small files. You have already set spark.sql.shuffle.partitions to 200. What is the most effective way to reduce the number of small files written?

A.Set spark.sql.files.maxRecordsPerFile to a high value to combine records into fewer files.
B.Use coalesce(1) before writing to reduce the number of output files.
C.Increase spark.sql.shuffle.partitions to 2000.
D.Enable optimized writes by setting spark.databricks.delta.optimizeWrite.enabled to true.
AnswerD

Optimized writes automatically coalesce small files during the write operation by adding a shuffle step that reduces the number of output files based on the data size. This is specifically designed to mitigate the small file problem in Delta Lake on Databricks. It balances file sizes without manual repartitioning, making it the most effective solution for reducing small files while maintaining parallelism.

Why this answer

Enabling optimized writes on Databricks triggers an automatic shuffle before writing to Delta, which coalesces data into fewer, larger files. This directly addresses the small file problem without sacrificing parallelism or requiring manual tuning. The other options either worsen the issue, introduce bottlenecks, or do not target the root cause of many small files from partitioned writes.

Exam trap

The trap here is thinking that increasing shuffle partitions or coalescing to one partition will solve small files, when optimized writes are the Databricks-specific feature designed for this.

255
MCQmedium

A data engineer is building a Spark Connect application and wants to attach a small lookup table to every task without a shuffle. The table is 20 MB and the cluster has default settings. Which approach should the engineer use?

A.Call `broadcast(lookup_df)` from `pyspark.sql.functions` and join it with the fact DataFrame.
B.Set `spark.sql.autoBroadcastJoinThreshold` to -1 and join normally.
C.Call `lookup_df.cache()` and then join, relying on the cache to avoid the shuffle.
D.Use `lookup_df.repartition(1)` before the join to force a single partition.
AnswerA

The `broadcast` hint tells the server-side optimizer to broadcast the small relation to all executors, avoiding a shuffle of the large fact table. Because the lookup is only 20 MB, it fits comfortably under the default auto broadcast join threshold, and the hint makes the intent explicit. This is the supported way to request a broadcast join in a Spark Connect application.

Why this answer

To broadcast a small relation in Spark Connect, the developer uses the `broadcast` function from `pyspark.sql.functions` as a join hint. The hint is serialized into the logical plan and evaluated by the server-side optimizer, which decides to replicate the small side to all executors. Caching, repartitioning to one partition, or disabling the broadcast threshold do not produce a broadcast join and can add unnecessary shuffle or memory pressure.

Exam trap

The trap here is thinking that caching or repartitioning a small DataFrame causes a broadcast join, when only the broadcast hint or the auto broadcast threshold influences the join strategy chosen by the server optimizer.

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

257
MCQeasy

You need to add a new calculated column 'discounted_price' to an existing DataFrame 'df' by multiplying 'price' by 0.9. Which DataFrame transformation accomplishes this correctly?

A.df.addColumn('discounted_price', col('price') * 0.9)
B.df.update('discounted_price', df.price * 0.9)
C.df.withColumn('discounted_price', col('price') * 0.9)
D.df.select(col('*'), col('price') * 0.9 as 'discounted_price')
AnswerC

withColumn returns a new DataFrame with the added or replaced column, and multiplying the price column by 0.9 computes the discounted value per row. This is the standard immutable transformation, matching the requirement to add discounted_price without mutating the original DataFrame.

Why this answer

Adding columns immutably is a core pattern in Spark DataFrame operations. Using the withColumn method creates a new DataFrame reference with the added transformation while leaving the original dataset unchanged, adhering to functional programming principles required for distributed data processing.

Exam trap

Candidates often try to modify the DataFrame in place, forgetting that Spark DataFrames are immutable; they fail to assign the result of 'withColumn' back to a variable.

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

259
MCQmedium

You are troubleshooting a Spark Connect application that fails to connect to a Databricks cluster. The error message indicates an authentication failure. Which of the following is the most likely cause?

A.The client machine's firewall is blocking outbound connections on port 443.
B.The Databricks cluster is not running or is in a terminated state.
C.The personal access token (PAT) has expired or is invalid.
D.The client is using an outdated version of the Spark Connect protocol.
AnswerC

Authentication failures in Spark Connect often stem from invalid or expired credentials. The personal access token is used to authenticate the client to the Databricks workspace. If the token is expired, revoked, or incorrect, the server rejects the connection. Checking the token's validity and permissions is the first step in troubleshooting.

Why this answer

Authentication failures in Spark Connect are typically due to invalid credentials, such as an expired or incorrect personal access token. The token must be valid and have the necessary permissions to access the cluster. Other issues like network problems or cluster state would produce different error messages, making the token the primary suspect.

Exam trap

The trap here is attributing an authentication failure to network or cluster state issues, when the error message explicitly points to credentials being the problem.

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

261
MCQhard

A developer is using the pandas API on Spark in a Databricks notebook. They have a pandas-on-Spark DataFrame psdf with a default index. They call psdf.sort_values('amount') and then psdf.head(10). They observe that the resulting index values are not sequential from 0 to 9, but instead appear as arbitrary integers. What is the most likely explanation for this behavior?

A.The default index type is 'distributed-sequence', which preserves the original row order and does not reassign new sequential indices after sorting.
B.The default index type is 'sequence', which creates a single-partition sequence index that is lost during shuffles such as sorting.
C.The default index type is 'distributed-sequence', but sort_values triggers a shuffle that resets the index to arbitrary values because sequence indices cannot survive shuffles.
D.The default index type is 'distributed', which assigns arbitrary but stable integers to rows and does not guarantee sequential values after operations like sort_values.
AnswerD

By default, pandas-on-Spark uses the 'distributed' index type when no index is specified. This type assigns each row a unique integer that is stable across operations but not necessarily sequential or contiguous. After sort_values, the index values are carried along with the rows, so they appear out of order. This is expected behavior and differs from pandas, where sorting resets the index only if reset_index is called.

Why this answer

The default index type in pandas API on Spark is 'distributed', which assigns stable but non-sequential integers to rows. When you sort a DataFrame, the existing index values travel with the rows, so the resulting index is not 0..n-1. To get sequential indices after sorting, you must explicitly call reset_index() or set the index type to 'distributed-sequence' before sorting.

Exam trap

The trap here is assuming that pandas-on-Spark mimics pandas' default behavior of producing a sequential index after sorting, when the default 'distributed' index intentionally preserves arbitrary but stable integers.

262
MCQmedium

A PySpark job on Databricks repeatedly calls `df.count()` and `df.show()` inside a loop across 40 iterations, and the Spark UI shows the identical lineage being recomputed on every iteration even though the source Delta table is unchanged. You want to avoid re-executing the upstream transformations without materializing the data to disk. What should you do?

A.Wrap the DataFrame in `spark.createDataFrame(df.rdd)` so the derived object keeps the computed rows in the driver.
B.Call `df.persist(StorageLevel.MEMORY_AND_DISK)` before the loop and `df.unpersist()` after it.
C.Set `spark.sql.adaptive.enabled` to true so the optimizer reuses prior query results automatically.
D.Add `.repartition(200)` to the DataFrame before the loop so each iteration reads a different partition set.
AnswerB

Caching the DataFrame before the loop stores the computed partitions in executor memory (spilling to disk when needed), so each subsequent action reads the cached blocks instead of walking the lineage again. Since the source is unchanged, the cached result stays valid for all 40 iterations, and unpersisting releases the memory once the loop finishes.

Why this answer

Because the same DataFrame is acted on repeatedly and the underlying Delta table does not change, persisting the DataFrame keeps the computed partitions available across iterations, so the lineage is evaluated once rather than 40 times. Repartitioning, enabling adaptive execution, or round-tripping through RDDs all leave the recomputation behavior intact.

Exam trap

The trap here is assuming that enabling an optimizer feature automatically reuses results across actions, when only an explicit cache or persist call stores computed partitions.

263
MCQmedium

A data engineer needs to join two large DataFrames, `sales` and `products`, on the `product_id` column. The `products` DataFrame is extremely small and fits entirely in a single executor's memory. To optimize performance and avoid a costly shuffle join across the network, which strategy should be applied using the Spark DataFrame API?

A.Call `sales.join(broadcast(products), "product_id")` to explicitly push the small DataFrame to all worker nodes.
B.Partition both DataFrames explicitly by `product_id` using `repartition(col("product_id"))` prior to joining.
C.Increase the shuffle partition count configuration via `spark.sql.shuffle.partitions` to a higher value.
D.Cache the `sales` DataFrame in memory before executing the standard inner join operation.
AnswerA

Wrapping the smaller DataFrame with the broadcast function forces the Catalyst optimizer to use a broadcast hash join. This eliminates the shuffle phase for the large sales dataset, significantly reducing overall execution time and resource contention across the cluster.

Why this answer

Broadcasting small datasets avoids expensive network shuffles by copying the entire small DataFrame to all worker nodes. This optimization drastically improves query performance for star schemas and dimension table lookups. Using the broadcast function explicitly guarantees that the Catalyst optimizer chooses a broadcast hash join instead of a sort-merge join.

Exam trap

Candidates often confuse broadcast joins with repartitioning. They might attempt to manually repartition both DataFrames, which triggers a massive, unnecessary shuffle, rather than using the broadcast hint to replicate the small table.

264
MCQhard

A data engineer is using Spark Connect from a local Python environment to connect to a Databricks cluster. They attempt to use the spark.sparkContext.broadcast() method to broadcast a large lookup dictionary for use in a UDF. The code fails. What is the most likely reason for this failure?

A.Broadcast variables are only supported when using the Scala API, not the Python API, in Spark Connect.
B.The broadcast variable size exceeds the maximum allowed by Spark Connect, causing a failure.
C.Broadcast variables are not supported in Spark Connect because the client does not have access to the SparkContext.
D.The broadcast variable must be created using the SparkSession.broadcast() method instead.
AnswerC

This is correct because Spark Connect clients do not have a SparkContext; they interact with the Spark server through a remote client. The broadcast() method is part of the SparkContext API, which is not available in Spark Connect. Therefore, attempting to use spark.sparkContext.broadcast() will fail because spark.sparkContext is not defined in the Spark Connect session.

Why this answer

Spark Connect clients do not have a SparkContext, so APIs like broadcast() that rely on it are unavailable. The client interacts with the remote Spark server through a thin client, and operations requiring direct driver access, such as creating broadcast variables, are not supported. Alternative approaches like using joins with small DataFrames should be used.

Exam trap

The trap here is assuming that Spark Connect supports all SparkContext APIs, when in fact it intentionally omits them to maintain decoupling.

265
MCQhard

A data engineer has two DataFrames: `orders` with columns `order_id` and `customer_id`, and `customers` with columns `customer_id` and `region`. They must produce a result containing only orders whose `customer_id` exists in `customers`, keeping every matching order row exactly once with no customer columns added. Which operation should be used?

A.orders.join(customers, 'customer_id', 'left_anti')
B.orders.join(customers, 'customer_id', 'left_semi')
C.orders.join(customers, 'customer_id', 'inner')
D.orders.union(customers.select('customer_id').withColumnRenamed('customer_id', 'order_id'))
AnswerB

A left semi join returns rows from `orders` that have a match in `customers`, includes only the left DataFrame's columns, and does not duplicate rows even when the right side has multiple matches. This precisely satisfies the requirement of filtering orders by existence without adding customer columns or multiplying rows.

Why this answer

A left semi join keeps only left-side rows that have a matching key on the right, returns just the left columns, and never duplicates rows regardless of duplicate keys on the right. An inner join would add the region column and could duplicate rows, a left anti join inverts the logic, and a union mixes incompatible data. The semi join matches the requirement exactly.

Exam trap

The trap here is defaulting to an inner join for existence filtering, ignoring that it brings along right-side columns and can duplicate rows when the right side has repeated keys.

266
MCQmedium

A developer is using Spark Connect to run a PySpark application against a remote Databricks cluster. The application calls df.cache() on a DataFrame that is used multiple times. Which statement accurately describes how caching behaves in this scenario?

A.The cache is stored in the client process memory, allowing subsequent operations to avoid round trips to the server.
B.The cache is automatically persisted to cloud object storage and shared across all Spark Connect clients connected to the same cluster.
C.Calling cache() has no effect because Spark Connect does not support caching operations on remote DataFrames.
D.The cache is managed on the server side, and the cached data is stored in the cluster's memory or disk according to the configured storage level.
AnswerD

In Spark Connect, cache() is a logical plan operation sent to the server. The server executes the caching, storing the data in the cluster's memory or disk based on the storage level. The client does not hold cached data. This means cache behavior and memory management are controlled by the remote Spark server, just as in a traditional Spark application.

Why this answer

In Spark Connect, cache() is a server-side operation. The logical plan includes the cache directive, and the remote Spark server executes it, storing the data in cluster memory or disk according to the storage level. The client does not hold cached data, so all subsequent operations still communicate with the server.

This preserves the semantics of caching while maintaining the thin-client architecture.

Exam trap

The trap here is assuming that because the client is remote, caching might happen locally or be unsupported, when in fact it is executed server-side as usual.

267
MCQhard

When using 'apply_batch' in Pandas-on-Spark, how does the function behave regarding the input data?

A.It applies the function to each row individually.
B.It applies the function to each partition as a Pandas DataFrame.
C.It forces a shuffle of all data to a single partition.
D.It requires the use of UDFs (User Defined Functions) with Python serialization.
AnswerB

This method operates by converting each partition of the Spark DataFrame into a standard Pandas DataFrame and applying the user-defined function. This allows developers to use the full power of the Pandas library on Spark data without needing to pull the entire dataset into the driver memory.

Why this answer

The 'apply_batch' function provides a bridge to perform custom Pandas operations on partitions of data. By passing a function that accepts a Pandas DataFrame, you can leverage native Pandas logic on individual Spark partitions. This is essential for complex logic not supported by the Spark engine, allowing for efficient data processing without leaving the Pandas-on-Spark ecosystem, provided the input data is partitioned appropriately for the task at hand.

Exam trap

Candidates mistakenly believe 'apply_batch' operates on the entire DataFrame at once or individual rows. It actually processes data at the partition level, which is critical for distributed performance.

268
MCQhard

A developer is building a PySpark job that reads a Parquet dataset with 2,000 files into `df`. They call `df.cache()` and then execute three separate actions in the same session. The Spark UI shows the Parquet files are read from storage three times, and the cache never appears in the Storage tab. What is the most likely cause?

A.`cache()` was applied to `df` but a variable reassignment such as `df = df.withColumn(...)` created a new logical plan that no longer references the cached node.
B.`cache()` only takes effect after `persist(StorageLevel.MEMORY_ONLY)` is called, so the default call does nothing.
C.`cache()` stores data in memory only, and the dataset is too large, so Spark silently drops the cached blocks.
D.The Parquet files are stored in a format that Spark cannot cache, so `cache()` is a no-op for Parquet sources.
AnswerA

This is correct because `cache()` marks a specific logical plan node. If `df` is later reassigned to a transformation like `withColumn`, the new DataFrame references an uncached plan, so actions re-read Parquet. The original cached node is never triggered, which explains both the repeated file reads and the empty Storage tab.

Why this answer

Caching in Spark marks a specific logical plan node as persistent. Any subsequent transformation that derives a new DataFrame from it produces a different plan that does not include the cached node, so actions against the new DataFrame bypass the cache entirely. The tell-tale signs here are repeated Parquet reads and an empty Storage tab.

To fix this, either cache after the final transformation or reuse the exact same cached DataFrame reference across actions.

Exam trap

The trap here is assuming `cache()` is a property of the data itself, when it is actually attached to a specific logical plan node.

269
MCQeasy

What is the primary role of the 'Executor' process in the Spark distributed architecture?

A.To serve as the central coordinator for the entire Spark application.
B.To execute tasks and manage local storage for RDDs.
C.To communicate with the cluster manager to request additional executors.
D.To define the DAG of stages for the Spark job.
AnswerB

Executors are the workhorses of the cluster. They run the code specified by the user in the form of tasks, manage memory for cached RDDs, and report their health back to the driver. This separation of concerns allows the cluster to scale horizontally while the driver maintains central control.

Why this answer

The executor is the worker process responsible for performing the actual data processing tasks. It manages memory for storage and execution, reports task status back to the driver, and interacts with local disk storage for caching or shuffling. By offloading these intensive tasks to multiple executors, Spark achieves horizontal scalability, allowing the system to process massive datasets by distributing the workload across a cluster of nodes.

Exam trap

Test-takers frequently mistake the Executor for the Driver, confusing task execution and local storage management with cluster coordination and master scheduling responsibilities.

270
MCQeasy

Which of the following describes the behavior of a 'Left Outer Join' in Spark SQL?

A.It returns only rows that have matches in both the left and right tables.
B.It returns all rows from the left table and matched rows from the right table.
C.It returns all rows from the right table and matched rows from the left table.
D.It excludes all rows from both tables that do not have a matching key.
AnswerB

The definition of a Left Outer Join is that it preserves the entire left side of the join. For rows that do not have a corresponding key in the right table, Spark fills the right-side columns with NULLs, allowing the developer to see the full list from the left dataset.

Why this answer

A Left Outer Join is fundamental in SQL for merging datasets where you want to keep all records from the left table, regardless of whether a match exists in the right table. Understanding this ensures that data integration processes correctly handle non-matching records, preventing data loss during joins. It is a critical concept for analysts and developers building reports that require comprehensive views of primary entities with optional supplemental info.

Exam trap

Candidates often confuse Left Outer Joins with Full Outer Joins, mistakenly thinking the result set includes non-matching rows from BOTH tables rather than just the left table.

271
MCQmedium

A data scientist is using Pandas API on Spark to process a large dataset. They call `psdf.apply(lambda row: row['a'] + row['b'], axis=1)` and notice extremely slow performance. Which statement best explains why this operation is inefficient and what alternative should be used?

A.`apply` with axis=1 is not supported in Pandas API on Spark and will always raise an error; the alternative is to use `psdf.apply` with axis=0.
B.`apply` with axis=1 is inefficient because it uses a Python UDF that processes rows individually, preventing Spark optimizations; using vectorized column operations like `psdf['a'] + psdf['b']` is much faster.
C.`apply` with axis=1 triggers a full shuffle and should be replaced with `groupby().applyInPandas()` for better performance.
D.`apply` with axis=1 is slow because it forces a conversion to a pandas DataFrame on the driver; the fix is to enable Arrow-based conversion with `spark.sql.execution.arrow.pyspark.enabled`.
AnswerB

`apply` with axis=1 in Pandas API on Spark is implemented via a Python UDF that iterates row by row, which incurs serialization overhead and blocks Spark's Catalyst optimizer from optimizing the expression. Vectorized column operations such as `psdf['a'] + psdf['b']` are translated into native Spark expressions and execute efficiently in the JVM. This alternative avoids Python UDF overhead and leverages distributed processing.

Why this answer

Row-wise `apply` with axis=1 in Pandas API on Spark is implemented using a Python UDF that processes each row individually, which is slow due to serialization and loss of Spark optimizations. The efficient alternative is to express the logic using vectorized column operations, such as `psdf['a'] + psdf['b']`, which are translated into native Spark expressions and execute in the JVM. This approach avoids Python UDF overhead and scales with the cluster.

Exam trap

The trap here is assuming that row-wise `apply` is optimized or that enabling Arrow will fix its performance, when the real issue is the row-by-row Python UDF execution that vectorized column operations avoid.

272
MCQmedium

You have a large DataFrame containing user transaction logs. You need to read this data and immediately repartition it by user_id to optimize downstream filtering operations. Which DataFrame API method should you use?

A.df.coalesce('user_id')
B.df.repartition('user_id')
C.df.partitionBy('user_id')
D.df.shuffle('user_id')
AnswerB

Repartitioning by a specific column triggers a wide transformation that performs a full shuffle across the cluster. This groups all rows with the same user_id into the same partition, significantly improving the performance of subsequent filters and joins on that key.

Why this answer

The repartition method creates a new set of partitions across the cluster network, which helps distribute data evenly to prevent skew. This is a critical transformation in Spark to optimize shuffle performance for downstream queries and aggregations by ensuring balanced workloads across executors.

Exam trap

Candidates often choose 'coalesce()' instead of 'repartition()' because they think it is faster, forgetting that 'coalesce' is designed to reduce partitions without a full shuffle, which can cause data skew.

273
MCQeasy

Which component in the Spark cluster architecture is responsible for communicating directly with the Cluster Manager (e.g., YARN, Mesos, K8s) to request and release resources?

A.Executor
B.Worker Node
C.Spark Driver
D.Task Scheduler
AnswerC

The Spark Driver contains the SparkContext and the scheduler backends. It acts as the application master, negotiating with the cluster manager to acquire executors. Without the driver, the Spark application would have no way to obtain the compute resources required to process data in a distributed manner.

Why this answer

The Spark Driver is responsible for resource negotiation with the underlying Cluster Manager. It maintains the SparkContext and coordinates with the resource manager to acquire executors for the application. Once executors are acquired, the driver schedules tasks on them.

This role is fundamental to the driver-executor model, as the driver serves as the central control point for the application's lifecycle and hardware interaction.

Exam trap

Candidates incorrectly select the 'Executor' or 'DAG Scheduler' as the component that negotiates resources, failing to realize the Driver acts as the sole intermediary.

274
MCQhard

A developer must combine two DataFrames, `left_df` and `right_df`, on a key column `id`. They need every row from `left_df` regardless of whether a match exists in `right_df`, and matching rows from `right_df` where available, with unmatched right-side columns filled as null. Which join invocation produces this result?

A.left_df.join(right_df, on="id", how="right_outer")
B.left_df.join(right_df, on="id", how="inner")
C.left_df.join(right_df, on="id", how="left_outer")
D.left_df.join(right_df, on="id", how="full_outer")
AnswerC

A left outer join preserves every row from left_df, attaching matching right-side columns where the key matches and filling them with null otherwise. This exactly satisfies the stated requirement that all left rows survive while unmatched right columns become null.

Why this answer

A left outer join retains every left row and fills right-side columns with null when no match exists, which is exactly what the scenario specifies. Inner join discards unmatched left rows, right outer join preserves the wrong side, and full outer join introduces unmatched right rows that were not requested.

Exam trap

The trap here is assuming that left_outer and full_outer are interchangeable when only one side's unmatched rows are relevant, which quietly adds or removes rows from the result.

275
MCQmedium

A developer has a Spark SQL DataFrame `df` with an array column named `scores` containing integers. They need to create a new column `passing` that is true only when every element in `scores` is greater than or equal to 70. Which Spark SQL higher-order function should they use?

A.exists(scores, s -> s >= 70)
B.filter(scores, s -> s >= 70)
C.forall(scores, s -> s >= 70)
D.transform(scores, s -> s >= 70)
AnswerC

The forall higher-order function returns true only if the lambda predicate evaluates to true for every element in the array. In this scenario, forall(scores, s -> s >= 70) yields a single boolean column that is true exactly when all scores are at least 70, which matches the requirement. It short-circuits on the first false, making it efficient.

Why this answer

The forall higher-order function is designed to test whether every element in an array satisfies a given predicate, returning a single boolean. In this scenario, it correctly produces a column that is true only when all scores are 70 or higher. transform maps each element, filter selects elements, and exists checks for at least one match, none of which yield the required all-elements condition.

Exam trap

The trap here is confusing higher-order functions that return arrays (transform, filter) with those that return booleans (exists, forall), and mixing up exists (any) with forall (all).

276
MCQmedium

A data scientist is using Spark SQL to compute a running total of sales amounts for each customer, ordered by transaction date. The query must return, for each row, the sum of all previous sales for that customer up to and including the current row. Which window specification should be used?

A.OVER (PARTITION BY customer_id ORDER BY transaction_date ROWS BETWEEN CURRENT ROW AND UNBOUNDED FOLLOWING)
B.OVER (PARTITION BY customer_id ORDER BY transaction_date ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
C.OVER (PARTITION BY customer_id ORDER BY transaction_date ROWS BETWEEN 1 PRECEDING AND CURRENT ROW)
D.OVER (PARTITION BY customer_id ORDER BY transaction_date RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW)
AnswerB

This window specification partitions by customer_id, orders by transaction_date, and defines the frame from the start of the partition up to the current row. SUM(sales_amount) over this window yields a running total per customer. The ROWS frame is explicit and ensures each row includes all preceding rows, which is exactly the requirement. It is the standard approach for cumulative sums.

Why this answer

A running total requires a window frame that starts at the beginning of the partition and ends at the current row. The specification with ROWS BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW, combined with PARTITION BY customer_id and ORDER BY transaction_date, achieves this by summing all rows up to and including the current one for each customer. Other frames either sum future rows, include tied rows, or limit to a small window, none of which produce the desired cumulative sum.

Exam trap

The trap here is confusing ROWS with RANGE, or using a frame that sums future rows instead of preceding ones, which yields incorrect cumulative totals when duplicate dates exist.

277
MCQmedium

Which process is responsible for tracking the location of data blocks cached in the executors?

A.The TaskScheduler
B.The BlockManager
C.The DAGScheduler
D.The Resource Manager
AnswerB

The BlockManager is responsible for the storage and retrieval of data blocks, whether in memory, on disk, or off-heap. The driver maintains the global mapping of all blocks, allowing the scheduler to make intelligent decisions based on data locality, which is essential for maximizing performance in distributed clusters.

Why this answer

The BlockManager is a critical component that lives on both the driver and the executors. While each executor's BlockManager handles the local storage of blocks, the master BlockManager on the driver maintains a registry of where every block is located across the entire cluster. This is vital for Spark's ability to schedule tasks near the data (data locality) and minimize unnecessary data shuffling during query execution.

Exam trap

Many candidates incorrectly attribute block tracking to the Driver's SparkContext or the Cluster Manager, overlooking the specialized internal role of the BlockManager in maintaining the distributed data registry.

278
MCQeasy

Which Spark component is responsible for maintaining the Directed Acyclic Graph (DAG) of stages and tasks?

A.The Executor
B.The Driver
C.The Cluster Manager
D.The Storage Manager
AnswerB

The Driver is responsible for translating user code into a DAG of stages. It optimizes this graph and decides the execution order of stages. By managing the DAG, the driver ensures that transformations are performed efficiently, minimizing shuffles and maximizing parallel processing performance across the distributed cluster environment.

Why this answer

The Spark Driver is the central coordinator. It generates the DAG of stages based on the user's transformations and the physical execution plan. This DAG is essential for Spark to optimize query plans through catalyst and execute tasks in the correct dependency order.

Understanding this architectural role is fundamental to knowing how Spark converts high-level DataFrame API calls into low-level distributed operations executed by the cluster.

Exam trap

Candidates often guess the Cluster Manager or the DAGScheduler component separately, failing to realize the Driver is the overarching process that hosts these internal scheduling and planning components.

279
MCQmedium

In a Databricks Spark cluster, which component is primarily responsible for scheduling tasks and managing the distribution of computation across the worker nodes?

A.The Executor
B.The Cluster Manager
C.The Driver
D.The Worker Node
AnswerC

The Driver maintains the state of the Spark application. It is responsible for analyzing, distributing, and scheduling tasks across the executors. By managing the Directed Acyclic Graph (DAG) and monitoring executor health, the driver ensures that jobs are executed efficiently and in the correct order based on stage dependencies.

Why this answer

The Driver process is the heart of the Spark application. It hosts the SparkContext, which interacts with the cluster manager to request resources. Once resources are allocated, the Driver converts the logical plan into a physical execution plan, splitting the job into stages and tasks, which are then distributed to the Executors.

Understanding this architecture is critical for debugging performance bottlenecks related to task scheduling and memory management in distributed environments.

Exam trap

Students often mistake the Cluster Manager (like YARN or K8s) for the task scheduler. While the manager allocates resources, the Driver performs the specific task scheduling and DAG management.

280
MCQeasy

Which command is used to display the logical and physical execution plans for a given Spark SQL query?

A.SHOW PLAN
B.DESCRIBE QUERY
C.EXPLAIN SELECT * FROM table
D.DEBUG SELECT * FROM table
AnswerC

The 'EXPLAIN' command is the standard way to inspect Spark's query execution plans. It helps developers understand how the Catalyst optimizer parses and plans the query, allowing them to identify bottlenecks, such as unnecessary shuffles or full table scans, before executing the actual data-intensive operations on the large dataset.

Why this answer

The 'EXPLAIN' command is the essential tool for inspecting the Catalyst query optimizer's work. It provides a breakdown of the logical plan (the raw representation of the query), the analyzed plan, the optimized plan, and the physical plan (how Spark will actually execute the tasks). Understanding this output is crucial for performance tuning, as it reveals how Spark handles joins, filters, and projections before the query actually runs on the distributed cluster nodes.

Exam trap

Test-takers frequently confuse the EXPLAIN command with DESCRIBE or SHOW commands, failing to realize that EXPLAIN specifically outputs the Catalyst optimizer's execution plans.

281
MCQhard

What is the consequence of having 'wide dependencies' in a Spark job regarding the Spark Architecture?

A.It allows for task pipelining within a single stage.
B.It forces the DAG scheduler to create a new stage boundary.
C.It enables data locality optimizations automatically.
D.It reduces the total number of tasks in the job.
AnswerB

Because wide dependencies involve shuffles, the DAG scheduler must finish all parent tasks before starting the child stage. This mandatory barrier allows Spark to guarantee that all intermediate data is successfully written and available to the next set of tasks across the cluster's network nodes.

Why this answer

Wide dependencies occur when a partition in the parent RDD contributes to multiple partitions in the child RDD, necessitating a shuffle. This forces a boundary in the DAG, triggering the creation of a new stage. Understanding this is essential because every shuffle introduces network I/O, disk I/O, and serialization overhead, which are the primary performance costs that developers must minimize when designing efficient Spark transformations and data models.

Exam trap

Candidates often assume wide dependencies are always bad and should be avoided at all costs. They miss the fact that they are sometimes necessary for operations like joins and aggregations.

282
MCQmedium

Which TWO of the following are benefits of using Delta Lake over standard Parquet files in Spark SQL?

A.Support for ACID transactions.
B.Ability to perform schema evolution and enforcement.
C.Faster raw I/O performance for single-column reads.
D.Automatic conversion of CSV files to optimized Parquet.
E.Compatibility with legacy Hive Metastore versions without modification.
AnswerA, B

ACID transactions ensure data integrity by allowing multiple concurrent readers and writers to interact with the data without corruption. This is a critical feature for data pipelines where consistency is paramount, and standard Parquet files do not provide this level of transactional guarantee by default.

Why this answer

Delta Lake extends the capabilities of standard Parquet by adding a transaction log and metadata layer. This enables ACID transactions, Time Travel, and schema enforcement, which are critical for robust data engineering. For a Databricks certified developer, knowing why Delta is the preferred format for the Lakehouse architecture is essential for building scalable, reliable, and maintainable data systems that surpass the limitations of raw file-based storage.

Exam trap

Candidates often mistakenly select performance-related features like 'automatic indexing' or 'auto-scaling' as benefits of Delta Lake, confusing general cloud platform capabilities with the specific ACID and schema features provided by the Delta format.

283
MCQmedium

A data engineer runs a Spark job on a Databricks cluster using the default FIFO scheduler. They notice that a long-running job is holding all cluster resources, and short ad-hoc queries submitted later are stuck waiting. The engineer wants to allow concurrent scheduling of multiple jobs within the same Spark application so that short jobs can run while the long job is still executing. Which Spark configuration should be set to enable this behavior?

A.spark.scheduler.mode=FAIR
B.spark.speculation=true
C.spark.dynamicAllocation.enabled=true
D.spark.scheduler.maxRegisteredResourcesWaitingTime=30s
AnswerA

Setting spark.scheduler.mode to FAIR enables the Fair Scheduler, which allows multiple jobs within the same Spark application to share cluster resources more evenly. This prevents a single long-running job from monopolizing all slots, so shorter jobs can start and complete without waiting for the long job to finish. This directly addresses the scenario.

Why this answer

The Fair Scheduler (spark.scheduler.mode=FAIR) enables multiple jobs within the same Spark application to share resources, so shorter jobs can run concurrently with a long-running job. The other options address resource scaling, registration timeouts, or straggler mitigation, none of which change intra-application job scheduling to allow concurrent execution.

Exam trap

The trap here is confusing dynamic allocation with job scheduling, assuming that adding more executors will automatically allow concurrent jobs, when the real issue is the FIFO scheduling policy.

284
MCQmedium

A developer is migrating a local pandas script to the Pandas API on Spark. The dataset is large and partitioned across many executors. The developer executes a custom row-wise operation using a standard Python lambda function inside a `.apply()` method without specifying return types or using vectorized operations. Why might this approach cause performance degradation in Databricks?

A.Spark automatically converts all pandas apply operations into native GPU-accelerated C++ code during the initial logical plan compilation phase.
B.The Catalyst optimizer completely bypasses the execution plan, forcing the cluster to fall back to a single-threaded local driver execution model.
C.Iterating through rows via Python lambdas forces high data serialization overhead between JVM and Python workers, destroying vectorized execution benefits.
D.Pandas API on Spark strictly prohibits the use of the `.apply()` method and immediately throws a compilation error during execution.
AnswerC

Row-wise Python functions require Python to deserialize every single record from JVM memory, process it individually, and serialize it back. This completely bypasses Apache Spark's tungsten memory management and columnar vectorization, leading to extreme network and CPU bottlenecks.

Why this answer

Standard pandas `.apply()` functions often execute row-by-row python processing rather than leveraging native Catalyst query optimizations. When using non-vectorized operations in Spark without explicit type hints, the engine must serialize data between JVM and Python workers repeatedly. This causes high serialization overhead, defeats distributed columnar optimization, and ultimately results in severe performance degradation compared to vectorized Spark expressions.

Exam trap

Candidates often assume standard pandas functions will automatically run efficiently in Spark. They fail to realize that row-wise lambda functions force expensive serialization between the JVM and Python.

285
Multi-Selectmedium

A developer is using Spark SQL to analyze a DataFrame that contains a column named tags, which holds an array of strings for each row. The developer needs to filter rows where the array contains the string 'urgent' and also produce a new column with the number of elements in the array. Which TWO Spark SQL expressions should be used in the query? (Choose two.)

Select 2 answers
A.ARRAY_MAX(tags)
B.ARRAY_CONTAINS(tags, 'urgent')
C.EXPLODE(tags)
D.COLLECT_LIST(tags)
E.SIZE(tags)
AnswersB, E

ARRAY_CONTAINS is the correct function to test whether an array column includes a specific value. It returns a boolean and works directly on array-typed columns in Spark SQL. In a WHERE clause, it filters rows where the tags array contains 'urgent'. This is the idiomatic and efficient way to perform membership tests on arrays without exploding them first, preserving row granularity.

Why this answer

ARRAY_CONTAINS provides a direct boolean test for whether an array includes a given value, making it ideal for filtering rows without altering their structure. SIZE returns the element count of an array, which satisfies the need for a new column with the number of tags. Together, they allow the query to filter and augment the DataFrame efficiently.

Other functions like EXPLODE change row cardinality, COLLECT_LIST aggregates, and ARRAY_MAX finds a maximum, none of which meet the specific goals.

Exam trap

The trap here is confusing array inspection functions with generator or aggregate functions, leading to incorrect row-level results or unnecessary query complexity.

286
MCQeasy

In the context of the Spark Driver, what is the 'DAG' and why is it important?

A.It is a physical storage format for saving Spark data.
B.It is a graph representing dependencies between stages of computation.
C.It is a configuration file that defines cluster resources.
D.It is a component that manages user security and access logs.
AnswerB

The DAG captures the lineage of transformations. Each node represents a transformation, and edges represent dependencies. This structure allows the Spark scheduler to optimize the execution by grouping stages and re-running only the necessary parts of the graph in case of failure, ensuring fault tolerance and efficient resource usage.

Why this answer

The Directed Acyclic Graph (DAG) represents the sequence of transformations applied to data. Spark builds this graph to optimize query plans before execution. Understanding this is essential because the DAG defines how stages are partitioned and executed.

If a developer understands the DAG, they can better structure their code—for example, by using narrow transformations instead of wide ones—to minimize shuffling and improve overall application performance by creating a more efficient execution path.

Exam trap

Candidates mistakenly describe the DAG as a physical data movement plan or a cache of data, rather than a logical dependency graph of transformations used for optimization.

287
MCQmedium

A data engineer is building a Spark SQL pipeline that must return the top 3 highest-paid employees within each department from a Delta table named `employees` with columns `dept`, `name`, and `salary`. The engineer wants a single query that produces one row per qualifying employee, ranked by salary descending within each department, without collapsing rows. Which approach should be used?

A.Use `GROUP BY dept` with the `MAX(salary)` aggregate and a `HAVING` clause limiting results to three rows.
B.Use `ORDER BY salary DESC` on the full table and apply `LIMIT 3`.
C.Use the `rank()` window function partitioned by `dept` and ordered by `salary DESC`, then filter on the rank column.
D.Use `DISTINCT` on `dept` and `salary`, then sort the result with `SORT BY salary DESC`.
AnswerC

A window function with `PARTITION BY dept ORDER BY salary DESC` assigns a rank within each department without collapsing rows, and filtering on the rank column keeps only the top three per department. This is the idiomatic Spark SQL pattern for top-N-per-group and is fully supported in Databricks SQL warehouses and clusters.

Why this answer

Window functions are the correct tool for top-N-per-group problems because they compute a value across a set of rows related to the current row while preserving all rows. Partitioning by department and ordering by salary descending yields a per-department rank, and filtering that rank to three returns exactly the desired rows in a single query.

Exam trap

The trap here is assuming that a global `ORDER BY ... LIMIT` or a `GROUP BY` aggregate can satisfy a per-group top-N requirement, when only a window function preserves row granularity while ranking within partitions.

288
MCQeasy

A developer notices that a Spark DataFrame job on Databricks is running slowly and the Spark UI shows that many tasks are reading from a Delta table with a large number of small files. The job performs a filter on a date column and then aggregates results. Which optimization technique will most directly improve read performance in this scenario?

A.Increase spark.sql.files.maxPartitionBytes to read more data per task.
B.Run OPTIMIZE on the Delta table to compact small files into larger ones.
C.Set spark.sql.autoBroadcastJoinThreshold to -1 to disable broadcast joins.
D.Enable Databricks Delta Cache to cache the table files on local SSDs.
AnswerB

OPTIMIZE compacts small files into larger, more efficient files, reducing the number of file opens and improving read throughput. This directly addresses the small files problem, which is a common cause of slow reads in Delta Lake. After compaction, the job will read fewer files, leading to faster scan and aggregation.

Why this answer

The small files problem in Delta Lake causes slow reads because each file requires a separate open and read operation. Compacting small files into larger ones with OPTIMIZE reduces the number of files, improving I/O efficiency. This is the most direct fix for the described symptom of many tasks reading small files, leading to faster filter and aggregation performance.

Exam trap

The trap here is thinking that caching or increasing partition size will fix slow reads caused by many small files, when the core issue is file count and compaction is needed.

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

290
MCQmedium

What is the function of the 'Executor' within the Spark execution model?

A.It acts as the central coordinator for all worker nodes.
B.It executes tasks assigned by the Driver and caches data.
C.It initializes the SparkSession for the user's application.
D.It defines the logical plan of the Spark application.
AnswerB

The executor is a JVM process running on a worker node that executes the code dispatched by the Driver. It manages its own local memory for caching and processing, reports task status back to the Driver, and ensures that the assigned tasks are completed using the allocated CPU and memory resources.

Why this answer

The executor is the component that performs the actual computation and stores data. It is important to know that executors are distributed entities—they carry out the work assigned by the Driver. By understanding that executors run tasks in parallel and manage the cache, developers can configure their cluster appropriately for the memory and compute needs of their specific data processing pipelines, ensuring high throughput and efficient resource utilization throughout the application's lifecycle.

Exam trap

Many candidates believe the Driver executes the actual data processing tasks, confusing job planning responsibilities with worker-level computation and caching.

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

292
MCQmedium

Refer to the exhibit. Which of the following is the most likely cause for this 'shuffle fetch failure' in a Databricks cluster?

A.The Driver ran out of memory during task scheduling.
B.The map-side executor was terminated before the reduce-side task could fetch its data.
C.The transformation involves a narrow dependency.
D.The input data format is unsupported.
AnswerB

When an executor terminates before its shuffle files are fetched, the reduce task receives a fetch failure. This commonly happens in dynamic environments where nodes are reclaimed. Implementing a persistent shuffle service or adjusting task retry policies can mitigate this issue by ensuring data availability for downstream tasks.

Why this answer

A shuffle fetch failure usually occurs when an executor loses the intermediate data that a downstream task is trying to read, often due to the executor being preempted or failing. In Databricks, this is a common symptom when autoscaling removes a node that held required shuffle files. Understanding this error is essential for debugging job stability and configuring cluster settings to ensure reliable data processing despite dynamic resource availability.

Exam trap

Candidates often blame network congestion or code errors for fetch failures. They fail to consider that autoscaling or preemptible instances in Databricks are the most frequent causes of lost shuffle files.

293
MCQeasy

Which of the following describes the 'Driver' process in a Spark application?

A.It executes tasks concurrently on distributed worker nodes.
B.It hosts the SparkContext and manages job scheduling.
C.It is responsible for storing data in a distributed cache.
D.It replaces the Cluster Manager to handle hardware resources.
AnswerB

The Driver process hosts the SparkContext, which is the main entry point for the Spark application. It manages the lifecycle of the application, translates the user's code into a DAG of tasks, and orchestrates the distribution of these tasks to the worker nodes for concurrent execution, serving as the central coordinator.

Why this answer

The Driver is the brain of the Spark application. It is the first point of contact and maintains the application state. Knowing that the Driver holds the SparkContext and manages the DAG is crucial, as this explains why high-latency tasks or memory-intensive operations on the Driver can degrade performance.

Developers must keep logic on the Driver lightweight to ensure that the application remains responsive and capable of coordinating work effectively across all distributed worker nodes.

Exam trap

Candidates often incorrectly attribute task execution or data storage to the Driver, forgetting that the Driver is strictly a control plane entity, not a worker node.

294
MCQmedium

A data engineer has a Pandas-on-Spark DataFrame `psdf` with a column `event_ts` stored as string timestamps. They run `psdf['event_ts'].astype('datetime64[ns]')` and then call `.dt.hour` on the resulting Series. In a Databricks notebook, what is the result of this operation?

A.The conversion and `.dt.hour` execute as Spark expressions, returning a new Pandas-on-Spark Series without collecting data to the driver.
B.The operation succeeds only if `spark.sql.execution.arrow.pyspark.enabled` is set to true; otherwise it falls back to a Python UDF that may fail on null timestamps.
C.The operation triggers an immediate collect of all rows to the driver, converts them to pandas, computes the hour locally, and returns a pandas Series.
D.The `.dt` accessor is unsupported in Pandas API on Spark, so the call raises an AttributeError before any Spark job is launched.
AnswerA

Pandas API on Spark implements `.astype('datetime64[ns]')` and the `.dt` accessor as distributed Spark column expressions. The string-to-timestamp cast maps to Spark's cast to TimestampType, and `.dt.hour` maps to the `hour` function, so no driver collection occurs. The result remains a Pandas-on-Spark Series backed by a Spark plan, which is the expected behavior for supported datetime operations.

Why this answer

Pandas API on Spark translates supported pandas operations into Spark logical plans. Casting a string column to datetime64 and accessing `.dt.hour` are both supported and map to Spark's timestamp cast and hour extraction. The result stays distributed as a Pandas-on-Spark Series, and no driver collection is triggered.

This preserves scalability and aligns with the library's goal of providing pandas-like syntax over Spark execution.

Exam trap

The trap here is assuming that any pandas operation on a Pandas-on-Spark object forces local execution, when supported datetime accessors are actually translated into distributed Spark expressions.

295
MCQmedium

You are debugging a PySpark DataFrame job on Databricks that performs multiple transformations and actions on a large delta table. You notice that the execution plan shows redundant computations where the same upstream DataFrame is evaluated repeatedly. Which transformation should you apply to optimize this workflow and avoid recomputing the upstream lineage?

A.Call df.broadcast() on the DataFrame before each downstream join operation.
B.Call df.repartition() to distribute the data evenly across partitions before execution.
C.Call df.cache() or df.persist() before referencing the DataFrame in multiple downstream actions.
D.Call df.coalesce() to reduce the number of shuffle partitions prior to the final write operation.
AnswerC

Caching serializes or stores the evaluated DataFrame partitions in memory or disk storage. When subsequent actions trigger execution, Spark retrieves the cached blocks directly instead of re-evaluating the entire upstream transformation DAG from source files.

Why this answer

Persisting or caching a DataFrame instructs Spark to store intermediate results in memory or disk across actions, preventing expensive recomputation of the entire lineage graph. This optimization is critical for iterative algorithms or workflows where a single DataFrame is referenced multiple times downstream.

Exam trap

Candidates often confuse caching with broadcasting or assume Spark automatically caches all reused DataFrames. Spark does not evaluate reference frequency and will recompute lineage unless explicitly told to persist.

Page 3

Page 4 of 4

All pages