PySpark Interview Questions And Answers That Actually Come Up

I've sat through too many PySpark interviews where the candidate could recite the documentation but fell apart when asked about broadcast variables or partition skew. This guide skips the textbook stuff and goes straight to what matters in practice. Let's start with something most people get wrong. PySpark lazy evaluation. Everyone says "it computes only when needed" like they read it once. Here's what actually happens: when you define a transformation like map or filter, nothing runs on the cluster. The computation graph just sits there in memory. It only executes when an action like collect, count, or save triggers it. This means you can chain twenty transformations and PySpark will optimize the whole thing into one job at the end. The practical implication is that debugging is annoying as hell. If your transformation has a bug, you won't know until the action runs. I once had a UDF that silently returned null for every row because of an unhandled exception inside it. The DataFrame still built fine. It was the collect call that took down the executor. My workaround was wrapping every UDF in a try-except that logs the error and returns a default value instead of crashing the whole stage. Use spark.udf.register() with explicit return types so errors surface earlier.

Partitioning is where most interviews trip people up. The default is 200 partitions, which is almost always wrong for real data. If you have a 50GB CSV and 200 partitions, each partition is 250MB. That's fine for a clean join. But if you're doing a groupBy on a high-cardinality column, you'll get massive skew. One partition might have 10% of your data while others barely have anything. The fix is repartition or coalesce depending on what you need. Repartition does a full shuffle, which is expensive but distributes evenly. Coalesce just merges existing partitions without a shuffle, so it's faster but only reduces partition count. Use coalesce when shrinking and repartition when you need to redistribute data evenly. I remember running a word count on a 200GB text corpus with the default partitions. The job took 47 minutes. After repartitioning to 800, it dropped to 11. The difference wasn't magic, it was just the Spark scheduler having enough tasks to keep all the executors busy instead of waiting on stragglers.

Core Concepts You Need to Explain Clearly

Driver versus executor. The driver is the process that runs your main() function, creates the SparkContext, and builds the DAG. It talks to the cluster manager, schedules tasks, and collects results. Executors are the worker processes that actually run the tasks and store data in memory or disk. Everything you write in PySpark gets serialized and sent to executors. If your objects aren't serializable, you get a pickling error. Pandas DataFrames aren't serializable across the wire by default, which is why groupby_map_partitions exists. Catalyst optimizer. This is Spark's query optimization engine. When you write a DataFrame query, Catalyst rewrites it before execution. It pushes filters down, combines predicates, reorders joins based on statistics, and caches intermediate results when beneficial. You can see what it does by calling df.explain(). I always tell candidates to run explain() on complex queries before the interview. Understanding why Catalyst chose a certain join strategy shows you actually use the tool instead of just reading about it. Join strategies in Spark. There are three main types: shuffle hash join, broadcast hash join, and sort merge join. Shuffle hash join does a shuffle on both sides, which is the most expensive. Broadcast hash join sends one small table to all executors via the network, then probes with the larger table. Sort merge join sorts both sides and merges them, good for large sorted datasets.

Get the Full Details

1. kpmg pyspark interview questions and answers | pyspark interview questions and answers - YouTube
1. kpmg pyspark interview questions and answers | pyspark interview questions and answers - YouTube

The counter-intuitive part: broadcast isn't always faster even when the small table fits in memory. If the broadcast variable takes significant time to serialize and ship, and your cluster has tons of cores, a shuffle hash join might win because it processes data in parallel from the start. I learned this the hard way with a lookup table of about 500MB. Broadcasting it to 200 executors took longer than the join itself. Switching to a regular join cut the runtime by half. The broadcast threshold default is 10MB, so you need to adjust spark.sql.autoBroadcastJoinThreshold when dealing with medium-sized tables.

Questions About Performance and Tuning

Memory management. Spark uses a single memory pool divided into storage and execution. Storage is for cached data and broadcast variables. Execution is for shuffles and sorts. If your cache evictions are high, you're using too much memory for storage. If you see speculative execution killing tasks, executors are running out of memory during shuffle. spark.sql.shuffle.partitions controls the number of partitions after a shuffle. Default is 200. For most production workloads with more than a few GB of data, you need to increase this. Setting it to 400 or 800 is typical. But don't just set it blindly high. More partitions means more task overhead. Each task has startup cost, and with hundreds of executors, the overhead compounds. Profile your job first, then tune. Serialization format. PySpark uses Java serialization by default, which is slow and produces large payloads. Switching to Kryo serialization through spark.serializer to org.apache.spark.serializer.KryoSerializer usually cuts shuffle data size by 3 to 5 times and speeds up serialization significantly. The catch is you have to register your custom classes with Kryo, or it falls back to Java serialization anyway. Add spark.kryo.registrator pointing to your registration class and you're set.

Accumulators and broadcast variables. Accumulators are write-only variables that only the driver can read. They're useful for counting errors or tracking metric totals across tasks. Broadcast variables send read-only data to all executors once instead of sending it with every task. Don't broadcast everything. If your broadcast variable changes per action, you're just adding overhead. Broadcast data that stays constant across multiple operations, like reference tables or configuration dictionaries.

pyspark interview questions and answers | regular expression in pyspark| regexp_replace | # ...
pyspark interview questions and answers | regular expression in pyspark| regexp_replace | # ...

SQL vs DataFrame API vs RDDs

