Courseiva

Databricks Certified Associate Developer for Apache Spark (Databricks-Spark-Assoc) — Questions 151–225

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

Page 2

Page 3 of 4

Page 4
151
MCQmedium

You are analyzing a large dataset using Pandas API on Spark. You have a Pandas-on-Spark DataFrame `psdf` that was created from a Spark DataFrame with multiple partitions. You call `psdf.head(10)` to quickly inspect the data. What does this operation return?

A.A Pandas-on-Spark DataFrame containing the first 10 rows.
B.A standard pandas DataFrame containing the first 10 rows.
C.A list of Row objects containing the first 10 rows.
D.A Spark DataFrame containing the first 10 rows.
AnswerB

This is correct. In Pandas API on Spark, `head(n)` collects the first n rows from the distributed DataFrame and returns them as a local pandas DataFrame. This allows you to inspect a small sample without triggering a full distributed computation. The operation is efficient because it only processes the necessary partitions and brings a limited amount of data to the driver.

Why this answer

The `head(n)` method in Pandas API on Spark returns a local pandas DataFrame containing the first n rows. It is intended for quick inspection of data, and it collects only the necessary rows to the driver. This behavior matches the pandas API, where `head` returns a DataFrame, but here it triggers a small action to bring data locally.

Exam trap

The trap here is assuming that `head` returns a distributed DataFrame, but it actually returns a local pandas DataFrame for immediate inspection.

152
MCQmedium

A Spark application reads a large CSV file and performs a series of transformations, including a filter and a join with a small lookup table. The job is running slowly, and the Spark UI shows that the CSV parsing stage is taking a long time. Which action would most improve performance?

A.Increase the number of partitions when reading the CSV by setting a lower spark.sql.files.maxPartitionBytes.
B.Set spark.sql.autoBroadcastJoinThreshold to -1 to disable broadcast join.
C.Convert the CSV to Parquet and read the Parquet file instead.
D.Cache the CSV DataFrame immediately after reading.
AnswerC

Parquet is a columnar format that supports predicate pushdown and column pruning, which drastically reduces I/O and parsing overhead compared to CSV. CSV parsing is row-based and requires reading and parsing every field, even if only a few columns are needed. Converting to Parquet allows Spark to read only the necessary columns and skip irrelevant data, significantly speeding up the initial read and subsequent transformations.

Why this answer

Converting the CSV to Parquet addresses the slow CSV parsing by leveraging Parquet's columnar storage, which enables column pruning and predicate pushdown. This reduces the amount of data read and parsed, directly improving the performance of the initial stage. The other options either do not target the parsing bottleneck or could worsen performance.

Exam trap

The trap here is assuming that increasing partitions or caching will fix slow CSV parsing, when the format itself is the primary bottleneck.

153
Multi-Selectmedium

A developer is using Spark on Databricks and notices that a particular job has many stages due to shuffle operations. They want to understand the role of the shuffle in the Spark execution model. Which two statements accurately describe the behavior of a shuffle operation in Spark? (Choose two.)

Select 2 answers
A.A shuffle eliminates the need for the driver to coordinate task scheduling across stages.
B.A shuffle writes intermediate data to disk on the map side and reads it over the network on the reduce side.
C.A shuffle creates a stage boundary, splitting the job into a map stage and a reduce stage.
D.A shuffle always results in exactly one partition per reducer task, regardless of the number of map tasks.
E.A shuffle is only triggered by actions, not by transformations.
AnswersB, C

This is correct. During a shuffle, map tasks write shuffle files to local disk (or memory if configured) and reduce tasks fetch these files over the network. This disk I/O and network transfer make shuffles expensive. Spark's shuffle manager (e.g., SortShuffleManager) handles this process, and the data is partitioned by key before being written, ensuring that all records for a given key end up on the same reducer.

Why this answer

A shuffle operation writes intermediate data to disk on the map side and transfers it over the network to reduce tasks, creating a stage boundary. These two characteristics are central to understanding why shuffles are expensive and how Spark's DAG is structured. The number of reduce partitions is configurable and not fixed to one per reducer, and shuffles are triggered by transformations, not actions.

Exam trap

The trap here is assuming that a shuffle guarantees a one-to-one mapping between map tasks and reduce partitions, or that actions directly cause shuffles.

154
MCQeasy

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

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

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

Why this answer

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

Exam trap

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

155
MCQeasy

You are monitoring a Structured Streaming query in Databricks and want to see the current status, including the number of input rows per second and the batch duration. Which of the following is the most direct way to access this information?

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

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

Why this answer

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

Thus, `lastProgress` is the correct choice.

Exam trap

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

156
MCQeasy

You are writing a Structured Streaming query that reads from a Kafka topic and outputs to the console. You want to see only the newly arrived data in each micro-batch, without aggregations. Which output mode should you use?

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

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

Why this answer

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

Exam trap

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

157
MCQmedium

A data engineer runs a PySpark job on a Databricks cluster. The job reads a 500 GB Parquet dataset, applies a filter, and writes the result. The engineer notices that during execution, all tasks of a particular stage complete quickly except for a handful that take far longer, and the Spark UI shows these tasks are processing partitions that contain far more records than others. Which Spark architecture concept best explains this behavior, and what is the most appropriate remediation?

A.Insufficient executor memory causing garbage collection pauses only on certain tasks; the engineer should increase spark.executor.memory.
B.Data skew across partitions during a shuffle; the engineer should apply salting or repartitioning to distribute records more evenly.
C.The DAG Scheduler is serializing stages incorrectly; the engineer should disable adaptive query execution to force static stage boundaries.
D.Too few partitions in the source Parquet files; the engineer should call coalesce(1) before writing the output.
AnswerB

Uneven record distribution across partitions after a shuffle causes a few tasks to process disproportionately large partitions, producing stragglers. Salting keys or repartitioning redistributes data so each task receives a comparable workload, directly addressing the imbalance observed in the Spark UI task duration metrics for this stage.

Why this answer

Skewed partition sizes after a shuffle produce a few long-running tasks while most finish quickly, which matches the Spark UI pattern described. Redistributing records through salting or repartitioning balances the workload across tasks, addressing the root cause rather than masking symptoms with memory tuning or reduced parallelism.

Exam trap

The trap here is assuming slow tasks always indicate a memory shortage, when uneven partition sizes from a shuffle are the more likely cause in this scenario.

158
MCQmedium

A developer has a PySpark DataFrame `df` with columns `order_id`, `customer_id`, and `order_ts` (timestamp). They need to return only the most recent order per customer, keeping all original columns, and they want to avoid a self-join or a manual sort-then-dropDuplicates approach. Which DataFrame operation should they use?

A.df.dropDuplicates(["customer_id"])
B.df.groupBy("customer_id").max("order_ts")
C.Use a Window partitioned by `customer_id` ordered by `order_ts` descending, add `row_number()`, then filter for row number equal to 1
D.df.orderBy("order_ts", ascending=False).limit(1)
AnswerC

A Window partitioned by `customer_id` and ordered by `order_ts` descending assigns rank 1 to the latest order per customer. Adding `row_number().over(window)` and filtering `row_number == 1` returns the full original row, preserving all columns. This is the idiomatic PySpark replacement for a self-join or sort-then-dedup pattern and scales with partitioning.

Why this answer

Selecting the latest row per group while preserving all columns is exactly what a Window with `row_number()` solves. Partitioning by the grouping key and ordering by the timestamp descending makes the most recent record rank 1, and filtering on that rank returns complete rows. Aggregations drop columns, global limits collapse groups, and dropDuplicates cannot express ordering.

Exam trap

The trap here is assuming that `dropDuplicates` on a grouping column keeps the latest record, when in fact it keeps an arbitrary row and ignores any ordering.

159
MCQmedium

A developer is debugging a Spark job and observes that a particular stage has 200 tasks, but only 10 executors with 2 cores each are available. What will happen to the remaining tasks in that stage?

A.The tasks will fail with an insufficient resources error
B.The tasks will be executed on the driver node to compensate for the lack of executors
C.The tasks will be queued and executed as cores become available, up to 20 concurrently
D.The tasks will be automatically coalesced to match the number of available cores
AnswerC

With 10 executors each having 2 cores, the cluster can run 20 tasks concurrently. The Task Scheduler will launch tasks as slots free up, so the 200 tasks are processed in waves. This is standard behavior: the number of concurrent tasks is limited by total cores, and the rest wait in the queue until resources are available.

Why this answer

The number of concurrently running tasks is bounded by the total number of cores across executors. Here, 10 executors with 2 cores each provide 20 slots, so 20 tasks run at a time while the remaining 180 wait in the scheduler's queue. As tasks finish, new ones are launched.

Spark does not fail or coalesce tasks due to resource scarcity.

Exam trap

The trap here is thinking that Spark dynamically adjusts partition count to fit available cores, when in reality it queues tasks and runs them in waves based on core availability.

160
Multi-Selectmedium

A data engineer has a PySpark DataFrame `events` with columns `user_id`, `event_time` (timestamp), and `payload` (string). They must produce a new DataFrame where each row is enriched with the `payload` value from the user's immediately preceding event, ordered by `event_time`, without collapsing any rows. Which TWO approaches accomplish this? (Choose two.)

Select 2 answers
A.Use `df.orderBy('event_time').rdd.zipWithIndex()` and manually look up the previous index per user in a driver-side dictionary.
B.Create a Window partitioned by `user_id` and ordered by `event_time`, then use `lag('payload', 1).over(windowSpec)` inside `withColumn`.
C.Apply `window('user_id', 'event_time')` with `lag('payload')` and rely on Spark to fill nulls for the first event.
D.Use `df.withColumn('prev_payload', lag('payload').over(Window.partitionBy('user_id').orderBy('event_time')))`.
E.Call `groupBy('user_id').agg(collect_list('payload'))` and then explode the resulting list back into rows.
AnswersB, D

This is correct because a Window partitioned by `user_id` and ordered by `event_time` restricts `lag` to each user's own timeline, and `lag('payload', 1)` returns the prior row's payload while preserving every row. `withColumn` adds the result as a new column without aggregating, which satisfies the requirement of not collapsing rows.

Why this answer

Enriching rows with a prior value without collapsing requires a Window function. Partitioning by `user_id` and ordering by `event_time` ensures the 'previous' row is the same user's immediately earlier event, while `lag` returns the payload from that row and preserves row count. Both `lag(...).over(windowSpec)` and the inline `Window.partitionBy(...).orderBy(...)` expression produce identical semantics; the alternatives either collapse rows, use invalid syntax, or lose per-user ordering.

Exam trap

The trap here is assuming that any 'previous row' lookup works globally, when Window functions must be partitioned by the entity whose sequence you care about.

161
MCQmedium

You need to combine two datasets in Spark SQL, retaining all records from the left table and matching records from the right table, while filling unmatched right columns with null values. Which join type should you use?

A.RIGHT OUTER JOIN
B.FULL OUTER JOIN
C.LEFT OUTER JOIN
D.INNER JOIN
AnswerC

