Courseiva

Databricks Certified Associate Developer for Apache Spark (Databricks-Spark-Assoc) — Questions 76–150

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

Page 1

Page 2 of 4

Page 3
76
MCQmedium

An engineer is developing a Structured Streaming job that reads from an Apache Kafka source and writes the output continuously toDelta Lake using outputMode("append"). The stream occasionally experiences late-arriving data. Which downstream behavior can the engineer expect regarding the Delta Lake table?

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

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

Why this answer

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

Exam trap

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

77
MCQmedium

A developer writes a Spark Connect application that calls `df.cache()` on a large DataFrame, then performs several transformations and an action. The developer expects the cached data to persist on the client for reuse across sessions. Which statement describes what actually happens?

A.The cache call is silently ignored because Spark Connect does not support caching DataFrames.
B.Caching is only supported for temporary views and fails when called directly on a DataFrame in Spark Connect.
C.The cache is stored on the client machine, so subsequent sessions on the same laptop can reuse it without recomputation.
D.The cache request is sent to the server, where the data is cached in the cluster's memory or disk, and it persists only for the lifetime of that server-side session.
AnswerD

In Spark Connect, `cache()` sends a plan to the server that marks the DataFrame for caching. The server stores the data according to the storage level within the cluster. The cache is tied to the server-side session and is not available to other client sessions or after the session ends.

Why this answer

Caching in Spark Connect is a server-side operation. When a client calls `cache`, the request is transmitted to the server, which stores the data according to the specified storage level within that session. The cache is not stored on the thin client and does not survive beyond the server-side session's lifetime, so reuse across separate client sessions is not possible.

Exam trap

The trap here is assuming the thin client holds cached data locally, when caching actually occurs on the server within the session's lifetime.

78
MCQhard

A data engineer needs to add a column `rank_in_dept` to a DataFrame `employees` that ranks each employee by `salary` descending within their `department`, but they must not collapse rows. They also want ties to receive the same rank with gaps afterward. Which expression correctly produces this column using the DataFrame API?

A.employees.withColumn('rank_in_dept', row_number().over(Window.partitionBy('department').orderBy(col('salary').desc())))
B.employees.withColumn('rank_in_dept', rank().over(Window.partitionBy('department').orderBy(col('salary').desc())))
C.employees.withColumn('rank_in_dept', dense_rank().over(Window.partitionBy('department').orderBy(col('salary').desc())))
D.employees.groupBy('department').agg(rank().over(Window.orderBy(col('salary').desc())).alias('rank_in_dept'))
AnswerB

`rank()` assigns the same rank to tied values and leaves gaps in the sequence after ties, which matches the requirement. Using `Window.partitionBy('department').orderBy(col('salary').desc())` scopes the ranking per department and orders by descending salary. `withColumn` preserves all original rows, so no data is collapsed. This is the precise combination for the scenario.

Why this answer

The requirement calls for a window function that ranks within each department, orders by descending salary, keeps all rows, and leaves gaps after ties. `rank()` combined with a window partitioned by department and ordered by descending salary satisfies every condition. `dense_rank` omits gaps, `row_number` breaks ties arbitrarily, and aggregating with a window function is not valid.

Exam trap

The trap here is treating `rank` and `dense_rank` as interchangeable, when only `rank` leaves gaps after tied values.

79
MCQhard

A Spark job reads a large Parquet file, performs a groupBy operation, and then writes the result. During execution, the job fails with an OutOfMemoryError on the Driver. Which component is most likely responsible for the memory issue?

A.The DAG Scheduler
B.The Driver
C.The Cluster Manager
D.The Executors
AnswerB

The Driver is responsible for coordinating the job, including collecting results from executors when actions like collect() or take() are called. If the result set is large, the Driver's memory can be overwhelmed. Additionally, the Driver holds the DAG and may store broadcast variables. In this scenario, the groupBy operation may produce a large result that is inadvertently collected to the Driver, causing an OutOfMemoryError.

Why this answer

An OutOfMemoryError on the Driver typically occurs when the Driver attempts to hold too much data in memory. This can happen during actions like collect() or when broadcasting large variables. In this scenario, the groupBy operation might produce a large aggregated result that is being collected to the Driver, exceeding its memory.

The Executors process data in a distributed manner, but the Driver aggregates final results, making it the likely culprit.

Exam trap

The trap here is assuming that OutOfMemoryError always relates to Executors, but the error explicitly mentions the Driver, which has a different role and memory constraints.

80
MCQeasy

A developer submits a Spark application to a Databricks cluster. The application creates a SparkSession, reads a CSV file, and calls count() on the resulting DataFrame. Which component is responsible for translating this logical operation into a physical execution plan and coordinating its execution across the cluster?

A.The Driver, which hosts the SparkSession and the DAG Scheduler.
B.The Cluster Manager, which allocates containers for the application.
C.The Executor processes running on worker nodes.
D.The Catalog, which stores table and column metadata.
AnswerA

The Driver hosts the SparkSession and runs the DAG Scheduler, which converts the logical plan into stages and tasks, then coordinates their execution. For the count() call, the Driver plans the read and aggregation, schedules tasks on executors, and gathers the final result.

Why this answer

The Driver is the control plane of a Spark application. It hosts the SparkSession, builds the logical and physical plans through the Catalyst optimizer, and uses the DAG Scheduler to break the plan into stages and tasks that executors run, then aggregates their results for the count().

Exam trap

The trap here is confusing the Cluster Manager's resource allocation role with the Driver's planning and coordination responsibilities.

81
MCQmedium

A developer runs a Spark application on a Databricks cluster in Standard access mode. The application reads a Parquet file, applies a filter, and calls `df.cache()` before an action. During execution, the driver logs show that a stage is retried because a task failed with an executor lost error. Which component is responsible for rescheduling the failed task on another executor within the same application?

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

The Task Scheduler is responsible for launching tasks on executors via the SchedulerBackend and for retrying failed tasks up to `spark.task.maxFailures`. When an executor is lost, the Task Scheduler detects the failure and reschedules the affected tasks on other available executors. It also handles speculative execution and locality preferences, making it the correct component for this scenario.

Why this answer

The Task Scheduler is the component that launches individual tasks on executors and handles their retries. When an executor is lost, the Task Scheduler receives the failure notification from the SchedulerBackend and reschedules the failed tasks on other executors, up to the configured maximum number of failures. The DAG Scheduler operates at the stage level, while the Cluster Manager and Catalyst Optimizer do not manage task-level retries.

Exam trap

The trap here is confusing the DAG Scheduler's stage-level retry logic with the Task Scheduler's task-level retry logic, leading to the wrong component being selected.

82
MCQeasy

A developer is writing a Spark application that will run on a Databricks cluster. They need to ensure that the driver program can communicate with the executors and that tasks are distributed correctly. Which component is responsible for coordinating the execution of tasks across the executors?

A.SparkContext
B.Worker Node
C.Cluster Manager
D.Executor
AnswerA

The SparkContext is the entry point for Spark functionality and resides in the driver program. It coordinates the execution of tasks by communicating with the cluster manager to acquire executors, and then sends tasks to those executors. It also manages broadcast variables and accumulators. In this scenario, the SparkContext is responsible for the overall coordination.

Why this answer

The SparkContext, located in the driver, is the central coordinator for a Spark application. It connects to the cluster manager to request executors, then schedules tasks on those executors via the Task Scheduler. It also manages shared variables and the overall job execution.

Without the SparkContext, the application cannot run or distribute tasks.

Exam trap

The trap here is confusing the cluster manager's resource allocation role with task coordination, which is actually performed by the SparkContext.

83
MCQeasy

A developer wants to use the pandas API on Spark in a Databricks notebook. They have an existing PySpark DataFrame `sdf`. Which code snippet correctly creates a pandas-on-Spark DataFrame from `sdf` while preserving the distributed execution plan?

A.import pandas as pd; psdf = pd.DataFrame(sdf)
B.import pyspark.pandas as ps; psdf = sdf.to_pandas_on_spark()
C.import pyspark.pandas as ps; psdf = ps.from_pandas(sdf)
D.import pyspark.pandas as ps; psdf = ps.DataFrame(sdf)
AnswerD

`ps.DataFrame(sdf)` creates a pandas-on-Spark DataFrame that wraps the existing Spark DataFrame. The underlying Spark plan is preserved, and operations on psdf will execute distributedly. This is the standard way to convert a PySpark DataFrame to a pandas-on-Spark DataFrame without collecting data to the driver. The import `pyspark.pandas as ps` is the correct module for the pandas API on Spark.

Why this answer

To convert a PySpark DataFrame to a pandas-on-Spark DataFrame while preserving the distributed plan, use `pyspark.pandas.DataFrame(sdf)`. This wraps the Spark DataFrame and allows pandas-like operations that execute on Spark. Other approaches either collect data to the driver or use incorrect methods that do not exist or are meant for different conversions.

Exam trap

The trap here is confusing the conversion direction: `from_pandas` converts a local pandas DataFrame, while `DataFrame(sdf)` wraps a PySpark DataFrame for distributed operations.

84
MCQmedium

A data engineer has two Spark SQL DataFrames: customers (customer_id, name) and orders (order_id, customer_id, amount). They want to retrieve every customer along with their orders, but they also want to include customers who have placed no orders, showing null for the order columns. Which operation should they use?

A.customers.join(orders, customers.customer_id == orders.customer_id, 'full')
B.customers.join(orders, customers.customer_id == orders.customer_id, 'inner')
C.customers.join(orders, customers.customer_id == orders.customer_id, 'left')
D.customers.join(orders, customers.customer_id == orders.customer_id, 'right')
AnswerC

A left outer join keeps all rows from the left DataFrame, customers, and matches rows from orders where possible. Customers with no orders appear with null values in the order columns, which is exactly the desired behavior. This satisfies the requirement to include every customer while optionally attaching their orders.

Why this answer

