A data engineering team is developing a Spark application to process sensitive financial data. They need to pass a large, read-only lookup table (approximately 500MB) containing currency exchange rates to all executor nodes. This table is used in a join operation within multiple tasks. Which Spark feature should be used to distribute this lookup table efficiently and minimize network I/O?
Answer and explanation
Correct answer: B
A broadcast variable is the correct choice for distributing a large, read-only dataset to all worker nodes. The driver serializes the variable and sends it to each executor only once, where it is cached in memory. Tasks on that executor can then access the data locally without causing repeated network transfers. Accumulators are used for aggregating results back to the driver. Caching the DataFrame would work, but broadcast is specifically designed for this 'side data' distribution pattern and is generally more efficient for joins. A UDF is for custom logic, not data distribution.
Question 2
In the context of the Apache Spark execution hierarchy, which of the following events will always trigger the creation of a new Spark Stage?
Answer and explanation
Correct answer: C
Spark creates a new stage at each shuffle boundary. Wide transformations, such as groupBy(), join(), or repartition(), require data to be redistributed (shuffled) across the network between executors. This redistribution marks the end of one stage and the beginning of another. Narrow transformations can be pipelined within a single stage. Actions trigger the execution of a job, which consists of one or more stages, but the action itself doesn't define the stage boundary; the shuffle does.
Question 3
Multiple answers
A Spark job processing a large dataset is experiencing performance degradation. Analysis of the Spark UI shows that one task in a particular stage is taking significantly longer than all other tasks. The stage involves a groupBy('user_id') operation. What is the most likely cause of this issue and the most appropriate solution? (Select TWO)
Answer and explanation
Correct answers: B, D
Question 4
A developer needs to read a large Parquet dataset partitioned by year, month, and day. To optimize read performance, they only want to load data for the first week of January 2023. Which Spark SQL query correctly applies partition pruning to achieve this?
SELECT * FROM sales_parquet WHERE _______
Answer and explanation
Correct answer: B
For partition pruning to be effective, the filter conditions must be applied directly to the partition columns. This allows Spark's Catalyst optimizer to read the file system metadata and skip reading the data files in partitions that do not match the filter. Applying functions to the partition columns (like concat or to_date) can prevent the optimizer from pushing down the predicate, forcing a full table scan. The correct approach is to filter directly on the raw partition columns.
Question 5
A streaming application needs to calculate a running count of events per user and output the updated count for each user as new data arrives in every micro-batch. Which Structured Streaming output mode is designed for this use case?
Answer and explanation
Correct answer: C
Update mode is specifically designed for use cases involving aggregations. In this mode, only the rows that were updated in the result table since the last trigger will be written to the sink. This is perfect for outputting running counts where only the users with new events in the current micro-batch will have their counts updated in the output. Complete mode would rewrite the entire result table, and Append mode is not supported for aggregations without a watermark.
Question 6
A retail analytics company is building a daily batch processing pipeline using Spark. The pipeline must ingest raw sales transaction data from a CSV file, enrich it with product dimension data from a Parquet file, calculate daily sales aggregates per product category, and write the final report to a Delta table, overwriting the previous day's report.
The raw sales data (sales_df) contains product_id, sale_amount, and transaction_time. The product dimension data (products_df) contains product_id and product_category. The final report must be partitioned by product_category for efficient querying by downstream business intelligence tools.
Which sequence of PySpark DataFrame operations correctly and most efficiently implements this logic?
Answer and explanation
Correct answer: B
This is the correct and most efficient sequence. First, the sales data is joined with the product data to add the product_category. A left outer join is appropriate to ensure no sales are dropped if a product is missing from the dimension table. Second, the enriched data is grouped by the product_category to calculate the sum of sales. Performing the join before aggregation is crucial to have the category available. Finally, the aggregated result is written to a Delta table, correctly using overwrite mode and partitioning by product_category for query performance.
Question 7
True or False: Spark Connect allows a Spark application's driver process to run on a separate machine from the Spark cluster, such as a developer's laptop or an IDE, while the execution of Spark jobs occurs on the remote cluster.
Answer and explanation
Correct answer: A
This statement is true. The primary purpose of Spark Connect is to decouple the client application (where the SparkSession is created and DataFrame logic is defined) from the Spark driver and cluster. It introduces a client-server architecture where the client sends unresolved logical plans to a Spark Connect server running on the cluster, which then translates them into Spark's physical plan for execution on the executors.
Question 8
A data scientist is working with a large Spark DataFrame and needs to apply a complex numerical computation that is already implemented and highly optimized in the scipy library. Which type of User-Defined Function is best suited for applying this scipy function to columns of a Spark DataFrame to maximize performance?
Answer and explanation
Correct answer: C
A Pandas UDF (also known as a vectorized UDF) is the best choice. It processes data in batches as pandas Series or DataFrames, which allows for efficient, vectorized computations using libraries like NumPy or SciPy. This avoids the high overhead of serialization/deserialization and row-by-row processing that occurs with standard Python UDFs, leading to significant performance improvements.
Question 9
Which of the following describes the role of the Driver in a Spark application running in cluster mode?
Answer and explanation
Correct answer: B
In cluster mode, the driver program is launched on a worker node within the cluster. Its main responsibilities are to host the main() method of the application, create the SparkSession, analyze and schedule jobs, and coordinate the execution of tasks on the executors. The executors are the processes that actually run the computation tasks.
Question 10
A developer is writing a DataFrame to a cloud storage location. The requirement is to write the data only if the target location does not already exist. If the location exists, the write operation should fail instead of overwriting or appending data. Which saveMode should be used?
Answer and explanation
Correct answer: D
SaveMode.ErrorIfExists is the correct option. It is the default behavior. If the target location already exists, a AnalysisException is thrown. Append adds data, Overwrite replaces it, and Ignore silently does nothing if the location exists.