Preserving all rows from the primary left table and augmenting them with matching values from the right table defines the left outer join. Unmatched columns are automatically padded with null values, exactly meeting the business requirement.

Why this answer

A left outer join ensures every single record from the left dataset appears in the final result set, regardless of whether a matching key exists in the right dataset. Missing matches on the right side are populated with null values, maintaining structural integrity for downstream analysis.

Exam trap

Candidates frequently confuse LEFT OUTER JOIN with RIGHT OUTER JOIN or FULL OUTER JOIN, accidentally swapping the order of tables or including unwanted nulls from the left side.

162
MCQeasy

A developer is using Pandas API on Spark and wants to convert a Pandas-on-Spark DataFrame `psdf` back to a standard pandas DataFrame for local analysis. Which method should they use?

A.`psdf.collect()`
B.`psdf.toPandas()`
C.`psdf.to_pandas()`
D.`psdf.toDF()`
AnswerB

`toPandas()` is the correct method to convert a Pandas-on-Spark DataFrame to a standard pandas DataFrame. It collects all data to the driver node and constructs a pandas DataFrame. This is suitable for small datasets that fit in memory, but can cause out-of-memory errors for large datasets.

Why this answer

The `toPandas()` method collects the distributed data into a single pandas DataFrame on the driver. It is the standard way to convert from Pandas API on Spark to local pandas. However, it should be used cautiously with large datasets to avoid memory issues.

Exam trap

The trap here is confusing `toPandas()` with other conversion methods like `toDF()` or `collect()`, which do not produce a pandas DataFrame.

163
MCQmedium

A developer is tuning a Databricks job and wants to know how many tasks will be created for the final stage of a job that reads a Parquet file with 200 partitions, applies a filter, and then calls coalesce(10) before writing the result. Assuming no other repartitioning or shuffles occur, how many tasks will the final write stage contain?

A.1 task, because coalesce always merges everything into a single partition
B.200 tasks, because the filter forces a shuffle before coalesce
C.200 tasks, because coalesce does not change the number of partitions
D.10 tasks, because coalesce(10) reduces the partition count to 10
AnswerD

coalesce(10) collapses the 200 input partitions into 10 output partitions using narrow dependencies, avoiding a shuffle. Each partition corresponds to one task in the stage that writes the data, so the final stage launches exactly 10 tasks. This matches the intent of coalesce: reducing partition count efficiently without redistributing data across the cluster.

Why this answer

coalesce(10) reduces the RDD from 200 partitions to 10 using narrow dependencies, so no shuffle occurs. Since the number of tasks in a stage equals the number of partitions in the RDD being processed, the final write stage launches 10 tasks. Filtering is also narrow and does not alter the partition count, so the coalesce result directly determines the task count.

Exam trap

The trap here is confusing coalesce with repartition, or assuming that filter triggers a shuffle, when in fact coalesce only merges partitions without a full shuffle and filter is narrow.

164
MCQhard

A Spark Structured Streaming job on Databricks reads from a Delta table and writes micro-batches to another Delta table with a 30-second trigger. After several hours, the batch duration grows from 4 seconds to over 60 seconds and the job falls behind. The source table is compacted regularly, and the cluster has enough CPU. Which tuning action is most likely to restore the original batch duration?

A.Increase spark.sql.shuffle.partitions to a very high value so each micro-batch uses more tasks and finishes faster.
B.Reduce the number of shuffle partitions and disable adaptive query execution so the streaming plan stays deterministic across micro-batches.
C.Enable Delta Lake optimized writes and tune the trigger interval so the job processes larger, less frequent micro-batches.
D.Check for stateful aggregation or deduplication without a watermark and add a watermark with a bounded state cleanup so state does not grow unbounded.
AnswerD

Unbounded state from a stateful operation without a watermark causes each micro-batch to scan an ever-growing state store, which steadily increases batch duration and can make the job fall behind. Adding a watermark lets Spark drop state older than the watermark, keeping state bounded and batch times stable. This directly addresses a gradual, hours-long degradation pattern.

Why this answer

A steady increase in batch duration over hours with a compacted source and adequate CPU strongly indicates unbounded state growth in a stateful streaming operation. Without a watermark, Spark retains all keys indefinitely, so each micro-batch reads and writes a larger state store. Adding a watermark with a bounded delay lets Spark evict old state, keeping per-batch work roughly constant and restoring stable batch times.

Exam trap

The trap here is treating a slowly growing streaming batch time as a parallelism problem and adding shuffle partitions instead of investigating state accumulation.

165
Multi-Selecthard

You are using Spark SQL to join two large Delta tables, orders and customers, on a common column customer_id. The orders table is partitioned by order_date, and the customers table is not partitioned. You need to ensure the join is efficient and minimizes shuffling. Which TWO actions should you take? (Choose two.)

Select 2 answers
A.Enable adaptive query execution (AQE) and set spark.sql.adaptive.enabled to true.
B.Use the MERGE command to combine the tables.
C.Partition the customers table by customer_id to match the orders table.
D.Broadcast the customers table if it is small enough to fit in memory.
E.Repartition both tables by customer_id before the join.
AnswersA, D

AQE can dynamically optimize joins by converting sort-merge joins to broadcast joins if one side is small after initial stages, and by coalescing partitions. Enabling AQE helps minimize shuffling and improves performance. It is a best practice for large joins in Spark SQL, especially when statistics are outdated or data skew exists.

Why this answer

Broadcasting the smaller table eliminates shuffling of the larger table. Enabling AQE allows Spark to dynamically optimize the join, potentially converting to a broadcast join or adjusting partitions. These two actions together minimize shuffling and improve efficiency.

Exam trap

The trap here is assuming that repartitioning both tables is always beneficial, or that partitioning the smaller table by the join key will help, when in fact it can cause more shuffling or small file problems.

166
MCQeasy

In Spark SQL, what is the primary difference between a temporary view and a global temporary view?

A.Temporary views persist after the cluster is terminated, while global temporary views do not.
B.Temporary views are accessible only within the current session, while global temporary views are accessible across all sessions.
C.Temporary views are written to the Hive Metastore, while global temporary views are only in memory.
D.Global temporary views are faster to query because they are cached by default.
AnswerB

Temporary views are bound to the specific Spark session that created them, ensuring isolation. Global temporary views are bound to the Spark application and are accessible by any session on the same cluster, allowing shared access to temporary datasets across different notebooks or users running on the same cluster.

Why this answer

Understanding scope is critical for managing data access in multi-user Databricks environments. Temporary views are session-specific, ensuring data isolation between users. Global temporary views are visible to all sessions on the cluster and are stored in the 'global_temp' database.

This distinction is vital for maintaining security and avoiding namespace collisions when multiple data engineers are working simultaneously on the same shared cluster infrastructure.

Exam trap

Candidates often incorrectly believe global temporary views are persistent across cluster restarts or shared across different workspaces, rather than just being session-agnostic within the same cluster.

167
Multi-Selectmedium

Which TWO factors influence the effective parallelism of a Spark application?

Select 2 answers
A.The number of RDD or DataFrame partitions.
B.The total amount of disk space on the worker nodes.
C.The number of available CPU cores in the executor pool.
D.The network latency between the driver and the cluster manager.
E.The size of the broadcast variables.
AnswersA, C

Partitions are the fundamental unit of parallelism in Spark. Each partition corresponds to one task. Increasing the number of partitions allows for more concurrent tasks, assuming there is sufficient CPU capacity on the executors to process them simultaneously, directly affecting the job's execution speed.

Why this answer

Effective parallelism in Spark is governed by the number of partitions created in the data and the number of available cores in the executor pool. If the partition count is too low, the cluster is underutilized. If the core count is too low, tasks are queued, creating bottlenecks.

Managing these two factors is the primary way to ensure that the workload is spread evenly across the available hardware resources for maximum throughput.

Exam trap

Candidates often focus only on the number of partitions. They fail to realize that having many partitions is useless if there are insufficient CPU cores to process them in parallel.

168
MCQeasy

Which environment variable is mandatory to establish a connection to a Databricks cluster using Spark Connect in a local Python environment?

A.SPARK_MASTER
B.SPARK_REMOTE
C.DATABRICKS_HOST
D.SPARK_CONNECT_URL
AnswerB