A left outer join preserves every row from the customers DataFrame while attaching matching order rows, so customers without orders appear with nulls in the order columns. Inner, right, and full outer joins either drop unmatched customers or add unmatched orders, so they do not meet the stated requirement.

Exam trap

The trap here is mixing up left and right outer joins by focusing on the table that has more rows rather than on which DataFrame must be fully preserved.

85
MCQhard

A developer must join a 4 TB `transactions` DataFrame against a 900 MB `merchants` DataFrame on `merchant_id`. The cluster has 40 executors each with 16 GB of memory, and the job currently shuffles the large side. They want to avoid the shuffle entirely. Which change should they make?

A.Increase `spark.sql.autoBroadcastJoinThreshold` only, without calling `broadcast`.
B.Repartition both DataFrames by `merchant_id` before the join.
C.Set `spark.sql.shuffle.partitions` to 40 so the shuffle matches the executor count.
D.Call `broadcast(merchants)` when joining so the small DataFrame is shipped to every executor.
AnswerD

Broadcasting the 900 MB side replicates it to every executor and lets each task probe the local copy, eliminating the shuffle of the 4 TB side. The tradeoff is memory: 900 MB per executor fits within the available heap, so this is the standard remedy for a large-to-small join and directly removes the network exchange the developer wants gone.

Why this answer

An explicit broadcast hint forces Spark to collect the small DataFrame and distribute it as a read-only hash map to every executor, so the large side is scanned in place with no exchange. Because the small table fits comfortably in executor memory, this converts a shuffle-heavy sort-merge join into a map-side lookup and removes the network bottleneck.

Exam trap

The trap here is believing that raising `spark.sql.autoBroadcastJoinThreshold` alone guarantees a broadcast join, when the optimizer's size statistics can override it.

86
MCQeasy

A developer is writing a PySpark script that must run a SQL statement against DataFrames already registered as temporary views named orders and returns. The developer wants the query to use Spark SQL syntax while returning a DataFrame that can be further transformed with the DataFrame API. Which call accomplishes this?

A.spark.sql("SELECT o.order_id, r.amount FROM orders o JOIN returns r ON o.order_id = r.order_id")
B.spark.catalog.sql("SELECT o.order_id, r.amount FROM orders o JOIN returns r ON o.order_id = r.order_id")
C.spark.executeSql("SELECT o.order_id, r.amount FROM orders o JOIN returns r ON o.order_id = r.order_id")
D.spark.sqlContext.runSql("SELECT o.order_id, r.amount FROM orders o JOIN returns r ON o.order_id = r.order_id")
AnswerA

spark.sql executes a SQL string against the current catalog and returns a DataFrame, so the result can be chained with DataFrame transformations such as filter or withColumn. Because orders and returns are registered as temporary views, they resolve in the session catalog. This is the standard way to mix SQL and the DataFrame API in one pipeline.

Why this answer

SparkSession.sql runs a SQL string against views registered in the session catalog and returns a DataFrame, allowing the result to be combined with DataFrame transformations. Because orders and returns are temporary views in the same session, the query resolves and returns a join result that can be further processed. The other calls reference methods that do not exist on those objects.

Exam trap

The trap here is confusing the catalog interface, which manages metadata such as tables and functions, with the SQL execution entry point on the SparkSession.

87
MCQeasy

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

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

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

Why this answer

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

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

Exam trap

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

88
MCQhard

Refer to the exhibit. Which performance indicator suggests that Task 15 is likely causing a performance bottleneck during the execution of a join operation?

A.Task 12 having 500MB on Local Disk.
B.Task 15 having 1.2GB Shuffle Read.
C.Task 22 having 400MB Shuffle Write.
D.The total volume across all tasks is too low.
AnswerB

A high shuffle read value suggests that the task is processing a disproportionately large amount of data compared to its peers. This is a classic sign of data skew, where a specific key is overrepresented, causing the executor assigned to that partition to take much longer to finish.

Why this answer

The 'Shuffle Read' metric indicates the volume of data transferred over the network to that specific task. A high shuffle read compared to other tasks often points to data skew, where one partition receives significantly more data than others. This is a critical insight for developers because skew can lead to uneven executor loads, where one executor works much longer than others, effectively slowing down the entire stage of the Spark job.

Exam trap

Candidates often misread task duration or spill metrics as the primary indicator of skew, ignoring the distinct volume imbalance shown by massive shuffle read sizes.

89
MCQeasy

Which clause is used in a SELECT statement to filter the results based on aggregated values?

A.WHERE
B.HAVING
C.FILTER
D.LIMIT
AnswerB

The HAVING clause is designed specifically to filter data after the GROUP BY and aggregation operations have been performed. This is the only way to apply predicates to the results of aggregate functions, which is a required capability for generating summarized insights from large datasets in Spark SQL.

Why this answer

The HAVING clause is essential for filtering the output of an aggregation. While WHERE filters rows before aggregation, HAVING filters groups after the aggregation is calculated. Distinguishing between these two clauses is a foundational concept in SQL and critical for developers using Spark SQL to summarize data.

Failure to use the correct filter leads to syntax errors or incorrect analytical results, highlighting the necessity of this fundamental SQL skill.

Exam trap

Candidates commonly confuse WHERE and HAVING clauses, attempting to filter aggregate metrics inside the WHERE clause before aggregation occurs.

90
MCQeasy

A data analyst wants to use Pandas API on Spark in a Databricks notebook but is unsure how to import it. Which import statement correctly enables the Pandas API on Spark?

A.`import databricks.pandas as dp`
B.`import pandas as pd`
C.`from pyspark.sql import PandasAPI`
D.`import pyspark.pandas as ps`
AnswerD

The correct module for Pandas API on Spark is `pyspark.pandas`. Importing it as `ps` is a common convention. This provides the pandas-like API that runs on Spark. In Databricks, this module is available by default, and this import statement is the standard way to access the API.

Why this answer

Pandas API on Spark is accessed via the `pyspark.pandas` module. Importing it as `ps` is the standard convention. This module provides a pandas-like interface that translates operations to Spark, enabling distributed processing.

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

Exam trap

The trap here is confusing standard pandas with Pandas API on Spark, or assuming a Databricks-specific module exists, when the correct import is the open-source `pyspark.pandas`.

91
MCQmedium

A Spark job reads a large CSV file, performs a groupBy aggregation, and then writes the result. The Spark UI shows that the job has multiple stages, and one stage has a large number of tasks. Which factor primarily determines the number of tasks in the stage that performs the aggregation?

A.The value of spark.sql.shuffle.partitions.
B.The number of partitions in the input RDD/DataFrame.
C.The number of cores per executor.
D.The number of executors in the cluster.
AnswerA

For operations that trigger a shuffle, such as groupBy, the number of tasks in the subsequent stage is equal to the number of shuffle partitions. By default, spark.sql.shuffle.partitions is set to 200, but it can be adjusted. This configuration directly controls the parallelism of the aggregation stage.

Why this answer

The number of tasks in a stage that follows a shuffle (like an aggregation) is determined by the number of shuffle partitions, which defaults to spark.sql.shuffle.partitions (200). The input partitions affect the initial stage, while executors and cores affect concurrency, not the total task count.

Exam trap

The trap here is assuming that the number of tasks always equals the number of input partitions, ignoring that shuffle operations introduce a new partitioning determined by spark.sql.shuffle.partitions.

92
Multi-Selectmedium

A developer is working with a Pandas API on Spark DataFrame `psdf` and wants to perform operations that are efficient in a distributed environment. Which two operations are considered efficient and do not require collecting data to the driver? (Choose two.)

Select 2 answers
A.`psdf.to_pandas()`
B.`psdf.merge(other_psdf, on='id')`
C.`psdf.head(20)`
D.`psdf['value'].apply(lambda x: x * 2)`
E.`psdf.groupby('category').agg({'value': 'sum'})`
AnswersB, E

`merge` in Pandas API on Spark is implemented as a distributed join in Spark. It shuffles data across the cluster but processes it in parallel, and the result is a distributed DataFrame. This is an efficient operation for large datasets, as it leverages Spark's join optimizations. No data is collected to the driver.

Why this answer

Distributed operations like `groupby().agg()` and `merge()` execute in parallel across the cluster and return distributed DataFrames, avoiding driver collection. In contrast, `apply` with a lambda, `head`, and `to_pandas` either collect data to the driver or force row-by-row processing, making them less efficient for large-scale data.

Exam trap

The trap here is assuming that any pandas-like method is equally distributed, when methods like `apply` and `to_pandas` actually break distribution and collect data to the driver.

93
MCQmedium

You are processing JSON data from a stream. You need to extract nested fields from the JSON structure. Which function is the most efficient and standard way to handle this in Structured Streaming?

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

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

Why this answer

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

Exam trap

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

94
MCQeasy

A developer is writing a Spark Connect application and wants to create a SparkSession that connects to a remote Databricks cluster. The developer has the cluster's connection string and an access token. Which method should be used to build the session?

A.SparkSession.builder.master("spark://<host>:<port>").getOrCreate()
B.SparkSession.builder.appName("RemoteApp").config("spark.connect.host", "<host>").getOrCreate()
C.SparkSession.builder.config("spark.remote", "spark://<host>:<port>").getOrCreate()
D.SparkSession.builder.config("spark.remote", "sc://<host>:<port>").getOrCreate()
AnswerD