This comes up constantly. The DataFrame API and Spark SQL share the same Catalyst optimizer underneath. There's no performance difference between writing a DataFrame chain and the equivalent SQL query. Use whatever is more readable for your team. I prefer DataFrame for transformations that require complex logic because the code chains are easier to follow. I use SQL for aggregations and window functions because the syntax is cleaner. RDDs are the lowest level abstraction. They give you full control but bypass Catalyst entirely, which means no optimization. The only time I reach for RDDs is when working with unstructured data that doesn't fit a schema, or when I need to do operations that DataFrame doesn't support directly. Even then, I convert to DataFrame as soon as possible to get the optimizer benefits back. One thing people miss: pandas UDFs. Standard PySpark UDFs serialize data row by row between Python and JVM, which is incredibly slow. Pandas UDFs use Apache Arrow to transfer entire columns at once in bulk. The performance difference is massive, often 10 to 100 times faster. The tradeoff is you lose row-level granularity inside the function since you're operating on pandas Series. But for most mathematical or string operations, this is a non-issue. Register them with @pandas_udf instead of @udf.

Common Pitfalls That Show Up in Interviews

Data skew is the number one production problem. When one key in your dataset has disproportionately more rows than others, you get straggler tasks that dominate your job duration. The standard workaround is salting. You add a random number to the skewed key on both sides of the join, which spreads the hot key across multiple partitions. Then you aggregate or reduce to combine the salted results. It adds complexity but it's the only reliable fix for severe skew. Small file problem. If your source data has thousands of tiny files, Spark creates thousands of tasks, and the overhead dwarfs the actual computation. Use coalesce or repartition early in your pipeline to consolidate files. Writing with .repartition(1) into a single file per partition also helps downstream jobs. Avoid collect() on large datasets. This is the most common rookie mistake. collect() brings all data to the driver, which means it all loads into driver memory. If your dataset is larger than available memory, the job crashes. Use take(n) if you just need a sample, or write to a file if you need the data elsewhere. I've seen candidates confidently explain lazy evaluation and then recommend collect() for debugging a 10GB dataset. That contradiction gets them filtered out immediately.

Checkpointing. Long dependency chains cause Spark to recompute from the beginning if a task fails partway through. Checkpointing writes intermediate results to durable storage, breaking the lineage chain. Use it for iterative algorithms or very long pipelines where recovery time matters more than the write overhead. I use it in ETL jobs that run for hours and have 50+ transformations. The checkpoint write adds maybe 15% overhead but saves 20 minutes on recovery when something fails.

3. pyspark interview questions and answers for experienced | databricks interview question ...
3. pyspark interview questions and answers for experienced | databricks interview question ...

What Good Answers Look Like

When asked about optimizing a slow Spark job, don't just list tuning parameters. Walk through your diagnostic process. Check the Spark UI first. Look for skew in task durations, examine the DAG for wide dependencies, check if there are excessive small tasks, and review GC logs if executors are slow. Then propose targeted fixes based on what you found. Generic answers about increasing partitions sound good until someone asks how you determined the right number. For questions about handling missing data, mention thatfillna() and dropna() are actions that trigger computation, so they're not free. For large datasets, consider filtering with where() before calling those methods to reduce the data volume first. Also note that fillna works differently on numeric and string columns, and with UDTs you might need to cast first. Window functions are another frequent topic. The difference between row_number(), rank(), and dense_rank() is standard knowledge. What separates good candidates is understanding the memory implications. Window functions require sorting data within partitions, which means shuffle and spill to disk if the partition is too large. If your window frame is unbounded, the entire partition stays in memory. For large windows, add a limit to the frame or repartition more aggressively before applying the window.

Understanding the difference between persist() and cache() matters too. They're the same method under the hood, but cache() defaults to MEMORY_ONLY while persist() lets you choose the storage level. Use MEMORY_AND_DISK if your data doesn't fit in RAM. Use DISK_ONLY for data you only access occasionally. MEMORY_ONLY_SER is faster on read but uses less space, at the cost of deserialization time.

Realistic Edge Cases Worth Mentioning

Schema evolution. When reading from Avro or Parquet files that change schema over time, PySpark handles it differently depending on the format. Parquet schema evolution is well-supported. Avro can be trickier with missing fields depending on your reader configuration. In production, I enforce schemas explicitly with .schema() when reading external data rather than letting Spark infer it. Inference is convenient but slow and can produce wrong types if the sample isn't representative. Handling timezone issues. Spark's default behavior for timestamps depends on your locale configuration. Inconsistent timezone handling between source data and Spark has caused me to waste an entire day tracking down off-by-one-hour errors in a daily aggregation job. Always set spark.sql.session.timeZone explicitly and make sure your source data timestamps are in the same timezone or properly converted before processing. Out of memory on the driver is more common than people expect. It usually happens when you collect large results, create enormous broadcast variables, or accumulate too many accumulators. If your driver is getting OOM, increase spark.driver.memory and investigate which operation is loading everything locally. The Spark UI's storage tab will show you exactly how much data is cached on the driver.

day 12 : pyspark scenario based interview questions and answers | join in pyspark | #interview # ...
day 12 : pyspark scenario based interview questions and answers | join in pyspark | #interview # ...

Dynamic partitioning. When writing Hive-style partitions, using dynamic partitioning avoids creating directories you don't need. Set spark.sql.sources.partitionOverwriteMode to dynamic, but be aware this can overwrite all partitions for a given partition key, not just the ones you inserted. In production pipelines where multiple jobs write to the same partitioned table, this causes data loss. Use static partitioning in shared environments or coordinate write schedules carefully. There's no single right answer for most of these questions. What interviewers are looking for is whether you've dealt with the problems, understand the tradeoffs, and can reason through alternatives. Saying "it depends" followed by a concrete explanation of when each approach is appropriate beats a memorized textbook answer every time.