SPARK_REMOTE is the primary configuration parameter for Spark Connect. It follows a specific format (sc://<workspace-url>:<port>;token=<token>;clusterId=<id>) that directs the client to the correct Databricks server. It is essential for establishing the gRPC channel required to transmit logical plans from the local machine to the cluster.

Why this answer

To connect to Databricks using Spark Connect, the SPARK_REMOTE environment variable must be set with the Databricks workspace URL and the compute resource identifier. This variable tells the SparkSession builder where to redirect the execution of commands. Without this properly formatted connection string, the Spark Connect client cannot authenticate or route the gRPC requests to the specific Databricks cluster intended for the computation.

Exam trap

Candidates often confuse SPARK_REMOTE with standard Spark configuration properties like spark.master. SPARK_REMOTE is specifically required for Spark Connect to establish the gRPC connection to the Databricks cluster.

169
MCQhard

A structured streaming DataFrame writes to a Delta table with a foreachBatch function that performs an upsert. After a cluster restart, the stream reprocesses some micro-batches and duplicate rows appear in the target table. The foreachBatch code already uses MERGE keyed on a unique id. Which change best prevents duplicates after restart?

A.Set the trigger to processingTime='1 minute' so micro-batches are larger and less likely to overlap.
B.Add a checkpointLocation to the writeStream options so the stream records its progress and resumes from the last committed offset.
C.Switch the output mode from append to complete so each micro-batch is fully replaced in the target table.
D.Add a deduplication step using dropDuplicates on the unique id inside the foreachBatch before the MERGE.
AnswerB

Structured streaming relies on the checkpoint location to persist progress, offsets, and state. Without it, a restarted query has no record of which micro-batches were committed and reprocesses data, producing duplicates even when the MERGE key is unique. Setting checkpointLocation gives exactly-once semantics for the sink when combined with an idempotent operation such as MERGE. This is the direct fix for reprocessing after restart.

Why this answer

Structured streaming tracks processed offsets and state in the checkpoint location. Without a checkpoint, a restarted query cannot know what it already committed and replays earlier micro-batches, creating duplicates even with an idempotent MERGE. Supplying checkpointLocation lets the query resume from the last committed offset, so each batch is processed once.

Combined with the existing MERGE on a unique id, this yields the intended exactly-once behavior on restart.

Exam trap

The trap here is believing a MERGE on a unique key alone guarantees exactly-once, when restart reprocessing is actually governed by the streaming checkpoint.

170
MCQmedium

You are processing a large dataset in Spark SQL and need to ensure that small files are avoided when writing data to Delta Lake. Which approach effectively minimizes small file generation during write operations?

A.Execute a DROP TABLE command before overwriting the existing table every time.
B.Increase the spark.sql.shuffle.partitions configuration to a very high value.
C.Enable 'autoOptimize' and 'optimizeWrite' at the Delta table level.
D.Use the 'repartition(1)' method on the DataFrame before writing to storage.
AnswerC

Enabling these properties allows Databricks to automatically coalesce small writes into larger files during the write operation itself. This significantly reduces the number of small files created by concurrent or frequent streaming writes, ensuring that data is laid out optimally for future analytical queries without manual intervention.

Why this answer

Optimizing file sizes is crucial for read performance and metadata management in Delta Lake. Using the OPTIMIZE command or enabling Auto Optimize are the standard patterns to address the small file problem. These techniques consolidate fragmented data into larger, performant files, reducing the overhead on the query engine and preventing performance degradation during subsequent read operations.

This is a fundamental skill for maintaining healthy, scalable data lakes on Databricks.

Exam trap

Candidates often suggest manual partitioning or repartitioning before every write. They fail to recognize that enabling built-in Delta Lake features like Auto Optimize is the more efficient, automated solution.

171
MCQeasy

A developer notices that a DataFrame transformation chain is executed twice: once for a count action used for logging and again for a write action. The source is a large Delta table and the repeated scan adds several minutes. Which action avoids the duplicate computation with the least risk?

A.Convert the DataFrame to a Pandas DataFrame for the count and then write from the original Spark DataFrame.
B.Call cache or persist on the transformed DataFrame before the count so the second action reuses the materialized data.
C.Set spark.sql.shuffle.partitions equal to the number of executor cores to speed up each execution.
D.Increase spark.sql.autoBroadcastJoinThreshold so more joins are broadcast and the plan becomes cheaper.
AnswerB

Caching the transformed DataFrame materializes it once, so the subsequent write action reads the cached partitions instead of recomputing the entire lineage from the Delta table. This directly eliminates the duplicate scan. It is a targeted change with predictable memory cost and no change to results. For a DataFrame reused across multiple actions, caching is the standard remedy.

Why this answer

When the same DataFrame lineage feeds two actions, Spark recomputes it for each action unless the intermediate result is materialized. Calling cache or persist after the expensive transformations stores the result so the later write reads it directly, removing the second full scan of the Delta table. The change is local, reversible, and does not alter results, making it the lowest-risk option among those presented.

Exam trap

The trap here is trying to tune shuffle or join settings for a problem that is actually caused by recomputing the same lineage across two separate actions.

172
MCQmedium

An analyst notices that a Spark SQL query filtering a Delta table with `WHERE order_date = '2024-03-15'` scans far more data than expected, even though the table is partitioned by `order_date`. The partition column was loaded as a string in `yyyy-MM-dd` format. Which explanation best accounts for the excessive scan?

A.The literal is compared against a partition column of a different type or format, so pruning cannot match partitions.
B.Partition pruning is disabled by default and must be enabled with `spark.sql.optimizer.enablePartitionPruning`.
C.The query lacks a `LIMIT` clause, which forces Spark to read every file in the table.
D.Delta tables never support partition pruning; only Parquet tables do.
AnswerA

If the partition column is stored as a string and the literal is cast or formatted differently, Catalyst may fail to derive a static partition predicate and fall back to scanning all partitions. Type or format mismatches between the filter literal and the partition column defeat pruning. Aligning the literal's type and format with the column restores efficient partition elimination.

Why this answer

Partition pruning depends on Catalyst proving that a filter predicate matches partition values, which requires the literal's type and format to align with the partition column. A string partition column compared to a differently typed or formatted literal prevents static pruning, causing a full scan. Correcting the literal's type and format lets the optimizer eliminate irrelevant partitions.

Exam trap

The trap here is blaming a missing configuration flag for poor pruning, when pruning is automatic and fails only when the predicate and partition column types or formats are incompatible.

173
MCQmedium

When joining two large tables in Spark SQL, which join type helps avoid expensive shuffles by loading a small table into memory on all executor nodes?

A.Shuffle Hash Join
B.Broadcast Hash Join
C.Sort Merge Join
D.Cartesian Product Join
AnswerB

Broadcast Hash Join optimizes performance by broadcasting the smaller table to all executors. Because the small table is available locally on every node, the join operation avoids large data shuffles, making it significantly faster for scenarios where one table can comfortably fit within the driver-defined broadcast memory threshold.

Why this answer

A Broadcast Hash Join is a highly efficient join strategy in Spark SQL. By marking the smaller table as a broadcast candidate, Spark sends a complete copy of that table to every executor node. This allows the join to be performed locally at each node, eliminating the need to shuffle the large table across the network, which is often the primary bottleneck in large-scale distributed join operations.

Exam trap

Candidates often select Sort-Merge Join or Shuffle Hash Join when trying to avoid network shuffles, forgetting that only broadcasting eliminates network movement for the small table.

174
MCQhard

A developer uses a Spark Connect session to create a temporary view with `df.createOrReplaceTempView("sales_v")` and then runs `spark.sql("SELECT * FROM sales_v")` in the same session. What is the scope of that temporary view?

A.It is registered in the server-side session catalog and is visible to operations within that session.
B.It is stored on the client machine and is visible only to the creating Python process.
C.It is written to the workspace metastore as a persistent view that survives cluster restarts.
D.It is registered globally on the cluster and is visible to all users and sessions.
AnswerA

When the client calls `createOrReplaceTempView`, the command is serialized and sent to the server, which registers the view in the session's catalog. Subsequent `spark.sql` calls in the same session can resolve the view because they execute against the same server-side session state. The view is not persisted beyond the session unless it is a global temporary view or a permanent object.

Why this answer

Creating a temporary view in Spark Connect sends a command to the server, which registers the view in the session catalog. Subsequent SQL in the same session resolves it because both operations share the same server-side session state. The view is not client-local, not cluster-global, and not persisted to the metastore; it vanishes when the session ends.

Exam trap

The trap here is assuming that a temporary view created through a Spark Connect client lives on the client, when it is actually registered in the server-side session catalog and scoped to that session.

175
MCQeasy

A developer wants to start a Structured Streaming query that reads from a Delta table and writes to another Delta table, and needs the query to process all existing data in the source table on its first run. Which option should be set on the read stream?

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

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

Why this answer

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

Exam trap

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

176
MCQeasy

A developer runs a Spark Connect client session against a Databricks cluster with `spark.conf.set("spark.sql.shuffle.partitions", "400")`. The cluster is configured with 8 worker nodes. Which component actually applies the shuffle partition setting to the physical plan?

A.The Databricks cluster-side Spark driver, which builds and executes the physical plan.
B.The local Spark Connect client process, which rewrites the plan before serialization.
C.The Databricks workspace control plane, which injects the value into the cluster's Spark configuration at startup.
D.The worker executors, which read the value from broadcast configuration during task scheduling.
AnswerA

With Spark Connect, the client sends unresolved logical plans and configuration over gRPC to the server, where the Spark driver performs analysis, optimization, and physical planning. The shuffle partition count is consumed during exchange planning on the driver, so the value takes effect against the 8-node cluster's execution. The client only declares the setting; the server enforces it.

Why this answer

Spark Connect splits responsibilities: the thin client builds unresolved logical plans and sends them, along with session configuration, to the server. The Databricks cluster-side Spark driver then analyzes, optimizes, and converts the plan into a physical plan, where settings such as the shuffle partition count influence exchange operators. Executors and the workspace control plane do not perform this planning step.

Exam trap

The trap here is assuming the local Spark Connect client performs planning or optimization, when in fact it only constructs and transmits unresolved logical plans and configuration to the server.

177
MCQmedium

What is the primary function of the Spark DAG Scheduler?

A.Managing low-level memory allocation for individual executors.
B.Converting logical transformation chains into stages of tasks.
C.Resource negotiation with the Cluster Manager.
D.Serializing data for network transmission between workers.
AnswerB

The DAG Scheduler analyzes the lineage of RDDs and DataFrames, identifying where shuffles occur. It groups operations that can be computed in parallel without data movement into stages. This allows Spark to build an efficient execution pipeline, reducing the need for disk I/O and increasing overall job speed.

Why this answer

The DAG Scheduler is responsible for translating the logical plan into a physical execution plan, specifically breaking the lineage into stages based on shuffle boundaries. By identifying wide dependencies that require data redistribution, it organizes the execution flow. This is fundamental for optimizing performance, as it minimizes data movement across the network by grouping together all narrow dependency transformations into a single executable stage before triggering a shuffle operation.

Exam trap

Candidates often confuse the DAG Scheduler with the Task Scheduler. They mistakenly believe the DAG scheduler manages low-level task execution on workers, rather than focusing on the logical stage boundary creation.

178
MCQeasy

A data engineer is configuring a Spark application on Databricks. They set `spark.executor.instances` to 4, `spark.executor.cores` to 5, and `spark.executor.memory` to 16g. The cluster has 5 worker nodes, each with 16 cores and 64 GB RAM. What is the maximum number of tasks that can run concurrently across all executors?

A.20
B.80
C.5
D.4
AnswerA

The maximum number of concurrent tasks equals the total number of executor cores in the cluster. With 4 executors and 5 cores per executor, the total is 4 × 5 = 20. Each core can run one task at a time, so up to 20 tasks can execute in parallel. This assumes no other resource constraints and that the cluster has enough capacity to launch all executors.

Why this answer

The number of concurrent tasks in a Spark application is determined by the total number of cores allocated to executors. With 4 executors and 5 cores each, the total is 20 cores, allowing up to 20 tasks to run simultaneously. This is a fundamental relationship in Spark's architecture: each core can process one task at a time, so the total task concurrency equals the sum of cores across all executors.

Exam trap

The trap here is multiplying the number of worker nodes by cores per executor instead of using the configured number of executors, or simply selecting the number of executors.

179
MCQhard

A data engineer observes that a Spark Structured Streaming job on Databricks processes micro-batches with steadily increasing latency over several hours. The Spark UI shows that the number of active tasks per batch stays constant, but each task processes a growing amount of state. Which architectural behavior explains this pattern?

A.Broadcast variables are being re-sent to executors on every micro-batch, adding network overhead.
B.The Driver is accumulating unbounded metadata from the DAG Scheduler across batches.
C.The Cluster Manager is throttling executor allocation, reducing parallelism per batch.
D.Stateful operators maintain growing keyed state in the executors as more keys arrive, increasing per-task work.
AnswerD

Stateful operations such as streaming aggregations or deduplication keep keyed state in executor memory and on disk. As new keys accumulate over hours, each task must read and update a larger state store, so per-task processing time rises even though the number of tasks stays constant, matching the observed latency growth.

Why this answer

Stateful streaming operators retain keyed state across micro-batches, and as distinct keys accumulate, each task must load, update, and write a larger state store. This raises per-task processing time while task parallelism stays fixed, which is exactly the pattern shown when active tasks are constant but per-task state keeps growing.

Exam trap

The trap here is attributing rising streaming latency to scheduling or resource allocation rather than to the growth of keyed operator state.

180
MCQhard

A Databricks job joins a 500 GB sales table with a 300 GB returns table on a customer_id key. A few customer_id values account for a large fraction of rows on both sides, and the job fails with executor OOM during the join. The developer wants to distribute the hot keys across more partitions without changing the query logic. Which technique should be used?

A.Broadcast the 300 GB returns table to every executor to avoid shuffling the skewed keys.
B.Enable Adaptive Query Execution with skew join handling so Spark splits skewed partitions at runtime.
C.Increase spark.sql.shuffle.partitions to 8000 to spread the skewed keys across more tasks.
D.Repartition both DataFrames by customer_id with 8000 partitions before the join.
AnswerB

AQE skew join handling detects partitions that are much larger than the median after the shuffle and splits them into smaller sub-partitions, each processed as a separate task. This distributes the hot customer_id values across multiple tasks and relieves the executor OOM without changing the join logic. It is the built-in mechanism for exactly this scenario.

Why this answer

AQE skew join handling identifies partitions that are disproportionately large after the shuffle and splits them into smaller sub-partitions, each handled by a separate task. This spreads the hot customer_id values across executors and resolves the OOM. Raising shuffle partitions or repartitioning by the join key keeps each hot key in one partition, and broadcasting a 300 GB table is infeasible.

Exam trap

The trap here is believing that increasing shuffle partitions or repartitioning by the join key will split a hot key, when hash partitioning always routes all rows for one key to a single partition.

181
Multi-Selectmedium

A PySpark DataFrame job on Databricks runs slowly. Inspection of the Spark UI shows that a shuffle stage writes 200 partitions but downstream stages process only a few, and the physical plan shows an Exchange before a filter. Which two changes are most likely to improve performance? (Choose two.)

Select 2 answers
A.Increase spark.sql.shuffle.partitions to a much larger number so each task handles fewer rows.
B.Repartition the DataFrame on the join key before the shuffle to colocate matching rows.
C.Reduce the number of shuffle partitions or coalesce the output so fewer, larger partitions are produced.
D.Cache the shuffled DataFrame with persist so downstream stages reuse the same partitions.
E.Apply the filter before the join or aggregation that triggers the Exchange so less data is shuffled.
AnswersC, E

When a shuffle produces many near-empty partitions, consolidating them reduces task launch overhead and improves per-task efficiency. Coalescing after the shuffle or lowering the partition count creates fewer, better-sized partitions for downstream stages. This matches the observed pattern of many partitions with little data. It addresses the scheduling waste visible in the UI.

Why this answer

The plan shows an Exchange feeding stages that process only a few partitions, meaning shuffle volume and partition sizing are both inefficient. Filtering earlier cuts the rows that ever reach the shuffle, and consolidating partitions removes the overhead of many tiny tasks. Together they reduce both the data moved across the network and the number of tasks scheduled, which is what the UI evidence points to.

Resource or caching changes would not target either cause.

Exam trap

The trap here is reaching for more shuffle partitions by reflex, when the UI actually shows too many near-empty partitions and excess shuffle input.

182
MCQmedium

When using Spark Connect, how does the client handle the authentication process with the Databricks workspace?

A.The client sends credentials as plaintext in the gRPC headers.
B.Authentication happens after the first query is executed.
C.The client token is provided as part of the connection string or environment variables.
D.Spark Connect relies on SSH keys stored on the local machine.
AnswerC

Authentication in Spark Connect is typically handled by providing a token in the SPARK_REMOTE connection string (e.g., token=...) or via Databricks profile configurations. The Spark Connect client library reads these tokens and includes them in the metadata of the gRPC requests for authentication against the Databricks compute resource.

Why this answer

Spark Connect uses the standard Databricks authentication mechanisms, primarily personal access tokens (PAT) or OAuth tokens, which are passed within the connection string or via environment variables. The client library handles the secure transmission of these credentials over the gRPC channel using TLS encryption, ensuring that the remote cluster validates the identity of the client before allowing any logical plan execution or data access.

Exam trap

Candidates often incorrectly assume that authentication is handled by the SparkSession object itself. In reality, it is managed via connection strings or environment variables passed to the client library.

183
MCQmedium

Which TWO factors contribute to the 'Data Locality' optimization in Spark?

A.The physical proximity of the storage to the compute nodes.
B.The number of executors running on the driver node.
C.The task scheduler's ability to query the block location.
D.The version of the Spark driver being used.
E.The total amount of memory assigned to the driver.
AnswerA, C

Spark's scheduler prioritizes tasks on nodes where the data is already physically located. This reduces network latency and improves performance by reading data from local disk or memory rather than across the network, making the storage-compute relationship vital for performance in large-scale data processing jobs.

Why this answer

Data locality is a critical optimization where Spark tries to schedule tasks on the node where the data resides. This prevents massive data movement across the network, which is the slowest part of a distributed system. By understanding how Spark respects the location of data blocks during task scheduling, developers can better partition their data and choose optimal file storage layouts for their workloads.

Exam trap

Candidates often assume data locality is a hardware feature managed by the storage layer alone, failing to recognize that the Spark scheduler must explicitly query block locations to decide where to place tasks.

184
MCQmedium

An engineering team wants to execute PySpark queries locally from an Integrated Development Environment (IDE) while offloading all distributed compute and data processing to a remote Databricks cluster. Which Spark Connect component architecture makes this workflow possible?

A.The local IDE runs a complete Spark driver instance inside an embedded JVM while the remote cluster acts purely as an executor pool for parallel task execution.
B.The client application communicates directly with worker nodes via standard JDBC connections to bypass the driver entirely for lower latency querying.
C.The local client translates PySpark DataFrame API calls into protocol buffer messages, transmitting them over gRPC to a remote Spark driver that handles query planning and execution.
D.The remote cluster pushes compiled JAR files back to the local client machine where all shuffle partitions are materialized and aggregated locally in memory.
AnswerC

Spark Connect's client-server split lets the local IDE host a thin client that converts DataFrame API calls into protocol buffers, sent via gRPC to a remote Spark driver. All query planning and distributed execution therefore occur on the Databricks cluster, not locally.

Why this answer

Spark Connect decouples the client application from the Spark driver using a gRPC-based client-server architecture. The local IDE runs a thin client that translates DataFrame operations into protocol buffer plans, streaming them over a network channel to the remote Spark driver for execution and optimization.

Exam trap

Candidates often confuse Spark Connect with legacy JDBC/ODBC thin clients or traditional cluster-mode submissions, mistakenly thinking the local JVM processes data transformations before sending them over the network.

185
Multi-Selectmedium

Which THREE components are involved in the process of executing a Shuffle operation?

Select 3 answers
A.Map-side output files on executor nodes.
B.Network communication between executors.
C.Reducer-side input fetching and aggregation.
D.The Cluster Manager's central storage.
E.The Driver's task result accumulation.
AnswersA, B, C

During a shuffle, the map task writes its intermediate output to local disk on the executor. These files serve as the source for downstream tasks. Without these files, reduce tasks would have no data to fetch, making persistent local storage essential for the shuffle's completion in distributed environments.

Why this answer

Shuffles require the coordination of map-side output, network transfer, and reduce-side aggregation. Understanding these components is critical because shuffles are the most expensive part of a Spark job due to disk I/O and network latency. When data must be reshuffled, Spark must store map outputs, have the executors communicate to fetch this data, and then process the incoming streams to complete the final aggregation or join operations.

Exam trap

Candidates often overlook the map-side output files, focusing only on the network transfer. They forget that shuffle data must be materialized to disk before it can be fetched by reducers.

186
MCQmedium

What is the primary purpose of the 'Cache' command in Spark SQL?

A.To permanently store the data on the disk for future use.
B.To keep the data in memory to avoid redundant re-computation.
C.To force the garbage collector to free up memory immediately.
D.To automatically partition the data across the cluster for faster access.
AnswerB

The primary goal of caching is to keep a computed DataFrame in memory so that subsequent actions triggered on that data do not need to re-execute the entire lineage of transformations. This is crucial for performance when the same data is used multiple times within a job.

Why this answer

Caching allows developers to persist frequently accessed data in memory across multiple actions. By avoiding repeated reads from storage and re-computation of transformations, caching significantly speeds up iterative workloads, such as machine learning training or complex dashboard refreshes. However, it must be used judiciously, as memory is a finite resource, and unnecessary caching can lead to OOM errors and overall performance degradation.

Exam trap

Examinees often confuse caching with permanent storage or assume it automatically speeds up every single-use query without needing an iterative context.

187
Multi-Selectmedium

A developer is using Spark on Databricks and wants to understand how the Driver and Executors communicate during a job. Which two statements accurately describe this interaction? (Choose two.)

Select 2 answers
A.The Driver collects results from Executors after tasks complete.
B.The Driver schedules tasks and sends them to Executors for execution.
C.Executors send heartbeat messages to the Driver to report their status.
D.The Driver and Executors share the same JVM for efficient communication.
E.Executors communicate with each other directly to share intermediate data during a shuffle.
AnswersA, B

When an action is triggered, the Driver collects results from Executors. For example, in a collect() action, Executors send their partial results back to the Driver, which aggregates them. This is a key interaction: the Driver initiates tasks, Executors run them, and then send results back to the Driver. This statement accurately describes the return communication path.

Why this answer

In Spark's architecture, the Driver coordinates job execution by scheduling tasks and sending them to Executors. Executors run these tasks and, upon completion, send results back to the Driver for actions that require aggregation. This two-way communication is fundamental.

Executors do not communicate directly with each other for shuffle data; they use a shuffle service. Heartbeats are sent to the Cluster Manager, not the Driver. The Driver and Executors run in separate JVMs in cluster mode.

Exam trap

The trap here is assuming Executors communicate directly with each other or that the Driver and Executors share a JVM, which is only true in local mode.

188
MCQmedium

A developer is building a Spark SQL pipeline in a Databricks notebook. They need to persist an intermediate DataFrame, built from a transformation of a Delta table, as a physical table in the current database so other notebooks in the same cluster can query it. They also want the table metadata to be managed by the metastore and the data to reside in the default warehouse directory. Which Spark SQL statement should they use?

A.CREATE OR REPLACE TABLE intermediate_sales AS SELECT * FROM sales WHERE region = 'EU'
B.CREATE OR REPLACE GLOBAL TEMP VIEW intermediate_sales AS SELECT * FROM sales WHERE region = 'EU'
C.CREATE OR REPLACE VIEW intermediate_sales AS SELECT * FROM sales WHERE region = 'EU'
D.CREATE OR REPLACE TEMP VIEW intermediate_sales AS SELECT * FROM sales WHERE region = 'EU'
AnswerA

This statement creates a managed Delta table in the current database, stores its data in the warehouse directory, and registers metadata in the metastore. Because it is a managed table, it persists beyond the session and is queryable from other notebooks on the same cluster, satisfying both the persistence and metastore-management requirements described in the scenario.

Why this answer

Creating an OR REPLACE TABLE with a SELECT statement materializes the transformed data as a managed Delta table, registers it in the metastore, and stores it in the default warehouse location. This matches the need for a persistent, cross-session-accessible physical table, unlike temporary views, global temporary views, or logical views, which do not persist materialized data.

Exam trap

The trap here is assuming that a global temporary view persists data across sessions or is stored in the metastore, when it is still session-scoped and stores no physical data.

189
MCQmedium

You have a Pandas API on Spark DataFrame `psdf` that was created from a Spark DataFrame with 200 partitions. You call `psdf.head(10)` in a Databricks notebook. What is the most likely performance characteristic of this operation?

A.It shuffles all data to a single partition before returning the first 10 rows, causing a full data movement.
B.It fails with an error because `head` is not supported on DataFrames with more than 100 partitions.
C.It collects only the necessary rows from the first partition(s) and returns quickly without scanning the entire dataset.
D.It triggers a full scan of all 200 partitions to ensure the first 10 rows are correctly ordered.
AnswerC

`head(10)` in Pandas API on Spark is optimized to fetch only the required number of rows. Spark executes a `limit` operation that reads from the first partition(s) until 10 rows are collected, avoiding a full scan. This makes it efficient even on large datasets, as it does not process all 200 partitions. The operation returns a small pandas DataFrame to the driver.

Why this answer

The `head` operation in Pandas API on Spark is designed to be efficient by leveraging Spark's `limit` operator, which reads only the necessary rows from the first partitions. It does not scan the entire dataset or perform a full shuffle. This makes it suitable for quickly inspecting large DataFrames without incurring heavy computation costs.

Exam trap

The trap here is assuming that any operation on a large partitioned DataFrame must scan all partitions, when in fact limit-based operations like `head` are optimized to read only what is needed.

190
MCQhard

A developer is using Spark Connect to connect to a Databricks cluster from a remote Python client. They need to run a custom Python function on a DataFrame column. They define the function and register it as a UDF using spark.udf.register(). After executing the job, they notice that the UDF fails with a ModuleNotFoundError for a library that is installed on their local machine but not on the cluster. What is the most likely cause and the appropriate solution?

A.The UDF is executed on the cluster, and the required library must be installed on all cluster nodes; the solution is to install the library on the cluster.
B.The UDF is executed on the client, so the library must be installed locally; the error indicates a local environment issue.
C.The UDF is executed in a separate Python process on the client, so the library must be installed in that process; the solution is to add the library to the client's Python path.
D.The UDF is executed on the driver node only, so the library must be installed on the driver; the solution is to restart the driver with the library.
AnswerA

This is correct because Spark Connect sends the UDF code to the server, where it is executed on the cluster's executors. Any dependencies used by the UDF must be available on the cluster nodes. The ModuleNotFoundError indicates that the library is not installed on the cluster. The appropriate solution is to install the library on the cluster, either via cluster libraries or init scripts.

Why this answer

In Spark Connect, UDFs are serialized and sent to the Spark server for execution on the cluster. Therefore, any Python libraries used by the UDF must be installed on the cluster nodes. The ModuleNotFoundError indicates that the library is missing on the cluster, so the solution is to install it there, not on the client.

Exam trap

The trap here is assuming that UDFs run on the client in Spark Connect, when they actually run on the remote cluster.

191
MCQmedium

A data scientist is working with a pandas-on-Spark DataFrame psdf that has a column 'category' with many unique values. They want to apply a custom Python function to each group to compute a complex statistic. They consider using psdf.groupby('category').apply(my_func). Which statement accurately describes the execution and potential performance implications of this operation?

A.It applies my_func in a distributed manner across partitions without shuffling, and the function must return a scalar or a pandas Series.
B.It uses the Spark Catalyst optimizer to translate my_func into native Spark expressions, so no Python UDF is involved and performance is optimal.
C.It triggers a full shuffle to group data, then applies my_func to each group as a pandas DataFrame on the executor, and the function's return type determines the result schema.
D.It applies my_func to each partition independently without grouping, and the results are concatenated, which is efficient for large datasets.
AnswerC

groupby().apply() performs a shuffle to co-locate rows of the same group. On each executor, it converts each group into a pandas DataFrame and applies my_func. The return type—whether scalar, Series, or DataFrame—is inferred to build the output schema. This can be powerful but may cause memory issues if a single group is large, because the entire group must fit in memory as a pandas object.

Why this answer

groupby().apply() in pandas API on Spark shuffles data to group rows, then applies the function to each group as a pandas DataFrame on executors. This enables complex group-wise logic but can be slow and memory-intensive because it involves a shuffle and materializes each group in memory. The return type of the function determines the output schema, and the operation is not optimized by Catalyst.

Exam trap

The trap here is assuming that groupby().apply() is as optimized as native Spark groupBy operations, when it actually uses a Python UDF and requires a full shuffle.

192
MCQmedium

A Databricks job reads a Parquet dataset, applies a chain of `withColumn` transformations, and writes it back with `.write.mode("overwrite").parquet(path)`. The Spark UI shows 6000 small output files totaling 50 GB, and a downstream reader is slow because of per-file overhead. You want to reduce the file count while keeping the write as a Spark-native operation. What should you do?

A.Convert the DataFrame to an RDD and call `saveAsTextFile` to let Spark choose the file layout.
B.Set `spark.sql.files.maxPartitionBytes` to a very large value before the write.
C.Add `.option("maxRecordsPerFile", 100000)` to the write call to merge small files.
D.Call `df.repartition(200)` immediately before the write to consolidate data into fewer output partitions.
AnswerD

Each output partition produces one file, so reducing the partition count directly reduces the file count. Repartitioning performs a full shuffle that evenly distributes rows across the target partitions, which yields larger, more balanced files that downstream readers can scan with far less per-file overhead.

Why this answer

Output file count equals the number of partitions at write time, so consolidating partitions with a shuffle-based repartition before writing produces fewer, larger files. Read-side settings and per-file record caps do not change the write fan-out, and dropping to RDD text output sacrifices the columnar format that makes the downstream reads fast.

Exam trap

The trap here is confusing a read-side tuning knob such as the maximum partition bytes with a write-side control over output file count.

193
MCQmedium

Refer to the exhibit. You are reviewing the logs for a Spark application and notice the warning regarding broadcasting a large task binary. What is the most likely cause and mitigation?

A.The partition size is too large; increase spark.sql.files.maxPartitionBytes.
B.A large variable is captured in a closure; use a Broadcast Variable.
C.The executors have insufficient memory; increase spark.executor.memory.
D.The cluster is out of network bandwidth; enable compression.
AnswerB

When a large object is referenced inside a transformation, Spark tries to serialize it with the task. A broadcast variable provides a mechanism to distribute the object efficiently once per node rather than once per task, preventing the warning and reducing the serialization burden on the driver.

Why this answer

The warning indicates that a large object, likely a high-dimensional collection or a large variable, is being captured in a closure and broadcast to all executors. This increases network pressure and memory usage. The mitigation is to use a broadcast variable or remove the reference to the large object from the closure to prevent serialization overhead and potential performance degradation.

Exam trap

Candidates often assume the solution is to increase the executor memory. However, the error is caused by serializing a large object in a closure, which remains an issue regardless of total heap size.

194
MCQmedium

Which TWO of the following statements correctly describe the role of the Spark Executor in a Databricks environment?

A.Executors are responsible for scheduling individual tasks on worker nodes.
B.Executors are responsible for executing the code logic assigned to them by the Driver.
C.Executors are responsible for storing data cached by the user in memory or on disk.
D.Executors manage the cluster-wide SparkContext.
E.Executors are responsible for physical cluster node provisioning.
AnswerB, C

This is the primary function of an executor. Once the driver sends task code to the executor, the executor performs the data processing, transformation, and aggregation operations. This separation of duties allows the driver to focus on orchestration while executors focus on parallelized data processing across the cluster.

Why this answer

Executors are the workhorses of a Spark cluster. They are responsible for executing the code logic (tasks) submitted by the driver and caching data in memory or on disk when requested by the application. Because they are the primary consumers of cluster resources, understanding their role is essential for capacity planning and ensuring that Spark jobs have enough memory and CPU to handle specific workloads without incurring out-of-memory errors.

Exam trap

Candidates often assume Executors are responsible for managing the cluster's overall health or scheduling, confusing their role as execution engines with the Driver's role as the application coordinator.

195
MCQmedium

A developer is building a Structured Streaming job in PySpark that reads from a Kafka topic and writes to a Delta Lake table. The job uses `outputMode("append")` and a 10-minute watermark on the event-time column `event_time`. A batch of late data arrives with events whose `event_time` is older than the watermark. What happens to these late events?

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

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

Why this answer

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

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

Exam trap

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

196
Multi-Selectmedium

A developer writes a Spark SQL query that groups orders by region and computes the total revenue per region, but also needs to return the number of distinct customers per region in the same result set. Which TWO expressions correctly compute the distinct customer count per region in a single GROUP BY region query? (Choose two.)

Select 2 answers
A.COUNT(customer_id) FILTER (WHERE customer_id IS NOT NULL)
B.COUNT(DISTINCT customer_id)
C.SUM(DISTINCT customer_id)
D.COLLECT_SET(customer_id)
E.APPROX_COUNT_DISTINCT(customer_id)
AnswersB, E

COUNT(DISTINCT customer_id) is a native aggregate that counts unique non-null customer_id values within each region group. Spark SQL supports this directly inside a GROUP BY region query, so it returns the distinct customer count per region without a subquery. It handles nulls by ignoring them, which matches typical distinct-count semantics in SQL.

Why this answer

To return a distinct customer count per region within a single GROUP BY region query, the developer can use COUNT(DISTINCT customer_id), which gives an exact deduplicated count per group, or APPROX_COUNT_DISTINCT(customer_id), which returns an approximate count using a sketch and is preferable when cardinality is high. Both are valid aggregates that operate per group and produce one value per region.

Exam trap

The trap here is treating COUNT with a FILTER clause as equivalent to a distinct count, when filtering nulls does not remove duplicate customer_id values within a region.

197
MCQmedium

Refer to the exhibit. Which action is the most likely cause of the error shown in the Spark job logs?

A.The executor nodes have insufficient memory for the shuffle operation.
B.The driver is attempting to aggregate a massive dataset into its memory.
C.The cluster configuration uses too many small partitions.
D.The broadcast join threshold is set too low for the current job.
AnswerB

The collect() method triggers the transfer of all partitions from worker nodes to the driver node. If the combined data exceeds the driver's allocated memory, the JVM throws an OOM error. This is a common architectural mistake when debugging or extracting large-scale distributed data to a single location.

Why this answer

The collect() action attempts to pull the entire DataFrame into the driver node's memory. When the dataset size exceeds the heap space allocated to the driver, a Java heap space error occurs. Developers must avoid collecting large datasets and instead use take() or head() for previews, or write results to cloud storage to maintain application stability when working with big data at scale.

Exam trap

Candidates often assume that calling 'collect()' is a safe way to inspect data, failing to realize it pulls the entire result set into the driver's memory, causing OOM errors.

198
Multi-Selecthard

A developer is using Pandas API on Spark and encounters a `compute.ops_on_diff_frames` error when combining two Pandas-on-Spark DataFrames. Which two actions can resolve this error? (Choose two.)

Select 2 answers
A.Use `psdf.spark.frame()` to extract the underlying Spark DataFrame and perform a join with Spark SQL.
B.Convert both DataFrames to pandas DataFrames using `to_pandas()` and then perform the operation locally.
C.Ensure both DataFrames originate from the same Spark DataFrame or are derived from a common ancestor without independent transformations.
D.Repartition both DataFrames to the same number of partitions before the operation.
E.Set `pyspark.pandas.options.compute.ops_on_diff_frames` to True to allow operations across different DataFrames.
AnswersC, E

Pandas API on Spark tracks lineage via an internal anchor. If both DataFrames share the same anchor, operations are allowed without the configuration flag. Deriving them from a common source or using `attach` to align anchors avoids the error and keeps execution efficient by avoiding unnecessary shuffles.

Why this answer

The error arises when operations combine DataFrames with different internal anchors. Enabling `compute.ops_on_diff_frames` explicitly allows such operations, while aligning anchors by deriving from a common source avoids the error altogether. Both approaches keep computation distributed and within the Pandas API on Spark.

Exam trap

The trap here is thinking that repartitioning or converting to pandas resolves the anchor mismatch, when the real solutions are either enabling the specific option or unifying the DataFrames' lineage.

199
MCQeasy

An analyst has a Spark SQL DataFrame named events with a string column event_time in the format 'yyyy-MM-dd HH:mm:ss'. They want to add a new column event_date containing only the date portion, keeping the original column intact. Which expression should they use in a select statement?

A.date_format(event_time, 'yyyy-MM-dd')
B.trunc(event_time, 'MM')
C.substring(event_time, 1, 10)
D.to_date(event_time, 'yyyy-MM-dd HH:mm:ss')
AnswerD

The to_date function parses the string using the supplied pattern and returns a date value containing only the date portion. Using the format string ensures correct parsing of the timestamp text. This produces the desired event_date column while leaving the original event_time column unchanged, which is exactly what the scenario requires.

Why this answer

Using to_date with the matching pattern correctly parses the timestamp string and returns a date-typed column containing only the date component. The other functions either return strings, operate on date types rather than strings, or rely on fragile character extraction, so they do not produce a proper date column as required.

Exam trap

The trap here is choosing substring because it visually looks like it extracts the date, while ignoring that it returns a string and is not robust to format variations.

200
MCQhard

A developer needs to add a computed column `discounted_price` equal to `price * 0.9` to an existing Delta table `products` and persist the change so all future queries see the new column. The table already contains data. Which statement should the developer run?

A.`ALTER TABLE products ADD COLUMNS (discounted_price DOUBLE)`
B.`UPDATE products SET discounted_price = price * 0.9`
C.`ALTER TABLE products ADD COLUMNS (discounted_price DOUBLE GENERATED ALWAYS AS (price * 0.9))`
D.`CREATE OR REPLACE TABLE products AS SELECT *, price * 0.9 AS discounted_price FROM products`
AnswerC

A generated column defined with `GENERATED ALWAYS AS` is computed automatically from the expression whenever rows are written, and adding it to a Delta table backfills values for existing rows. This persists the derived value in the table so all future queries see `discounted_price` populated. It is the declarative way to express the computed column in Spark SQL on Delta.

Why this answer

Adding a persisted computed column to a Delta table is done with `ALTER TABLE ... ADD COLUMNS` using a `GENERATED ALWAYS AS` expression. Spark computes the value from the base column for existing and future rows, so queries see a populated `discounted_price` without a manual update.

Statements that assume the column already exists or recreate the table either fail or risk data integrity.

Exam trap

The trap here is assuming that adding a plain column automatically populates it from an expression, when a bare `ADD COLUMNS` leaves existing rows null and requires a generated-column clause to derive values.

201
MCQmedium

A job is reading a huge amount of data from a table, but only uses three columns. Which optimization technique will provide the most significant I/O performance benefit?

A.Caching the DataFrame.
B.Column pruning.
C.Increasing the number of partitions.
D.Broadcasting the table.
AnswerB

Column pruning forces Spark to read only the columns required by the transformation. By ignoring unused columns at the storage layer, you significantly reduce the amount of data transferred from disk to memory, which is the primary performance gain for wide tables in distributed analytics and ETL applications.

Why this answer

Column pruning is the practice of reading only the required columns from a data source. In formats like Parquet, this reduces the total amount of data read from disk and transferred across the network. This is a crucial optimization for Databricks developers to reduce I/O bottlenecks and improve overall pipeline speed, especially when dealing with wide tables containing hundreds of unused columns in analytical query workloads.

Exam trap

Candidates often select repartitioning or caching instead of column pruning, misunderstanding that I/O bottlenecks depend on the amount of data read from disk.

202
MCQmedium

A developer has a local Python script that connects to a Databricks cluster using Spark Connect and creates a DataFrame from a small list of tuples. They then call .collect() on the DataFrame and receive the results. Which statement accurately describes how the data and operations are processed in this scenario?

A.The DataFrame is created and processed entirely on the client; the cluster is only used for storage.
B.The client serializes the logical plan and sends it to the Spark server, which executes the plan and returns the collected results.
C.The client executes the DataFrame operations locally using a built-in Spark engine, then syncs the results to the cluster.
D.The client sends the raw data to the cluster, which then returns a Python object that the client uses to perform further operations locally.
AnswerB

This is correct because Spark Connect decouples the client from the Spark driver. The client builds a logical plan for the DataFrame operations and sends it over gRPC to the Spark server, which executes the plan on the cluster. The results are then returned to the client when an action like collect() is called.

Why this answer

Spark Connect uses a client-server architecture where the client builds a logical plan and sends it to the Spark server via gRPC. The server executes the plan on the cluster and returns the results. This decouples the client from the Spark driver, allowing remote execution without a local Spark context.

Exam trap

The trap here is assuming that Spark Connect runs a local Spark engine on the client, when in fact all computation is performed on the remote Spark server.

203
MCQmedium

Which Spark configuration property determines the maximum amount of memory the Spark Driver can request for itself when running on a Kubernetes cluster?

A.spark.executor.memory
B.spark.driver.memory
C.spark.memory.fraction
D.spark.kubernetes.driver.limit.memory
AnswerB

This property explicitly sets the heap size for the Spark Driver process. In a containerized environment like Kubernetes, the cluster manager uses this value to allocate resources for the driver pod, ensuring the driver has sufficient memory to manage the job's metadata and DAG scheduling requirements.

Why this answer

The 'spark.driver.memory' property is the primary configuration used to define the heap size for the Spark Driver process. In Kubernetes environments, this value is translated into the container resource requests for the driver pod. Proper sizing prevents OOM errors during large collect operations or when handling massive metadata for complex query plans, ensuring the driver maintains stability throughout the application lifecycle.

Exam trap

Candidates often confuse 'spark.driver.memory' with 'spark.executor.memory', assuming the same property applies to both the driver and the worker processes.

204
MCQmedium

Which of the following Spark SQL configuration settings should be adjusted to prevent the 'Driver OOM' error when collecting massive amounts of query results to the driver node?

A.spark.driver.maxResultSize
B.spark.executor.memory
C.spark.sql.shuffle.partitions
D.spark.memory.fraction
AnswerA

This configuration sets the limit on the total size of serialized results of all actions that can be returned to the driver. By increasing this value, you allow larger result sets, though it is often safer to rewrite the query to write the results to storage rather than collecting them locally.

Why this answer

Collecting data to the driver node is a dangerous operation in distributed computing because the driver has limited heap memory. Spark provides configuration limits to prevent users from accidentally crashing the driver by pulling too much data from the executors. Adjusting these settings—or better, avoiding collecting altogether—is a key skill for Databricks developers to maintain cluster stability in production environments.

Exam trap

Candidates often confuse 'spark.driver.maxResultSize' with cluster-level memory configurations like 'spark.executor.memory'. They incorrectly try to increase executor memory to solve a driver-side data collection crash.

205
MCQeasy

A data analyst needs to run a Spark SQL query that returns the top 5 highest-paid employees from a table named employees, ordered by salary descending. Which query correctly returns exactly five rows?

A.SELECT TOP 5 * FROM employees ORDER BY salary DESC
B.SELECT * FROM employees ORDER BY salary DESC LIMIT 5
C.SELECT * FROM employees ORDER BY salary DESC FETCH FIRST 5 ROWS ONLY
D.SELECT * FROM employees LIMIT 5 ORDER BY salary DESC
AnswerB

This query sorts all employees by salary in descending order and then limits the result to the first five rows. In Spark SQL, LIMIT after ORDER BY returns the top N rows according to the sort, which exactly matches the requirement to get the five highest-paid employees.

Why this answer

To get the top 5 highest-paid employees, the query must sort by salary descending and then apply LIMIT 5. Spark SQL supports ORDER BY followed by LIMIT, which returns the first five rows after sorting, exactly matching the requirement.

Exam trap

The trap here is assuming that other SQL dialects' TOP or FETCH FIRST syntax works in Spark SQL, or that LIMIT can precede ORDER BY.

206
MCQeasy

Which SQL function is used to concatenate strings while allowing you to specify a custom separator, handling null values by ignoring them?

A.concat()
B.concat_ws()
C.array_join()
D.format_string()
AnswerB

The 'concat_ws' function accepts a separator as its first argument and handles NULL values by ignoring them. This behavior is crucial for data cleaning and string formatting tasks where missing values should not invalidate the entire concatenated output, ensuring a consistent and readable result string for downstream consumers.

Why this answer

The 'concat_ws' function stands for 'concatenate with separator'. It is specifically designed to take a delimiter as the first argument, followed by a variable number of strings. A key feature of 'concat_ws' is its ability to skip NULL values automatically, preventing the entire result from becoming NULL—a common issue with the standard 'concat' function.

This makes it ideal for building CSV-like strings or concatenating fields that may contain missing data.

Exam trap

Many candidates confuse concat_ws() with standard concat(), overlooking how standard concat turns the entire output null if any single column contains a null value.

207
MCQhard

A Spark job reads a large Parquet dataset, performs a filter, and then a groupBy aggregation. The job's DAG shows two stages: one for the filter and one for the aggregation. The first stage has 200 tasks, and the second stage has 200 tasks. The job is running on a cluster with 10 executors, each with 8 cores. The engineer observes that the second stage takes significantly longer than the first. Which of the following is the most likely cause for the increased duration in the second stage?

A.The second stage involves a shuffle, which requires data to be repartitioned across the network, causing additional I/O and serialization overhead.
B.The second stage is reading from disk, while the first stage reads from memory, causing slower performance.
C.The second stage has a larger number of partitions than the first stage, causing more tasks to be scheduled.
D.The second stage has more tasks than the first stage, leading to higher scheduling overhead.
AnswerA

The groupBy aggregation triggers a shuffle, where data is redistributed across executors based on the grouping key. This involves network transfer, disk I/O, and serialization, making the second stage slower than the filter stage, which is narrow and operates on data locally. The shuffle is the primary reason for the increased duration.

Why this answer

The groupBy aggregation requires a shuffle, which redistributes data across the cluster based on the grouping key. This involves network transfer, disk I/O, and serialization/deserialization, all of which add significant overhead compared to narrow transformations like filter. Even with the same number of tasks, the shuffle makes the second stage slower.

Exam trap

The trap here is focusing on the number of tasks or partitions as the cause of slowness, rather than recognizing the inherent cost of a shuffle operation.

208
Multi-Selecthard

A Databricks engineer is diagnosing why a Spark job's shuffle phase writes a very large amount of data to disk. The engineer wants to reduce shuffle overhead by changing how the job is structured and configured. Which TWO actions are most likely to reduce the volume of shuffle data written? (Choose two.)

Select 2 answers
A.Increase spark.sql.shuffle.partitions to a much larger value without changing the query plan.
B.Pre-aggregate data with reduceByKey before a subsequent join so fewer records participate in the shuffle.
C.Use broadcast hash join instead of sort-merge join when one side of the join is small enough to fit within the broadcast threshold.
D.Call persist(MEMORY_ONLY) on the DataFrame before the shuffle stage.
E.Enable the Kryo serializer instead of the default Java serializer for the shuffle.
AnswersB, C

Combining values locally with reduceByKey before the shuffle reduces the number of records that must be repartitioned, lowering both shuffle write and read volume. This map-side aggregation is the classic optimization for reducing network traffic in wide transformations.

Why this answer

Eliminating a shuffle through broadcast hash join and reducing record counts through map-side pre-aggregation both cut the actual bytes that must be written and read across the network. Partition-count tuning, serializer choice, and caching change performance characteristics without removing the underlying data movement.

Exam trap

The trap here is treating a larger shuffle partition count as a way to reduce shuffle data, when it only splits the same volume into more files.

209
MCQhard

When executing a Spark SQL query, what does the Catalyst optimizer perform during the 'Analysis' phase?

A.It translates the logical plan into a series of physical execution steps.
B.It checks the existence of tables and columns in the catalog.
C.It pushes down predicates to the data source to minimize I/O.
D.It chooses the most efficient join strategy based on table size.
AnswerB

The Analysis phase is explicitly responsible for verifying that all referenced objects, such as tables and columns, exist in the underlying metadata catalog. It resolves these names into concrete references, which is a prerequisite for all further optimization and execution steps in the Spark SQL pipeline.

Why this answer

The Analysis phase is the first step in the query optimization process where Spark resolves identifiers and validates the schema. It checks if tables and columns exist in the catalog and resolves data types. Without this step, Spark would not be able to build a logical plan for the query, as it would not know which data to fetch or how to process the specified columns and tables correctly.

Exam trap

Candidates often confuse the Analysis phase with physical optimization or logical plan generation, forgetting that checking catalog existence happens first.

210
MCQmedium

In the Databricks Spark environment, what is the role of the 'Shuffle Service'?

A.It manages the distribution of data across HDFS clusters.
B.It enables executors to fetch data from terminated executors.
C.It converts narrow transformations into wide ones.
D.It caches all RDD partitions in memory.
AnswerB

By decoupling the lifecycle of the shuffle data from the executor process, the Shuffle Service allows shuffle data to persist even after the executor that created it has been terminated. This is vital for maintaining fault tolerance and ensuring jobs do not fail during dynamic cluster scaling.

Why this answer

The External Shuffle Service allows executors to be decommissioned or removed without losing intermediate shuffle data. By offloading shuffle file management to a persistent service, Spark ensures that if an executor terminates, other nodes can still fetch the data required to complete the shuffle. This is critical in Databricks for dynamic allocation and auto-scaling, as it maintains stability despite the frequent addition and removal of worker nodes during cluster runtime.

Exam trap

Candidates often think the Shuffle Service is required for all Spark jobs. They fail to realize it is specifically for maintaining shuffle data when executors are dynamically removed or decommissioned.

211
MCQeasy

Which component manages the lifecycle and allocation of executors in a Databricks cluster?

A.The SparkContext
B.The Databricks cluster manager
C.The user's notebook session
D.The Hadoop YARN service
AnswerB

The cluster manager handles the lifecycle of the executor processes, including adding or removing nodes based on load or termination requests. This ensures that the infrastructure matches the cluster configuration defined by the user, providing a stable environment for the Spark driver to execute its tasks.

Why this answer

The Databricks cluster manager is responsible for requesting resources from the cloud provider, initiating the executor processes on those nodes, and monitoring their health. This architectural layer provides the abstraction that allows users to simply define a cluster size, while Databricks handles the complex underlying infrastructure provisioning and lifecycle management required to run distributed Spark applications reliably.

Exam trap

Test-takers frequently mistake the Databricks cluster manager for the Apache Spark Driver or cluster-agnostic cloud services, missing that Databricks provides a specialized layer for resource provisioning and executor lifecycle management.

212
MCQhard

A developer has a Delta table events with a high-cardinality column user_id and a low-cardinality column country. A query filters on country = 'US' and also on user_id IN (...). The developer runs EXPLAIN and sees a full scan of all files. Which statement about data skipping and the Delta table's statistics correctly explains why the filter on country is not skipping files?

A.Delta data skipping only works on columns used in JOIN conditions, so a WHERE filter on country is never eligible for file pruning regardless of statistics.
B.Data skipping requires the filtered column to be the first column in the table's partitioning scheme, so country must be a partition column for any file pruning to occur.
C.Data skipping is disabled by default and must be enabled with spark.databricks.delta.dataSkipping.enabled before any file pruning can happen on any column.
D.Delta data skipping uses per-file min/max statistics, and if the table was written without collecting statistics or the file count is small, the country filter cannot eliminate files even though the value is low cardinality.
AnswerD

Delta data skipping relies on per-file statistics such as min/max and null counts stored in the transaction log. If statistics were not collected, or if there are too few files for skipping to matter, the optimizer cannot prune files for country = 'US'. Low cardinality alone does not guarantee skipping; the statistics must exist and the file layout must be granular enough to separate values.

Why this answer

Delta data skipping prunes files using per-file statistics such as min/max and null counts recorded in the transaction log. A filter on country can only eliminate files if those statistics were collected at write time and the file layout separates country values. Low cardinality does not by itself enable skipping, and skipping is not limited to join conditions or partition columns, nor does it require a special enablement flag for basic operation.

Exam trap

The trap here is assuming that a low-cardinality filter column automatically triggers file skipping, when skipping actually depends on the presence and usefulness of per-file statistics in the transaction log.

213
MCQhard

You are building a Structured Streaming pipeline that reads from a Delta table source and applies a stateful deduplication using dropDuplicates on a composite key. After several hours, the job fails with an error indicating that the state store has grown too large. You need to bound the state size while still removing duplicate events that arrive within a reasonable window. Which approach should you take?

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

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

Why this answer

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

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

Exam trap

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

214
Multi-Selecthard

Which THREE of the following are valid ways to monitor or debug Spark SQL query performance in Databricks?

Select 3 answers
A.Use the 'EXPLAIN' statement to inspect the physical plan of the query.
B.Examine the Spark UI to review stage-level details and task metrics.
C.Manually restart the cluster every time a query takes more than 10 seconds.
D.Use the Databricks Query Profile to visualize the query execution tree.
E.Increase the driver memory to the maximum available for all jobs.
AnswersA, B, D

The EXPLAIN command is essential for viewing the logical and physical plans generated by the Catalyst optimizer. It reveals how Spark intends to execute the query, showing joins, filtering, and scans, which helps developers identify if the plan is optimized or if performance is hindered by inefficient operations.

Why this answer

Monitoring tools like the Spark UI, Query Profile, and Explain plans are critical for developers to understand how their SQL code executes. They provide visibility into shuffle sizes, stage durations, and physical plan generation. Effectively using these tools is the difference between writing performant code and creating bottlenecks.

They allow developers to identify skewed data, suboptimal joins, and unnecessary data scanning, which are common issues in large-scale Spark SQL workloads.

Exam trap

Candidates often incorrectly include 'DESCRIBE HISTORY' as a performance debugging tool. While useful for auditing, it is not a primary tool for analyzing query execution plans or stage-level bottlenecks.

215
MCQmedium

A data engineer is working with the Pandas API on Spark and needs to convert a Spark DataFrame named `sdf` into a pandas DataFrame so it can be processed locally on the driver node. Which method should the engineer use to execute this conversion?

A.Call `sdf.to_pandas()` to collect the data from the distributed Spark DataFrame into a standard single-node pandas DataFrame on the driver.
B.Call `sdf.collect_as_pandas()` to execute the query and retrieve rows into a pandas DataFrame object on the driver node.
C.Call `sdf.to_spark()` followed by `.toPandas()` to leverage standard PySpark conversion mechanisms for better cluster stability.
D.Call `sdf.pandas_api()` to transform the distributed collection into a local pandas structure ready for machine learning tasks.
AnswerA

This method is the designated API function for converting a Pandas API on Spark DataFrame into a standard pandas DataFrame. It triggers immediate computation across the cluster and materializes the final dataset entirely within the driver node's local memory space.

Why this answer

The to_pandas() method explicitly collects data from distributed executors back to the driver node, returning a standard single-node pandas DataFrame. This operation requires sufficient driver memory to hold the entire dataset, making it crucial to apply appropriate filtering or sampling beforehand to prevent out-of-memory errors in large-scale cluster environments.

Exam trap

Many candidates confuse distributed Pandas API on Spark methods with PySpark DataFrame methods like toPandas(), assuming both share identical syntax and behavior across all underlying execution engines.

216
MCQeasy

What happens when an action is called on a Spark DataFrame?

A.The data is immediately written to disk.
B.The entire transformation graph is executed.
C.The Spark context is shut down.
D.The cache is automatically cleared.
AnswerB

Actions are the only operations that force Spark to evaluate the lazy transformation chain. By triggering the DAG scheduler, Spark processes the data and returns a result to the driver or writes it to a sink. This execution model allows for significant query-level optimizations before runtime starts.

Why this answer

An action triggers the execution of the DAG. Spark creates a job, splits it into stages, and launches tasks on executors to process the data. Until an action is called, Spark only builds a logical plan (lazy evaluation).

This is fundamental for query optimization; it allows the Catalyst Optimizer to analyze the entire plan before execution, ensuring that unnecessary operations are pruned and the most efficient physical plan is generated for the cluster.

Exam trap

Candidates often confuse transformations with actions, mistakenly believing that operations like select(), filter(), or withColumn() trigger immediate computation in Spark when they actually just build up the logical plan.

217
MCQhard

A Spark application is running in cluster mode on Databricks. The driver program is running on a worker node, and the application has been running for several hours. Suddenly, the driver node experiences a hardware failure and crashes. What happens to the running tasks and the application?

A.The application fails completely, and all running tasks are lost; the application must be restarted from scratch.
B.The executors continue to run tasks and complete the job, and a new driver is automatically elected from the executors.
C.The running tasks continue on executors until they finish, and then the application terminates gracefully.
D.The cluster manager restarts the driver on another node, and the application resumes from the last checkpoint.
AnswerA

In Spark's architecture, the driver is the central coordinator. If the driver crashes, the entire application fails because there is no other component to manage the DAG Scheduler, Task Scheduler, or SparkContext. Running tasks on executors will be killed or eventually time out. Without a driver, the application cannot continue, and it must be restarted. Databricks may attempt to restart the driver if configured with automatic restart, but otherwise, it fails.

Why this answer

The driver is the single point of failure in a Spark application. It hosts the SparkContext, DAG Scheduler, and Task Scheduler. If the driver crashes, the application fails, and all running tasks are lost.

While the cluster manager may restart the driver, the application does not automatically resume from checkpoints unless explicitly configured. Thus, the correct outcome is complete failure and restart from scratch.

Exam trap

The trap here is assuming that Spark has built-in driver fault tolerance or automatic recovery from checkpoints, which is not the default behavior.

218
Multi-Selecthard

Which THREE of the following are valid ways to create a DataFrame from an existing table in Spark SQL?

Select 3 answers
A.spark.table('table_name')
B.spark.sql('SELECT * FROM table_name')
C.spark.read.table('table_name')
D.spark.open('table_name')
E.spark.from_table('table_name')
AnswersA, B, C

The 'spark.table()' method is a direct and efficient way to create a DataFrame by looking up the table name in the catalog. It is highly readable and is the standard way to retrieve a persistent table as a DataFrame for further transformation using the Spark API.

Why this answer

Spark provides multiple entry points to access table data. You can use the 'spark.table()' method for direct catalog access, 'spark.sql()' to execute a SELECT statement, or the 'spark.read.table()' method. All three provide the same underlying Dataset/DataFrame API, allowing for flexible programmatic interaction with metadata managed by the Hive Metastore or Unity Catalog, ensuring consistency across different development styles and API usage patterns in Databricks.

Exam trap

Candidates often select invalid or overly verbose DataFrame creation syntaxes, such as attempting to use non-existent SparkSession methods or confusing RDD loading commands with native table readers.

219
MCQmedium

Which component in the Spark architecture is responsible for maintaining the state of the Spark application and coordinating the execution of tasks across the cluster?

A.The Cluster Manager
B.The Spark Driver
C.The Executor
D.The Spark Master
AnswerB

The Driver serves as the engine's control plane. It converts the user program into tasks, schedules them on executors, and monitors progress. By maintaining the Directed Acyclic Graph (DAG) and task metadata, it manages the application lifecycle and ensures all transformations are executed in the correct dependency order.

Why this answer

The Driver process is the central coordinator in Spark. It runs the main() method, creates the SparkContext, and performs RDD graph scheduling and task distribution. Understanding the Driver's role is crucial because it is the primary bottleneck for metadata operations and task scheduling in a Spark cluster, and failure here results in the loss of the application's state and active execution context.

Exam trap

Candidates often confuse the Driver with the Cluster Manager or Executors. They mistakenly believe the Cluster Manager coordinates task execution, whereas the Driver is the actual brain managing the SparkContext and task scheduling.

220
MCQmedium

Which of the following describes the behavior of a 'Broadcast Hash Join' in Spark SQL?

A.The larger table is shuffled to match the smaller table's partitions.
B.The smaller table is sent to all executors, avoiding a shuffle of the large table.
C.Both tables are shuffled to a common partition based on the join key.
D.The join is performed entirely on the driver node.
AnswerB

By broadcasting the smaller table, Spark eliminates the need to move the large dataset across the network. The large table's partitions are processed in parallel on each executor against a full local copy of the small table. This is the most efficient join type for small-to-large table operations.

Why this answer

A Broadcast Hash Join is a highly efficient join strategy where the smaller table is sent to all worker nodes, allowing the join to occur locally in memory. This eliminates data shuffles, which are usually the slowest part of a distributed join. Understanding when the optimizer chooses this—and when to force it—is essential for optimizing SQL performance in Databricks environments where network bandwidth is a common bottleneck.

Exam trap

Students often assume both tables are broadcast or that a broadcast join requires shuffling both datasets, overlooking the core design of sending only the small table.

221
Multi-Selecthard

A data engineer is building a PySpark application that must validate incoming records in a DataFrame `raw` before loading them into a curated table. They want to apply user-defined validation logic that cannot be expressed with built-in functions, and they want the result to remain a DataFrame column of Boolean values. Which TWO approaches allow this? (Choose two.)

Select 2 answers
A.Use `df.transform` to apply a Python function that returns a Boolean column.
B.Use a Scala UDF registered in the Spark session and invoke it from a Python DataFrame column expression.
C.Use `df.rdd.map` to apply the Python function, then convert the RDD back to a DataFrame with a schema.
D.Use `pandas_udf` with a scalar return type and apply it via `withColumn`.
E.Define a Python function and register it with `spark.udf.register`, then call it from `selectExpr` or `expr`.
AnswersD, E

A `pandas_udf` with a scalar return type operates on batches using Arrow for data transfer and produces a column of the declared type. Applied through `withColumn`, it yields a Boolean column suitable for validation. This approach is generally faster than row-at-a-time Python UDFs because it amortizes serialization across batches.

Why this answer

Both a registered Python UDF invoked through SQL expressions and a scalar `pandas_udf` applied with `withColumn` produce Boolean DataFrame columns using custom Python logic. The RDD round-trip leaves the DataFrame API and requires manual schema handling. `transform` alone does not perform row-wise custom logic, and a Scala UDF is not a standard path in a PySpark-only application.

Exam trap

The trap here is assuming `df.transform` executes custom Python row logic, when it only applies a DataFrame-to-DataFrame function.

222
MCQhard

A developer is troubleshooting a Spark job that fails with an OutOfMemoryError on the driver. The job collects a large DataFrame to the driver using .collect() and then processes it locally. The developer wants to avoid the driver OOM while still obtaining the results. Which approach is most appropriate?

A.Replace .collect() with .take(1000) to limit the number of rows returned.
B.Increase spark.driver.memory to a very large value to accommodate the collected data.
C.Use .foreach() to process each row on the driver instead of collecting the entire DataFrame.
D.Write the DataFrame to a distributed storage system and then read it back in smaller chunks for local processing.
AnswerD

Writing the DataFrame to distributed storage (e.g., Delta Lake, Parquet) and then reading it in batches avoids bringing the entire dataset to the driver at once. This leverages Spark's distributed nature for the heavy lifting and allows the driver to process manageable chunks. It is a scalable pattern that prevents driver OOM while still enabling access to all data.

Why this answer

Collecting a large DataFrame to the driver is an anti-pattern because it moves all data to a single node, risking OOM. The scalable solution is to persist the DataFrame to distributed storage and then read it in smaller partitions or batches for local processing. This maintains distribution and avoids overwhelming the driver.

Exam trap

The trap here is believing that increasing driver memory is a sustainable fix, when the real solution is to avoid collecting large data to the driver altogether.

223
MCQeasy

What is the result of applying the COALESCE function in Spark SQL when multiple arguments are provided?

A.It returns the sum of all arguments.
B.It returns the first non-null argument.
C.It returns the last non-null argument.
D.It throws an error if any argument is null.
AnswerB

COALESCE iterates through its arguments in order and returns the first one that is not null. If all arguments are null, it returns null. This is the idiomatic way in Spark SQL to perform null-replacement or provide fallback values for columns containing missing or null data in source systems.

Why this answer

The COALESCE function is a standard SQL tool for handling null values. It returns the first non-null argument from a list. This is extremely useful in data engineering for providing default values or merging columns that have sparse data.

Understanding how to use it helps in writing cleaner SQL that handles missing values without needing complex CASE WHEN statements, which improves code readability and maintainability.

Exam trap

Candidates often confuse COALESCE with ISNULL or NVL, mistakenly believing it returns the first null value encountered or performs a conditional check rather than returning the first non-null argument found.

224
MCQhard

A Spark job performing a join between a 10 GB table and a 5 MB lookup table is running slowly, and the physical plan shows a SortMergeJoin. You want to avoid the shuffle. What should you do?

A.Repartition both DataFrames on the join key before the join.
B.Increase spark.sql.shuffle.partitions to 2000.
C.Use a cross join and filter afterward.
D.Set spark.sql.autoBroadcastJoinThreshold to a value larger than 5 MB, such as 10 MB, and ensure the small table is broadcast.
AnswerD

The auto broadcast join threshold controls the maximum size of a table that can be broadcast to all executors. The default is 10 MB, but if the small table is slightly above the threshold or statistics are missing, Spark may choose SortMergeJoin. Explicitly setting the threshold higher or using broadcast() hint forces a BroadcastHashJoin, eliminating the shuffle of the large table. This is the correct approach to avoid the shuffle in this scenario.

Why this answer

Broadcasting the small table eliminates the need to shuffle the large table, converting the join to a BroadcastHashJoin. This is achieved by ensuring the small table's size is below the auto broadcast join threshold or by using an explicit broadcast hint. The other options either do not change the join strategy or introduce unnecessary shuffles.

Exam trap

The trap here is thinking that increasing shuffle partitions or repartitioning will avoid the shuffle, when in fact they only change how the shuffle is performed.

225
MCQhard

Which property of RDDs (Resilient Distributed Datasets) is primarily responsible for Spark's fault tolerance during cluster execution?

A.Data replication across multiple nodes.
B.The RDD Lineage graph.
C.The Spark Driver's checkpointing of all intermediate tasks.
D.The use of a centralized data warehouse.
AnswerB

Lineage tracks the sequence of transformations applied to the data. If a node fails, Spark uses this graph to recompute only the lost partitions, ensuring the job completes successfully without needing a complete restart. This design is what makes Spark highly resilient in large, distributed compute environments.

Why this answer

Lineage (the dependency graph) is the core mechanism of RDD fault tolerance. Because RDDs are immutable and record their transformation history, Spark can reconstruct lost partitions by recomputing them from their parent RDDs. This is superior to traditional replication models because it avoids the high cost of copying data over the network, allowing Spark to maintain resilience while maximizing performance and minimizing storage overhead across the distributed cluster infrastructure.

Exam trap

Candidates often confuse RDD Lineage with Data Replication. They mistakenly believe Spark handles fault tolerance by copying data to other nodes, rather than recomputing lost partitions from the recorded lineage.

Page 2

Page 3 of 4

Page 4

All pages