In Spark Connect, the connection to a remote server is configured using the spark.remote configuration property, which specifies the URL of the Spark Connect server (e.g., sc://host:port). The builder pattern with getOrCreate() is the standard way to create a session. This method correctly sets the remote endpoint and initializes the session.

Why this answer

To create a Spark Connect session, the developer must configure the spark.remote property with the Spark Connect server URL, which uses the sc:// scheme. The builder pattern with getOrCreate() is then used to instantiate the session. Other methods like master() or incorrect configuration keys will not establish a Spark Connect connection and may result in a local session or an error.

Exam trap

The trap here is confusing the Spark Connect configuration property spark.remote with traditional Spark master settings, or using the wrong URL scheme.

95
MCQmedium

A Spark application is submitted to a Databricks cluster. The application uses a broadcast variable to distribute a small lookup table to all Executors. Which component is responsible for broadcasting this variable?

A.The Driver
B.The DAG Scheduler
C.The Cluster Manager
D.The Executors
AnswerA

Broadcast variables are created on the Driver and then distributed to Executors. The Driver serializes the variable and sends it to each Executor, where it is cached for read-only use. This avoids shipping a copy with every task. The Driver initiates the broadcast and manages its distribution, making it the correct component.

Why this answer

Broadcast variables are created on the Driver via the SparkContext. The Driver serializes the variable and distributes it to all Executors, where it is cached. This mechanism reduces data transfer by avoiding sending the variable with each task.

The Driver is the initiator and manager of this process, while Executors are recipients. The Cluster Manager and DAG Scheduler are not involved in broadcasting data.

Exam trap

The trap here is thinking that Executors or the Cluster Manager handle broadcasting, but the Driver is the component that initiates and manages broadcast variables.

96
MCQmedium

Which configuration parameter should be adjusted to change the default number of partitions when reading from a shuffle-heavy operation?

A.spark.executor.memory
B.spark.default.parallelism
C.spark.sql.shuffle.partitions
D.spark.driver.memory
AnswerC

This parameter explicitly defines the number of partitions to use when shuffling data for joins or aggregations. By tuning this value, developers can control the level of parallelism, which is crucial for balancing the trade-off between task overhead and executor utilization during resource-intensive stages of a Spark job.

Why this answer

The `spark.sql.shuffle.partitions` configuration is the primary setting used to control the parallelism of shuffle-based operations in Spark SQL. Adjusting this value allows developers to match the level of parallelism to the specific data volume and cluster resources. Understanding how this setting impacts performance is essential for fine-tuning Spark applications to prevent task contention and ensure that the cluster is fully utilized during data-intensive stages of a job.

Exam trap

Candidates often confuse spark.sql.shuffle.partitions with spark.default.parallelism or source-specific reading options when trying to control shuffle stages.

97
Multi-Selectmedium

You are building a Structured Streaming job that reads from a Kafka source and writes to a Delta table. You need to ensure that the job can recover from failures and process data exactly once. Which TWO of the following are required to achieve exactly-once semantics? (Choose two.)

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

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

Why this answer

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

Exam trap

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

98
MCQhard

A data engineer is using Pandas API on Spark to process a large dataset. They call `psdf.to_pandas()` on a DataFrame that is 50 GB in size. What is the most likely outcome?

A.The operation succeeds only if the DataFrame is cached in memory beforehand.
B.The operation fails with an OutOfMemoryError because the entire dataset is collected into the driver's memory.
C.The operation completes successfully but takes a long time due to the large data volume.
D.The operation automatically partitions the data and returns a list of smaller pandas DataFrames.
AnswerB

`to_pandas()` collects all data from the distributed Spark DataFrame to the driver node as a single pandas DataFrame. A 50 GB dataset will almost certainly exceed the driver's memory capacity, leading to an OutOfMemoryError or a crash. This is a common pitfall when working with large datasets in Pandas API on Spark, as pandas is not distributed.

Why this answer

`to_pandas()` is a collect operation that brings all data to the driver as a pandas DataFrame. For large datasets, this will exceed driver memory and cause an OutOfMemoryError. It is intended for small results, not for converting massive distributed DataFrames.

Alternative approaches like sampling or aggregating before collection should be used.

Exam trap

The trap here is underestimating the memory implications of `to_pandas()` on large datasets, assuming that Spark's distributed nature protects the driver from memory issues.

99
MCQmedium

Which of the following correctly describes the relationship between a Spark Job and a Spark Stage?

A.A job is a single task that runs on one executor.
B.A stage is a set of parallel tasks that do not require a shuffle.
C.Stages are independent and never depend on each other.
D.A job must contain exactly one stage.
AnswerB

Stages are defined by shuffle boundaries. A single stage consists of a set of tasks that can be executed in parallel without any data exchange between them. Once data must be shuffled, a new stage is triggered, making this the correct definition of the relationship between tasks and stages.

Why this answer

In Spark's execution architecture, a job is composed of multiple stages. Stages are defined by shuffle boundaries. When a job is submitted, the DAG scheduler breaks it into stages based on wide transformations that require data movement across the network.

Understanding this hierarchy is essential for diagnosing performance issues, as it allows developers to identify which specific parts of their code lead to expensive shuffle operations, thus enabling better query optimization.

Exam trap

Students often confuse the hierarchy, mistakenly believing that a single stage contains multiple jobs or that tasks within a stage require network shuffles, overlooking that shuffle boundaries actually define stages.

100
MCQmedium

A developer is building a feature that must, for each `customer_id`, concatenate the distinct `product` values from many rows into a single comma-separated string. The result must contain each product only once per customer. Which single approach produces this result?

A.df.groupBy("customer_id").agg(concat_ws(",", collect_set("product")).alias("products"))
B.df.groupBy("customer_id").agg(concat_ws(",", collect_list("product")).alias("products"))
C.df.groupBy("customer_id").pivot("product").agg(count("product"))
D.df.select("customer_id", concat_ws(",", "product")).groupBy("customer_id").agg(first("product"))
AnswerA

collect_set gathers distinct product values into an array per customer, and concat_ws joins that array into a comma-separated string. Because collect_set deduplicates, each product appears once per customer, exactly matching the requirement to build a distinct list in a single aggregated column.

Why this answer

Producing a distinct comma-separated list per customer requires an aggregation that deduplicates first. collect_set returns a set of unique product values, which concat_ws then joins into one string. collect_list keeps duplicates, and the other approaches either collapse to a single value or reshape the data into columns, none of which yield the required distinct concatenated string.

Exam trap

The trap here is reaching for collect_list out of habit, forgetting that it preserves duplicates, whereas the distinct requirement demands collect_set.

101
MCQhard

A developer has two DataFrames, `orders` (columns `order_id`, `customer_id`) and `customers` (columns `customer_id`, `customer_name`). They need a result containing every order, with the matching customer name where one exists and null where the customer is not found. Which single join configuration guarantees this?

A.orders.join(customers, on="customer_id", how="left")
B.orders.join(customers, on="customer_id", how="right")
C.orders.join(customers, on="customer_id", how="inner")
D.orders.join(customers, on="customer_id", how="full")
AnswerA

A left join keeps every row from orders regardless of whether a match exists in customers. When no matching customer_id is found, the customer_name column is populated with null, precisely the behavior the scenario requires. The on parameter also avoids a duplicate customer_id column in the result.

Why this answer

Preserving all rows of the left DataFrame while enriching with matching values from the right is the definition of a left join. Unmatched orders retain their fields and receive nulls for customer_name. Inner, right, and full joins either drop unmatched orders or introduce extra customer-only rows, so they do not meet the requirement.

Exam trap

The trap here is treating an inner join as harmless because most orders do match, when even a single unmatched order is silently dropped and the requirement explicitly demands every order be kept.

102
MCQhard

A developer writes a Spark Connect client that creates a DataFrame, calls `df.collect()`, and then reuses the same DataFrame for a second `df.count()`. The cluster is remote. What happens on the second action?

A.The server automatically detects the repeated plan and returns a cached count without recomputation.
B.The client sends the logical plan again and the server re-executes the computation unless the DataFrame was cached.
C.The client raises an error because a DataFrame cannot be used for multiple actions in Spark Connect.
D.The client reuses cached results from the first action because the DataFrame is immutable.
AnswerB

Spark Connect is lazy: each action sends the accumulated logical plan to the server, which plans and executes it. Without an explicit `cache()` or `persist()`, the second action re-runs the computation from source. The DataFrame object on the client is only a plan builder; it holds no materialized data, so the server must recompute the result for `count()`.

Why this answer

In Spark Connect, a DataFrame is a client-side logical plan builder, not a materialized dataset. Each action serializes the current plan and sends it to the server, which plans and executes it. Because no cache or persist was invoked, the second action recomputes the result from the source.

Reusing the DataFrame object is valid; it simply does not reuse prior results.

Exam trap

The trap here is assuming that reusing the same DataFrame object across actions reuses results, when results are only reused if the DataFrame was explicitly cached or persisted on the server.

103
MCQeasy

A developer runs a PySpark job that joins a 10 GB DataFrame with a 50 MB lookup DataFrame. The job takes far longer than expected, and the Spark UI shows a SortMergeJoin with a large shuffle read and write for both sides. The developer wants the smallest change that most improves performance. Which action should the developer take?

A.Increase spark.sql.shuffle.partitions so the sort-merge join shuffle uses more tasks and completes faster.
B.Broadcast the 50 MB lookup DataFrame using broadcast() or by raising spark.sql.autoBroadcastJoinThreshold above 50 MB.
C.Repartition the 50 MB lookup DataFrame by the join key before the join to co-locate matching rows.
D.Cache the 10 GB DataFrame with persist(StorageLevel.MEMORY_AND_DISK) before the join.
AnswerB

Broadcasting the small side sends a copy of the 50 MB lookup to every executor and converts the join to a BroadcastHashJoin, eliminating the shuffle of both DataFrames. This is the minimal change that removes the expensive sort-merge shuffle and typically yields the largest speedup for a small-to-large join.

Why this answer

The Spark UI shows a SortMergeJoin with large shuffle on both sides, which is unnecessary when one input is only 50 MB. Broadcasting the small lookup table turns the join into a BroadcastHashJoin, removes the shuffle and sort of both DataFrames, and requires only a hint or a threshold adjustment.

Exam trap

The trap here is tuning shuffle partitions or caching when the join strategy itself is the problem and a broadcast would eliminate the shuffle entirely.

104
MCQeasy

A data engineer has a batch DataFrame `df` with a `status` string column and wants to keep only rows where `status` equals "active". The engineer wants the filter applied as early as possible in the plan and does not want a shuffle. Which operation best fits this requirement?

A.df.select("status").distinct()
B.df.where(col("status") == "active")
C.df.groupBy("status").count().filter(col("status") == "active")
D.df.orderBy("status").filter(col("status") == "active")
AnswerB

where() is an alias for filter() and applies a row-wise predicate as a narrow transformation, so no shuffle occurs and each partition can drop non-matching rows locally. The predicate is also available to the Catalyst optimizer for predicate pushdown into compatible sources. This satisfies both the correctness requirement of keeping only active rows and the performance requirement of avoiding a shuffle.

Why this answer

Filtering rows by a column predicate is exactly what where() or filter() does, and it is a narrow transformation that runs within each partition without exchanging data. Aggregation, distinct, and sorting all introduce shuffles that the scenario rules out. The where() call also exposes the predicate to Catalyst so it can be pushed into the source when supported, making it both correct and efficient.

Exam trap

The trap here is reaching for a groupBy or distinct pattern to isolate a value, when a simple row predicate already returns the desired records without any shuffle.

105
MCQhard

A Spark job reads a large Parquet dataset, performs a groupBy on a high-cardinality column, and writes the result to a Delta table. The job fails with a FetchFailedException on a particular executor. The Spark UI shows that the executor had sufficient memory but the shuffle fetch failed due to a connection reset. Which configuration change is most likely to resolve this issue?

A.Increase spark.executor.memory to prevent the executor from being killed during shuffle.
B.Increase spark.reducer.maxSizeInFlight to allow larger shuffle blocks to be fetched.
C.Set spark.shuffle.io.maxRetries and spark.shuffle.io.retryWait to higher values to handle transient network issues.
D.Set spark.sql.adaptive.enabled=false to disable adaptive query execution and avoid shuffle re-computation.
AnswerC

FetchFailedException due to connection reset often indicates transient network problems or shuffle service timeouts. Increasing shuffle I/O retries and retry wait allows the reducer to retry fetching blocks after a failure, improving resilience. This directly addresses the connection reset by giving the fetch more attempts and time to succeed.

Why this answer

A FetchFailedException with connection reset during shuffle fetch is typically caused by transient network issues or shuffle service timeouts. Increasing spark.shuffle.io.maxRetries and spark.shuffle.io.retryWait makes the shuffle fetch more resilient by retrying failed attempts. This is the most direct configuration change to handle temporary network glitches without altering the overall job logic or resource allocation.

Exam trap

The trap here is assuming that memory or query planning changes will fix a shuffle fetch failure, when the error is network-related and requires retry tuning.

106
MCQmedium

A data analyst is using the Pandas API on Spark to compute summary statistics. They call `psdf.describe()` on a large DataFrame and notice the job takes much longer than expected. They want to understand why this operation is more expensive than a similar operation on a small local pandas DataFrame. What is the primary reason?

A.`describe()` is a lazy operation that only builds a plan, so the delay is from plan construction rather than execution.
B.`describe()` triggers a full scan and computes multiple aggregations that may require shuffling data across partitions.
C.`describe()` caches the DataFrame in memory by default, and the caching step is what takes the extra time.
D.`describe()` converts the entire DataFrame to a local pandas DataFrame on the driver before computing statistics.
AnswerB

`describe()` computes count, mean, stddev, min, max, and percentiles for numeric columns. Percentiles in Spark are computed with approximate algorithms that require a shuffle to gather distribution information across partitions. The full scan plus the multi-aggregation plan and the shuffle for quantiles explain the increased runtime on a distributed DataFrame.

Why this answer

`describe()` on a Pandas API on Spark DataFrame runs a distributed Spark job that scans all rows and computes multiple aggregates. Percentiles require approximate quantile algorithms that shuffle data across partitions. This is fundamentally more expensive than local pandas, which operates on in-memory data on a single machine without network shuffles.

Exam trap

The trap here is assuming that pandas-like syntax implies local, in-memory execution rather than distributed Spark jobs.

107
MCQmedium

A developer is building a Python application that connects to a Databricks cluster using Spark Connect. The application uses the `databricks-connect` package and is configured with the cluster ID and authentication credentials. During a test run, the developer calls `spark.sql("SELECT * FROM sales")` and then `df.show()`. What happens when the `show()` action is executed?

A.The client uses a local SparkContext to connect to the cluster's driver and executes the query as if it were local.
B.The client downloads the entire `sales` table to the local machine and executes the SQL query using a local Spark session.
C.The client sends the logical plan to the Spark Connect server, which executes it on the cluster and returns the result rows back to the client.
D.The client compiles the SQL query into a JVM bytecode and sends it to the cluster for execution.
AnswerC

Spark Connect uses a client-server architecture where the client builds an unresolved logical plan and sends it via gRPC to the Spark Connect server running on the cluster. The server optimizes and executes the plan, then streams results back to the client. This decouples the client from the driver, enabling remote execution and improved stability.

Why this answer

In Spark Connect, the client builds a logical plan and sends it to the Spark Connect server via gRPC. The server executes the plan on the cluster and returns results. The client does not download data or execute locally.

This architecture decouples the client from the driver, enabling remote execution and improved stability.

Exam trap

The trap here is assuming that Spark Connect executes queries locally or downloads data, when it actually delegates execution to the remote server.

108
MCQhard

You are using Pandas API on Spark to process a large dataset. You have a Pandas-on-Spark DataFrame `psdf` and you apply a custom Python function using `psdf.apply(func, axis=1)`. The function is computationally intensive and you notice that the job is running slowly with many tasks. What is the most likely reason for the performance issue?

A.The `apply` function with `axis=1` is executed row-by-row using a Python UDF, which incurs high serialization and execution overhead, and prevents Spark from optimizing the query.
B.The `apply` function with `axis=1` is executed in a distributed manner, but it requires a full shuffle of the data before applying the function.
C.The `apply` function with `axis=1` forces the entire DataFrame to be collected to the driver, and then applies the function locally, causing memory issues.
D.The `apply` function with `axis=1` is not supported in Pandas API on Spark and falls back to a single-node pandas execution, causing a bottleneck.
AnswerA

This is correct. When you use `apply` with `axis=1`, Pandas API on Spark translates it into a Python UDF that processes each row individually. This involves serializing each row from the JVM to Python, executing the function, and deserializing the result. This overhead is significant and prevents Spark's Catalyst optimizer from optimizing the logic. It is generally recommended to avoid row-wise `apply` for large datasets.

Why this answer

Using `apply` with `axis=1` in Pandas API on Spark often results in poor performance because it is implemented via a Python UDF that processes rows one by one. This incurs high serialization overhead and prevents Spark's optimizer from improving the plan. For large datasets, it is better to use vectorized operations or built-in functions.

Exam trap

The trap here is assuming that `apply` is as efficient as in pandas, but in a distributed setting, row-wise Python UDFs are slow.

109
MCQmedium

A streaming DataFrame job on Databricks writes to a Delta table every 10 seconds. Over several hours, the number of files in the target directory grows into the hundreds of thousands, and downstream reads slow dramatically. The job uses foreachBatch with a write that produces many small files per micro-batch. Which action should be taken to reduce the small-files problem for this streaming write?

A.Call repartition(1) on the streaming DataFrame before the write in foreachBatch.
B.Configure the write with optimizeWrite enabled (or set the target table property delta.autoOptimize.optimizeWrite = true) so Spark coalesces output files per partition before committing.
C.Increase the trigger interval from 10 seconds to 10 minutes so each micro-batch writes more data.
D.Set spark.sql.shuffle.partitions to 1 so all output is written by a single task.
AnswerB

Optimize Write coalesces the many small files produced by each micro-batch into fewer, larger files before they are committed to the Delta table. This directly attacks file proliferation at write time, reducing the file count and improving downstream read performance without changing the streaming logic.

Why this answer

The streaming job creates many small files because each micro-batch writes multiple files per partition. Optimize Write coalesces those files into fewer, larger ones at commit time, directly reducing file proliferation and improving downstream read performance without sacrificing parallelism or increasing latency.

Exam trap

The trap here is reaching for repartition(1) or a single shuffle partition to reduce file count, which fixes file count by destroying write parallelism.

110
Multi-Selecthard

You are processing a large, highly skewed PySpark DataFrame in Databricks and want to optimize a forthcoming join operation against a small lookup dimension table. Which TWO strategies are valid and effective DataFrame API techniques to optimize this join performance? (Choose TWO)

Select 2 answers
A.Use the broadcast() function on the small dimension table to distribute it to all worker nodes and avoid a shuffled hash join.
B.Apply a broadcast join hint directly to the large skewed fact table to force executor nodes to cache its partitions in memory.
C.Introduce a salt column with random integers to the join keys of both DataFrames to evenly distribute skewed keys across multiple tasks.
D.Increase the spark.sql.shuffle.partitions configuration to an extremely high number like 10000 to eliminate data skew entirely.
E.Convert the large DataFrame into a local Pandas DataFrame using toPandas() to perform the join operations locally on the driver node.
AnswersA, C

broadcast() ships the small dimension table to every executor, converting the join into a broadcast hash join and eliminating the shuffle of the large skewed DataFrame. This directly avoids the expensive shuffled hash join that skew would otherwise make prohibitively slow.

Why this answer

Optimizing joins involving skewed datasets requires avoiding shuffles for small tables or breaking up skewed keys. Broadcasting the small table eliminates the shuffle phase entirely, while salting keys distributes skewed join keys across multiple tasks, preventing memory bottlenecks on single executors during heavy shuffle exchanges.

Exam trap

Candidates often assume broadcasting large fact tables or relying purely on automatic AQE join optimization solves all skew problems without manual intervention.

111
MCQhard

A developer needs to join `orders` (large) with `customers` (small, a few thousand rows) on `customer_id`. They want to broadcast the small table and confirm the broadcast actually took effect. Which combination of actions is correct?

A.Increase `spark.sql.autoBroadcastJoinThreshold` to a very large value and rely on the optimizer without checking the plan
B.Call `customers.repartition(1)` before the join so the small table is in one partition
C.Call `orders.join(broadcast(customers), "customer_id")` and verify by inspecting the physical plan for `BroadcastHashJoin`
D.Use `orders.join(customers.hint("shuffle_hash"), "customer_id")` and check the plan for `ShuffledHashJoin`
AnswerC

Wrapping the small DataFrame in `broadcast()` hints Spark to use a broadcast join, and inspecting the physical plan via `explain()` confirms whether `BroadcastHashJoin` appears. If the plan shows it, the hint was honored; if not, a shuffle join is occurring. This is the correct way to force and then verify a broadcast join.

Why this answer

Broadcasting a small table avoids shuffling the large one, and the `broadcast()` hint is the explicit way to request it. Verification matters because hints can be ignored when the table exceeds the threshold or when the join type is unsupported. Inspecting the physical plan for `BroadcastHashJoin` is the standard confirmation, making the combination of hint plus plan inspection the correct answer.

Exam trap

The trap here is assuming that increasing the broadcast threshold or repartitioning the small table guarantees a broadcast join, when only the plan confirms it.

112
MCQmedium

A developer has a DataFrame `events` with columns `user_id`, `event_type`, and `payload` (a JSON string). They need to extract the `device` field from `payload` into a new column without changing the other columns or the row count. Which approach is correct?

A.events.withColumn("device", explode(from_json("payload", schema)))
B.events.select("user_id", "event_type", from_json("payload", schema).getField("device"))
C.events.rdd.map(lambda r: r.payload["device"]).toDF()
D.events.withColumn("device", get_json_object("payload", "$.device"))
AnswerD

`get_json_object` extracts a scalar value from a JSON string using a JSONPath expression. Wrapping it in `withColumn` adds the `device` column while preserving all existing columns and the row count. This is the simplest correct approach when only a single field is needed from the JSON string without parsing the entire document.

Why this answer

Extracting a single scalar from a JSON string column is exactly what `get_json_object` does, and `withColumn` adds the result without altering other columns or row counts. Parsing the whole payload with `from_json` is unnecessary when only one field is needed, and any use of `explode` or RDD conversion changes the shape of the DataFrame.

Exam trap

The trap here is reaching for `from_json` and `explode` reflexively, when a single-field extraction with `get_json_object` is simpler and preserves the DataFrame shape.

113
MCQmedium

You are processing a streaming dataset of sensor readings. You need to calculate the average temperature every 10 minutes, allowing data to arrive up to 2 minutes late. Which windowing approach correctly handles this requirement in Structured Streaming?

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

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

Why this answer

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

Exam trap

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

114
MCQeasy

A data engineer wants to run a Structured Streaming query that reads from a Kafka topic and writes aggregated counts to a console sink for debugging. The query uses a grouping aggregation on a tumbling event-time window. Which output mode must be used so that only rows that changed since the last trigger are emitted?

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

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

Why this answer

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

Exam trap

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

115
MCQhard

A developer is building a Structured Streaming job that reads from a Kafka topic and writes to a Delta table. The job must handle late data up to 15 minutes and ensure that aggregations are updated correctly. The developer adds a watermark of 15 minutes on the event time column. What is the effect of this watermark on the aggregation state and output?

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

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

Why this answer

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

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

Exam trap

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

116
MCQmedium

Which component in the Spark architecture is responsible for scheduling tasks and managing the execution of jobs on the cluster?

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

The Driver creates the SparkSession, translates transformations and actions into a DAG, and orchestrates the execution of tasks across the worker nodes. It maintains information about the state of the executors and ensures that data processing occurs efficiently by optimizing the execution plan before dispatching it to executors.

Why this answer

The Driver process is the central coordinator in Spark. It hosts the SparkContext, which converts user code into a Directed Acyclic Graph (DAG) and schedules tasks across the executors. Understanding this role is vital because it explains why the Driver can become a bottleneck if it handles too much data locally or manages excessive partitions, directly impacting the overall job latency and system stability in distributed environments.

Exam trap

Candidates frequently mix up the roles of the Driver and the Executors, incorrectly attributing task scheduling and job coordination responsibilities to worker nodes.

117
MCQmedium

Which file format is best suited for performance-critical Spark applications that require efficient schema enforcement and column pruning?

A.CSV
B.JSON
C.Parquet
D.XML
AnswerC

Parquet is a highly optimized columnar format that allows Spark to skip reading unnecessary columns and push down filter operations to the storage layer. This minimizes I/O and CPU overhead, making it the industry standard for performance-critical, large-scale data processing workflows within Databricks and the Apache Spark ecosystem.

Why this answer

Parquet is a columnar storage format that natively supports predicate pushdown and column pruning, allowing Spark to read only the necessary data from disk. This drastically reduces I/O throughput and improves query performance significantly. Understanding why columnar formats are superior to row-based formats for analytics is a foundational concept for Databricks developers designing high-performance data lakes and efficient ETL pipelines at scale.

Exam trap

Candidates often select CSV or JSON formats, confusing human readability with performance-critical analytics features like native columnar pruning.

118
MCQmedium

You are developing a Structured Streaming job that reads from a Delta table and writes to another Delta table. You need to ensure that the streaming query can recover from failures and continue processing without data loss or duplication. Which of the following must be configured?

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

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

Why this answer

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

Therefore, the checkpoint location is the essential configuration.

Exam trap

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

119
MCQmedium

What happens when a Spark job triggers a 'shuffle' operation during execution?

A.Data is automatically cached in memory on all executors.
B.Data is redistributed across the executors based on key distribution.
C.The Driver node collects all data to perform the operation.
D.The job terminates immediately due to network bandwidth limits.
AnswerB

Shuffles are necessitated by wide transformations where data needs to be aggregated or joined. Spark moves data partitions across the network so that all values for a specific key reside on the same executor, ensuring that the subsequent operation can perform the calculation correctly across the entire dataset.

Why this answer

A shuffle involves re-partitioning data across the cluster, requiring significant network I/O and disk activity. Recognizing this is crucial for performance tuning because shuffles are often the most expensive parts of a Spark job. By understanding how data is redistributed, developers can avoid unnecessary shuffles, choose better join strategies, and configure partition counts to reduce latency and prevent bottlenecks that occur when data must be moved between executors.

Exam trap

Candidates confuse a shuffle with a simple broadcast join or a partition re-balance, failing to identify that a shuffle specifically involves network-wide data redistribution across executors.

120
MCQhard

Refer to the exhibit. What is the best way to resolve this error?

A.Increase the spark.driver.maxResultSize setting.
B.Write the results to a data lake instead of collecting them.
C.Disable the Spark broadcast join threshold.
D.Increase the number of executor instances.
AnswerB

Writing results to persistent, distributed storage (like Parquet files on S3/ADLS) is the architecture-correct way to handle large outputs. This avoids sending all data to the driver, allowing the Spark executors to perform the work in parallel and keeping the driver's memory footprint small and stable throughout the execution.

Why this answer

This error occurs when the result set returned to the driver from the executors exceeds the limit defined by `spark.driver.maxResultSize`. The best practice is to stop trying to bring huge datasets back to the driver, and instead write the results to a distributed storage system like S3 or ADLS. This prevents the driver from becoming a memory bottleneck and ensures the application remains scalable for large-scale data processing tasks.

Exam trap

Candidates often attempt to increase 'spark.driver.maxResultSize' to fix the error, rather than changing their code pattern to avoid bringing massive data back to the driver node.

121
MCQeasy

A Spark application is running on a Databricks cluster with 3 worker nodes, each having 4 cores. The application uses the default configuration. How many tasks can run concurrently across the cluster?

A.4
B.3
C.7
D.12
AnswerD

In Spark, each task runs on one core. The total number of cores across all executors determines the maximum number of concurrent tasks. With 3 worker nodes, each with 4 cores, there are 12 cores available. Assuming default configuration where each node runs one executor using all cores, the cluster can run 12 tasks concurrently.

Why this answer

The maximum number of concurrent tasks in a Spark cluster is equal to the total number of cores available across all executors. In this scenario, with 3 worker nodes each having 4 cores, the total is 12 cores. Therefore, up to 12 tasks can run in parallel, assuming each task uses one core and there is no dynamic allocation or other constraints.

Exam trap

The trap here is summing the number of nodes and cores instead of multiplying them, or forgetting that concurrency is based on total cores across the cluster.

122
MCQeasy

A developer is using Spark on Databricks and wants to monitor the progress of a job. They need to understand how the driver coordinates with executors. Which component is responsible for scheduling tasks onto executors and tracking their status?

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

The Task Scheduler is responsible for scheduling individual tasks onto executors based on data locality and resource availability. It tracks task status and retries failed tasks. This component directly manages the execution of tasks on the cluster, making it the correct answer for coordinating with executors.

Why this answer

The Task Scheduler in Spark is responsible for assigning tasks to executors and monitoring their execution. It works closely with the DAG Scheduler, which breaks the job into stages, but the Task Scheduler handles the actual task-level scheduling and status tracking. This makes it the component that directly coordinates with executors.

Exam trap

The trap here is confusing the DAG Scheduler with the Task Scheduler, as both are involved in scheduling but at different levels of granularity.

123
MCQeasy

A developer notices that a PySpark DataFrame transformation chain runs a full scan of a Delta table each time a new action is invoked, even though the source data has not changed. The developer wants to persist the intermediate DataFrame in memory across actions. Which method should be used?

A.Call DataFrame.collect() after each transformation to materialize the results on the driver.
B.Call DataFrame.cache() before the first action so the DataFrame is stored in memory on first computation.
C.Set spark.sql.adaptive.enabled to true so the optimizer reuses prior scan results automatically.
D.Call DataFrame.checkpoint() to truncate the lineage and write the data to a reliable file system.
AnswerB

cache() marks the DataFrame for in-memory persistence with the default MEMORY_AND_DISK storage level, so the first action materializes it and subsequent actions reuse the cached partitions instead of rescanning the Delta table. This directly addresses the repeated full scans described in the scenario and is the idiomatic way to persist across multiple actions.

Why this answer

cache() persists the DataFrame across actions using the default MEMORY_AND_DISK level, so the first action computes and stores partitions and later actions read from the cache rather than rescanning the Delta table. Checkpointing writes to disk and cuts lineage but does not keep data in memory, while AQE and collect() do not provide reusable in-memory persistence.

Exam trap

The trap here is conflating checkpointing, which truncates lineage to disk, with caching, which keeps the materialized DataFrame available in memory for repeated actions.

124
MCQeasy

Which action should be taken to optimize a Spark application that performs multiple operations on the same DataFrame and shows evidence of redundant re-computations in the Spark UI DAG visualization?

A.Increase the number of shuffle partitions.
B.Use the cache() or persist() method on the DataFrame.
C.Implement a custom partitioner for all joins.
D.Convert the DataFrame to a RDD.
AnswerB

Caching stores the materialized result of a DataFrame. When the application accesses the data again, Spark retrieves it from the cache rather than re-computing the entire lineage. This is the correct technique for eliminating redundant work when multiple downstream operations depend on the same intermediate DataFrame.

Why this answer

Persisting (or caching) a DataFrame instructs Spark to store the computed results of the DataFrame in memory or on disk. This is vital when the same DataFrame is accessed multiple times across different stages, as it prevents Spark from re-executing the entire lineage graph from the source, significantly reducing processing time and resource consumption in complex, multi-step analytical pipelines.

Exam trap

Candidates often mistake cache() for an action. They forget that cache() is a lazy transformation and must be followed by an action like count() or write() to actually materialize the data in memory.

125
MCQeasy

A developer has a DataFrame `raw` and needs to permanently persist it as Parquet partitioned by `region`, overwriting any existing data at that path, without registering it in the metastore. Which call achieves this?

A.raw.rdd.saveAsTextFile("/data/out")
B.raw.write.format("parquet").saveAsTable("out")
C.raw.createOrReplaceTempView("out").write.parquet("/data/out")
D.raw.write.mode("overwrite").partitionBy("region").parquet("/data/out")
AnswerD

This uses the DataFrameWriter with `mode("overwrite")` to replace existing data, `partitionBy("region")` to lay out Hive-style directories, and the `parquet` format sink to write files. It writes directly to the filesystem path without touching the metastore, which is precisely the requirement stated in the scenario.

Why this answer

The DataFrameWriter path with `mode("overwrite")`, `partitionBy("region")`, and `.parquet` writes typed Parquet files into per-region directories and replaces prior contents in one step. Because it targets a filesystem path rather than a table name, no metastore entry is created, satisfying the "without registering it in the metastore" constraint.

Exam trap

The trap here is confusing `saveAsTable`, which registers metastore metadata, with a direct path-based `write`, which does not.

126
MCQmedium

What is the primary benefit of the Catalyst Optimizer in the Spark SQL architecture?

A.It manages the physical cluster resources.
B.It automatically generates the most efficient execution plan.
C.It serializes data for storage in parquet files.
D.It performs automatic garbage collection on executors.
AnswerB

Catalyst uses rules to rewrite the logical plan, applying techniques like filter pushdown and join reordering. This ensures that the physical execution plan is as efficient as possible, reducing unnecessary data scanning and shuffling, which leads to significantly faster job completion times in complex SQL and DataFrame operations.

Why this answer

Catalyst optimizes logical plans through rule-based and cost-based transformations, such as predicate pushdown and constant folding. By simplifying the query plan before it reaches the physical execution layer, it drastically reduces the amount of data processed. This is critical for performance because it minimizes I/O and CPU usage, ensuring that queries are executed using the most efficient physical operators possible, which is essential for scaling across large datasets.

Exam trap

Many test-takers confuse the Catalyst Optimizer with cluster resource managers or physical execution schedulers, failing to recognize its specific role in query plan optimization.

127
MCQhard

Which TWO of the following techniques effectively reduce the shuffle volume in a Databricks Spark job?

A.Apply filter transformations as early as possible in the DataFrame lineage.
B.Increase the memory allocated to the Spark driver.
C.Select only necessary columns before performing wide transformations.
D.Enable dynamic allocation of executors.
E.Set the spark.sql.shuffle.partitions to 1.
AnswerA, C

Early filtering, or predicate pushdown, reduces the total number of rows processed. By discarding irrelevant records before any shuffle operation occurs, you significantly decrease the amount of data written to disk and transferred over the network, leading to faster execution times and lower resource consumption during the shuffle phase.

Why this answer

Reducing shuffle volume is essential for performance, as shuffling moves data across the network, which is the most expensive operation in Spark. Techniques like predicate pushdown and column pruning minimize the amount of data read and processed before the shuffle phase occurs. Mastering these optimizations ensures that jobs remain scalable as datasets grow, preventing network congestion and I/O saturation during complex transformations or aggregate operations on large distributed DataFrames.

Exam trap

Candidates mistakenly believe that adding a repartition command or caching early reduces shuffle volume, confusing memory persistence with network transfer reduction.

128
MCQhard

A data engineer is using Spark Connect to run a job on a Databricks cluster. They notice that when they call `df.count()`, the operation takes longer than expected. They suspect that the client is transferring data unnecessarily. Which statement best explains the data transfer behavior of `df.count()` in Spark Connect?

A.The server executes the count and returns only the count value to the client, minimizing data transfer.
B.The client sends a count request to the server, and the server returns a lazy evaluation plan that the client must execute.
C.The client sends the entire DataFrame to the server, which then counts the rows and returns the result.
D.The client downloads all rows to count them locally, then discards the data.
AnswerA

`df.count()` is an action that triggers execution on the server. The server computes the count and returns a single integer to the client. Only the result is transferred, not the underlying data. This is efficient and typical for aggregate actions in Spark Connect.

Why this answer

In Spark Connect, actions like `count()` are executed on the server. The client sends the logical plan, and the server computes the result, returning only the scalar value. This minimizes data transfer and leverages the cluster's processing power.

The client does not download data for aggregation.

Exam trap

The trap here is assuming that the client downloads data to perform aggregation, when actually the server computes and returns only the result.

129
MCQhard

A data scientist is writing a Spark Connect application that requires custom user-defined functions (UDFs). How are UDFs handled when executing code through Spark Connect?

A.Python UDFs are executed locally on the client machine before sending results to the server.
B.Custom UDFs are completely unsupported in Spark Connect because client environments are strictly isolated.
C.User-defined functions are serialized and transmitted to the server where they execute on the cluster.
D.UDF definitions must be pre-installed as wheel files on every cluster worker node prior to execution.
AnswerC

Spark Connect's thin client cannot execute UDFs locally; the client serialises the function and ships it to the server, where it is deserialised and run on the cluster's executors. This preserves distributed execution despite the decoupled client-server architecture.

Why this answer

Spark Connect transmits Python UDFs by serializing the function and sending its bytecode definitions over the gRPC channel to the server. The server then deserializes and executes these functions within the cluster environment, ensuring compatibility and secure execution without needing identical local Python binary environments on client machines.

Exam trap

Test-takers frequently assume that Spark Connect executes custom UDFs locally on the client machine, confusing client-side code definition with server-side execution.

130
MCQhard

What is the primary function of the 'Shuffle Service' in a Spark cluster when using dynamic allocation?

A.To compress shuffle data before it is written to the disk.
B.To allow executors to retrieve shuffle data from removed executors.
C.To rebalance data partitions across the cluster during a shuffle.
D.To increase the speed of network transfers during a shuffle.
AnswerB

The External Shuffle Service runs as a separate process on each node, independent of the Spark executors. When an executor is removed, its shuffle files remain accessible through the service, preventing the need to recompute shuffle stages and maintaining application stability during scaling events.

Why this answer

The External Shuffle Service allows executors to be decommissioned without losing shuffle files needed by downstream stages. In dynamic allocation, Spark frequently scales the number of executors based on workload. Without this service, if an executor holding shuffle map output files were terminated, the downstream tasks would fail because their input data would be permanently lost, forcing expensive recomputations of the upstream shuffle stages.

Exam trap

Candidates assume the Shuffle Service is for performance speed, ignoring its critical role in fault tolerance when executors are dynamically removed during a job.

131
MCQmedium

What is the purpose of the 'Broadcast Variable' in the Spark architecture?

A.To share mutable state across all executors.
B.To send a read-only variable to every node efficiently.
C.To aggregate intermediate results from tasks.
D.To partition data for balanced shuffle operations.
AnswerB

By using an efficient peer-to-peer distribution mechanism, the broadcast variable ensures each node receives the data only once. This is far more efficient than including the variable in the task closure, which would send the data repeatedly, creating a massive bandwidth bottleneck for every single task launched.

Why this answer

Broadcast variables allow the Driver to send a read-only copy of a large variable to every executor, rather than sending a copy with every single task. This drastically reduces network traffic and memory usage when joining small lookup tables with large datasets. Understanding this is critical for performance, as it prevents the redundant transmission of data and optimizes join operations, ensuring that the cluster remains efficient and avoids network congestion during large-scale operations.

Exam trap

Test-takers frequently confuse broadcast variables with accumulators or normal shuffle joins, incorrectly thinking they allow workers to write and synchronize shared mutable state back to the driver.

132
MCQmedium

A developer is building a Spark Connect application that runs on a laptop and connects to a remote Databricks cluster. During development, the laptop loses network connectivity for a few minutes while a long-running DataFrame transformation is executing. The developer notices the local Python process raises a gRPC error and the job is no longer tracked. Which statement best explains this behavior?

A.Spark Connect automatically switches to a local SparkSession fallback when the remote connection fails, preserving the job.
B.The disconnect only affects result collection; the remote cluster continues the job to completion and stores the result for later retrieval by any client.
C.Spark Connect uses a gRPC channel between the client and the Spark server; if that channel is disrupted, the client loses the logical plan and the server may cancel the associated execution.
D.The local client retains a full copy of the RDD lineage and can recompute the result locally after reconnecting, so no work is lost.
AnswerC

Spark Connect decouples the client from the driver through a gRPC-based protocol. The client sends unresolved logical plans over this channel and holds a session handle. When network connectivity drops, the gRPC stream breaks, so the client can no longer track or control the running query, and the server may abort the associated execution due to the lost session.

Why this answer

The correct answer reflects Spark Connect's client-server architecture, where a gRPC channel carries logical plans and control messages. A network interruption breaks that channel, so the client cannot track or manage the running query, and the server may cancel the execution because the session is no longer reachable. This is different from classic Spark, where the driver and client are co-located and a local network blip does not sever the control path.

Exam trap

The trap here is assuming that Spark Connect behaves like a classic local SparkSession, where losing the network does not affect a running job because the driver is local.

133
MCQhard

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

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

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

Why this answer

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

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

Exam trap

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

134
MCQmedium

A developer is writing a Spark Connect application that runs on a local laptop and connects to a Databricks cluster. The application defines a Python function and registers it with `spark.udf.register` for use inside a `select` expression. When the code runs, the function executes on the server. Which statement describes how the UDF is handled in this scenario?

A.The UDF is rejected because Spark Connect does not support user-defined functions of any kind.
B.The UDF is serialized and shipped to the Spark Connect server, where it is deserialized and executed within the server's Python worker processes.
C.The UDF is converted into a SQL expression by the client and embedded directly in the query plan without server-side Python execution.
D.The UDF runs only on the client machine, and its outputs are transmitted back to the server as literal values.
AnswerB

Spark Connect serializes the UDF definition and sends it to the server, where it is deserialized and executed in the server's Python workers. This preserves the familiar PySpark UDF programming model while keeping execution on the cluster, so the local client process does not need to run the function itself.

Why this answer

Spark Connect keeps the client thin by serializing operations, including UDF definitions, and sending them to the server for execution. The server deserializes the UDF and runs it in its Python worker processes, so distributed execution and data locality remain on the cluster while the developer keeps the familiar PySpark UDF API.

Exam trap

The trap here is assuming Spark Connect executes UDFs locally on the client, when in fact the UDF is serialized to the server for execution.

135
MCQeasy

In a Spark application running on a Databricks cluster, the driver program creates a SparkSession and defines a series of transformations. When an action is triggered, the driver requests resources from the cluster manager. Which component is responsible for negotiating and acquiring these resources on behalf of the Spark application?

A.The Task Scheduler
B.The DAG Scheduler
C.The cluster manager
D.The Executor
AnswerC

The cluster manager (e.g., YARN, Kubernetes, or Databricks' internal manager) is responsible for allocating resources such as executors to the Spark application. The driver communicates with the cluster manager to request containers or pods, which then launch executors. This is a core part of Spark's architecture.

Why this answer

The cluster manager is the component that allocates resources to the Spark application. The driver requests resources, and the cluster manager launches executors accordingly. The DAG Scheduler and Task Scheduler operate within the driver to plan and schedule tasks, while executors run the tasks but do not negotiate resources.

Exam trap

The trap here is confusing the role of the cluster manager with that of the driver's internal schedulers, which plan tasks but do not acquire cluster resources.

136
MCQmedium

A data engineer is writing a PySpark job that must read Parquet data, apply transformations, and then write the result. They want to avoid recomputation if the resulting DataFrame is referenced multiple times in later stages. Which action best accomplishes this?

A.Call df.checkpoint() without setting a checkpoint directory
B.Call df.count() to force computation before reuse
C.Call df.rdd to convert the DataFrame to an RDD
D.Call df.cache() or df.persist() before the DataFrame is reused
AnswerD

Caching or persisting materializes the DataFrame's computed partitions so later actions reuse them instead of re-executing the full lineage. This directly addresses the goal of avoiding recomputation when the same DataFrame feeds multiple downstream operations, and it remains lazy until an action triggers the actual caching.

Why this answer

Marking the DataFrame with `cache` or `persist` stores its computed partitions so subsequent actions reuse them rather than replaying the full lineage. Forcing a count, converting to an RDD, or checkpointing without a configured directory either fails to store reusable data or introduces errors and overhead. Caching is the standard mechanism for avoiding recomputation across multiple downstream uses.

Exam trap

The trap here is believing that triggering an action such as `count()` persists the data, when in fact only `cache` or `persist` retains the computed partitions for reuse.

137
Multi-Selectmedium

You are designing a Structured Streaming job that must read from a file source and write to a Delta table. You want the job to be resilient to failures and to continue processing only new files after a restart. Which two actions should you take? (Choose two.)

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

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

Why this answer

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

Exam trap

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

138
MCQmedium

An enterprise data engineering team is migrating legacy PySpark client applications to use Spark Connect to improve client stability and isolate resource consumption. A developer initializes the Spark session pointing to a remote cluster. Which specific mechanism does Spark Connect use to communicate execution plans between the client application and the server cluster?

A.It leverages standard JDBC protocol connections over secure sockets to transmit compiled execution plans.
B.It uses Apache Arrow flight protocol directly for executing all distributed transformations across the worker nodes.
C.It serializes DataFrame execution plans into protocol buffers and transmits them over a gRPC communication channel.
D.It establishes a Py4J gateway bridge over a dedicated network socket to invoke remote JVM methods seamlessly.
AnswerC

Spark Connect decouples the client from the driver by encoding unresolved logical plans as protocol buffers, then streaming them over gRPC to the server, which handles planning and execution. This satisfies the isolation requirement, since the thin client holds no JVM or cluster resources.

Why this answer

Spark Connect implements a gRPC-based client-server architecture. The client serializes dataframe transformations into protocol buffers, which are transmitted over gRPC streams to the Spark driver server. This decoupled protocol ensures that memory pressure on the client does not directly crash the driver, providing better isolation and resource management in modern distributed data architectures.

Exam trap

Candidates often confuse Spark Connect with traditional JDBC/ODBC or direct PySpark driver-worker communication, overlooking the specific use of protocol buffers over a gRPC transport layer.

139
MCQhard

A data engineer has a DataFrame `orders` with columns `order_id`, `customer_id`, and `amount`, and a small lookup DataFrame `tiers` with `customer_id` and `tier`. The engineer wants to attach the tier to every order. Some orders have a customer_id that is not present in tiers, and those orders must still appear with a null tier. Which operation produces this result?

A.orders.join(tiers, on="customer_id", how="inner")
B.orders.union(tiers)
C.orders.join(tiers, on="customer_id", how="right")
D.orders.join(tiers, on="customer_id", how="left")
AnswerD

A left outer join preserves every row from the left DataFrame and fills columns from the right with nulls when no match exists. That is exactly the required behavior: all orders remain, and orders whose customer_id is missing from tiers get a null tier. The join key is shared, so the output contains one customer_id column plus order_id, amount, and tier.

Why this answer

The requirement is to keep all orders while enriching them with tier where available, filling null otherwise. A left outer join does precisely this: it preserves the left side unconditionally and matches right-side rows on the key, leaving nulls for misses. Inner join drops unmatched orders, right join preserves the wrong side, and union performs no key-based matching at all.

Exam trap

The trap here is confusing which side an outer join preserves, so a right join looks similar but actually retains the lookup rows and discards unmatched orders.

140
MCQmedium

A data engineer runs the following statement in a Databricks notebook: CREATE OR REPLACE TEMP VIEW high_value_customers AS SELECT customer_id, SUM(amount) AS total FROM sales GROUP BY customer_id HAVING SUM(amount) > 10000. Later, the same engineer opens a new notebook attached to the same cluster and tries to run SELECT * FROM high_value_customers. What will happen?

A.The query fails with an analysis error only if the underlying sales table was dropped; otherwise the temporary view remains globally accessible.
B.The query succeeds only if the second notebook calls REFRESH TABLE high_value_customers before selecting from it, because temp views are lazily registered.
C.The query fails with TABLE_OR_VIEW_NOT_FOUND because the temporary view is scoped to the SparkSession that created it and is not visible to a different notebook session.
D.The query succeeds because temporary views are stored in the cluster-wide Spark catalog and shared across all notebooks on that cluster.
AnswerC

Temporary views live in the session-scoped catalog of the SparkSession that created them. A second notebook attached to the same cluster receives a separate SparkSession, so the name high_value_customers cannot be resolved and Spark raises TABLE_OR_VIEW_NOT_FOUND. Persisting the result as a managed or external table would be required for cross-notebook visibility.

Why this answer

Temporary views are registered in a session-scoped catalog, so a view created in one notebook is invisible to other notebooks even when they run on the same cluster. A second notebook gets its own SparkSession and cannot resolve the view name, producing TABLE_OR_VIEW_NOT_FOUND. To share results across sessions, the engineer must persist them as a table in the metastore or use a global temporary view in the global_temp database.

Exam trap

The trap here is assuming that two notebooks attached to the same cluster share one SparkSession and therefore share temporary views, when in fact each notebook session has its own session-scoped catalog.

141
Multi-Selectmedium

A data engineer is writing a Spark SQL query that joins a `transactions` table to a `customers` table on `customer_id`. The engineer wants to ensure that rows from `transactions` with no matching customer are still returned, with nulls for customer columns, and also wants to exclude duplicate rows that arise from the join. Which TWO clauses should the engineer include? (Choose two.)

Select 2 answers
A.Use a `LEFT OUTER JOIN` between `transactions` and `customers`.
B.Apply `SELECT DISTINCT` to the final result set.
C.Use an `INNER JOIN` instead of an outer join.
D.Use a `FULL OUTER JOIN` between the two tables.
E.Add a `GROUP BY` on all columns from both tables.
AnswersA, B

A `LEFT OUTER JOIN` preserves every row from the left `transactions` table and fills unmatched customer columns with nulls. This satisfies the requirement to retain transactions lacking a matching customer. Inner joins would drop those rows, so the outer join is essential to the scenario's stated goal of not losing left-side records.

Why this answer

Preserving unmatched left-side rows requires an outer join oriented to the left table, and eliminating duplicated combinations requires a distinct projection. Together, a left outer join plus `SELECT DISTINCT` returns every transaction, matches customers where possible, and collapses repeated pairs into single rows. Inner and full outer joins change which unmatched rows survive and do not meet the stated retention rule.

Exam trap

The trap here is treating deduplication as something a join type provides, when duplicate elimination requires a separate distinct or aggregate operation after the join is performed.

142
Multi-Selectmedium

A data engineer is using Spark Connect from a remote Python client to interact with a Databricks cluster. The engineer wants to understand which operations are executed on the server side versus the client side. Which two statements correctly describe this behavior? (Choose two.)

Select 2 answers
A.Calls such as df.filter() and df.select() build a logical plan on the client and are sent to the server for execution.
B.The client maintains a local SparkContext that runs tasks in parallel with the remote server to reduce latency.
C.User-defined functions defined with @udf are executed in the client Python process to avoid shipping code to the server.
D.df.collect() brings the result set to the client, so the returned data is materialized in the client process memory.
E.df.show() executes entirely on the client by sampling data from a local cache maintained by Spark Connect.
AnswersA, D

In Spark Connect, DataFrame transformations like filter and select are lazy and only construct an unresolved logical plan on the client. That plan is serialized and sent to the server, where it is analyzed, optimized, and executed. The client does not process data locally for these operations, so the heavy lifting occurs on the remote Spark server.

Why this answer

The correct statements highlight the division of labor in Spark Connect: transformations build a logical plan on the client, while actions like collect trigger server-side execution and return results to the client. UDFs are shipped to the server, and there is no local SparkContext or client-side data cache. Understanding this split is essential for predicting where code runs and where memory is consumed.

Exam trap

The trap here is assuming that client-side Python code, such as a UDF body, executes locally, when Spark Connect actually ships it to the server for execution.

143
MCQhard

A Spark job is running on Databricks and experiences a stage where tasks are taking much longer than expected. The Spark UI shows that some tasks have significantly higher shuffle read sizes than others, and the stage is skewed. Which Spark feature can automatically mitigate this skew by splitting large partitions into smaller ones?

A.Columnar Shuffle with Kryo Serialization
B.Dynamic Resource Allocation
C.Adaptive Query Execution (AQE) with skew join optimization
D.Speculative Execution
AnswerC

AQE dynamically reoptimizes query plans based on runtime statistics. When it detects skewed partitions during a shuffle, it can split large partitions into smaller sub-partitions, balancing the load across tasks. This reduces the impact of skew and improves stage performance, making it the correct feature for this scenario.

Why this answer

Adaptive Query Execution (AQE) in Spark 3.x can dynamically detect and handle skew during shuffle operations. When enabled, it can split large skewed partitions into smaller ones, distributing the load more evenly across tasks. This directly addresses the skew observed in the stage, reducing task duration and improving overall job performance.

Exam trap

The trap here is confusing skew mitigation with other performance features like Dynamic Resource Allocation or Speculative Execution, which do not split partitions to balance load.

144
MCQeasy

Which clause is used in a Spark SQL query to limit the number of rows returned by a query, and in which logical order is it executed?

A.LIMIT, executed after ORDER BY.
B.TAKE, executed before ORDER BY.
C.FETCH, executed before WHERE.
D.TOP, executed after GROUP BY.
AnswerA

In SQL, LIMIT is applied after the final result set has been determined, including any sorting required by the ORDER BY clause. This ensures that if you request the top 10 rows, you receive the 10 rows with the highest or lowest values as determined by the specified column ordering.

Why this answer

The LIMIT clause is a common SQL operation used to restrict result sets. In Spark SQL, it is essential to understand that it is applied after sorting if an ORDER BY is present, or simply on the result stream if not. Using LIMIT is a critical performance practice to avoid overwhelming the driver when previewing large datasets in a notebook environment.

Exam trap

Candidates often assume LIMIT is applied before sorting, which would result in non-deterministic data. They fail to realize the logical execution order is crucial for consistent results.

145
MCQeasy

A Databricks job fails with an OutOfMemoryError on the driver. The job collects a large DataFrame to the driver for local processing. Which action should you take to resolve this?

A.Increase the executor memory.
B.Increase spark.driver.maxResultSize.
C.Set spark.sql.shuffle.partitions to a higher value.
D.Replace collect() with take(n) or write the DataFrame to storage.
AnswerD

collect() transfers the entire DataFrame to the driver, which can cause an OutOfMemoryError if the data is large. Using take(n) retrieves only a limited number of rows, or writing the DataFrame to storage avoids bringing all data to the driver. This directly addresses the driver memory issue by reducing the amount of data transferred. It is the correct approach.

Why this answer

The OutOfMemoryError on the driver is caused by collecting a large DataFrame. The correct fix is to avoid transferring all data to the driver by using take() for a sample or writing the results to storage. Increasing driver or executor memory or adjusting shuffle partitions does not address the fundamental problem of excessive data transfer to the driver.

Exam trap

The trap here is increasing driver memory or maxResultSize instead of reducing the amount of data collected to the driver.

146
MCQhard

When a Spark job is stuck in a shuffle phase, what is the most effective first step to identify the root cause of the performance bottleneck?

A.Restart the cluster to clear any cached data.
B.Check the Spark UI 'Stages' tab for uneven task durations.
C.Increase the spark.executor.memory value.
D.Enable speculative execution immediately.
AnswerB

The Spark UI allows you to see the duration and data size of every task in a stage. If one task takes significantly longer than others, it is a clear indicator of data skew. This allows the developer to isolate the specific partition causing the delay and implement a remediation strategy.

Why this answer

The Spark UI provides a detailed breakdown of stages, tasks, and memory usage. Examining the 'Stages' and 'SQL' tabs reveals which tasks are taking the longest, which is indicative of skew or resource contention. Learning to interpret the Spark UI is the single most important skill for a Databricks developer, as it turns opaque 'stuck' jobs into actionable data regarding partition distribution and executor behavior.

Exam trap

Candidates immediately restart the cluster or rewrite business logic instead of inspecting task metrics in the Spark UI to isolate data skew.

147
MCQhard

A streaming DataFrame `df` has a watermark defined on `eventTime` with a delay of 5 minutes. The query uses `withWatermark("eventTime", "5 minutes")` and writes to a Delta table in `append` output mode. A record with event time 10:00 arrives when the watermark is at 10:10. What happens to this record?

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

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

Why this answer

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

This matches the expected watermark behavior.

Exam trap

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

148
MCQhard

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

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

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

Why this answer

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

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

Exam trap

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

149
MCQhard

Refer to the exhibit. Based on the error log, what is the most likely cause of the job failure?

A.The Spark Driver failed to schedule the tasks properly.
B.The executor ran out of memory while processing a task.
C.There was a network failure between the Driver and the worker.
D.The cluster manager was unable to find available nodes.
AnswerB

The error 'Container killed by YARN for exceeding memory limits' directly points to the executor consuming more memory than allocated. This often happens during heavy data manipulation, such as large joins or aggregations, where the memory demand exceeds the limits set for the JVM heap or off-heap memory.

Why this answer

The error explicitly states that the container was killed by YARN for exceeding memory limits. This is a common issue in Spark when executors are assigned tasks that require more memory than what is available in the configured heap. Understanding this error is crucial because it indicates a need to either increase the `spark.executor.memory` setting or optimize the code to reduce the memory footprint of individual tasks during data processing.

Exam trap

Students often assume network timeouts or disk failures cause container terminations, missing explicit YARN memory limit violations indicated in executor error logs.

150
Multi-Selecthard

A developer is using Pandas API on Spark in a Databricks notebook and needs to combine two Pandas-on-Spark DataFrames that originate from different Spark DataFrame ancestors. They encounter a `compute.ops_on_diff_frames` error. Which two actions will resolve this error? (Choose two.)

Select 2 answers
A.Call `.cache()` on both DataFrames before the operation to align their internal anchors.
B.Enable the Spark configuration `spark.databricks.pandas.enableArrow` to allow cross-frame operations.
C.Set `spark.sql.execution.arrow.pyspark.enabled` to true and retry the operation.
D.Convert one of the DataFrames to a pandas DataFrame using `.to_pandas()` and then create a new Pandas-on-Spark DataFrame from it.
E.Use `psdf1.spark.frame()` and `psdf2.spark.frame()` to obtain the underlying Spark DataFrames, then join them using PySpark DataFrame APIs.
AnswersD, E

By collecting one frame to the driver with `.to_pandas()` and then reconstructing a Pandas-on-Spark DataFrame via `spark.createDataFrame` or `ps.from_pandas`, the new object gets a fresh Spark lineage that does not conflict with the other frame. This breaks the diff-frames constraint, allowing the combination to proceed, though it requires the collected data to fit in driver memory.

Why this answer

The `compute.ops_on_diff_frames` error occurs when Pandas-on-Spark objects derive from different Spark DataFrame ancestors. Two valid resolutions are to break the lineage conflict: either collect one frame to pandas and recreate a Pandas-on-Spark object, or drop to the underlying Spark DataFrames and combine them with PySpark operations. Both approaches produce a single consistent lineage, allowing the operation to complete.

Exam trap

The trap here is thinking that a performance or serialization setting such as Arrow can resolve a logical lineage conflict, when the error is fundamentally about differing Spark DataFrame ancestors.

Page 1

Page 2 of 4

Page 3

All pages