AWS Training
Modules Listen Certification
0:00 0:00

← PySpark

PS0 Interview questions

Read these out loud. Under 90 seconds each. The pattern that scores in PySpark interviews is always the same: give the mechanism, not the setting. Anyone can name a config key; the signal is whether you can explain what it changes and why.


Warm-up

Q. What's the difference between a transformation and an action?

Transformations build a plan — select, filter, join, withColumn. They cost essentially nothing when the line runs. Actions execute that plan — write, count, collect, show. The practical consequence is that your traceback points at the action, not at the line with the bug, so you debug backwards from wherever execution was triggered.

Q. Why is collect() dangerous?

It pulls every row into the driver, which is a single process. On anything large that's an out-of-memory death, and the failure is usually ugly — the driver goes unresponsive and the job dies with an error that doesn't obviously point at the cause, because the process that would have reported it is the one that died. My rule is nothing unbounded ever comes back to the driver.

Q. repartition versus coalesce?

repartition does a full shuffle and can increase or decrease the partition count. coalesce only reduces, and avoids the shuffle. The bit people miss is that coalesce avoids the shuffle by narrowing upstream parallelism too — so coalesce(1) after an expensive computation can run that whole computation in a single task. If I need one output file and upstream parallelism, I use repartition(1) and pay for the shuffle.

Q. What is a wide transformation?

One where output partitions depend on many input partitions — groupBy, join, distinct, orderBy. It requires a shuffle, so data crosses the network and hits disk, and it's a stage boundary. Counting the wide transformations in a script tells me roughly how many times the data gets written and re-read.


Depth

Q. Why are Python UDFs slow?

Three reasons. First, every row is serialised out of the JVM into a Python worker process and back. Second, that Python process has its own memory outside the JVM heap, so you can run out of memory somewhere you weren't monitoring. And third — this is the one that usually costs most — the optimiser goes blind. A UDF is opaque to Catalyst, so it can't push the predicate down or prune columns based on it.

That third one is why a UDF inside a filter is especially bad: instead of the scan skipping partitions, Spark reads everything and filters in Python. The extra I/O dwarfs the UDF cost.

The follow-up: "So how would you avoid one?"

In order: use a built-in from pyspark.sql.functions — people reach past those constantly, and things like regexp_extract, from_json and transform cover most of what UDFs get written for. If it genuinely needs Python, use a pandas UDF so the crossing happens in Arrow batches instead of per row. And if it's a lookup, replace it with a broadcast join, which doesn't cross the boundary at all.


Q. What does Apache Arrow do for PySpark?

It's an in-memory columnar format used to transfer data between the JVM and Python processes efficiently. The thing I'd point out is what its existence tells you — that boundary is expensive enough that a dedicated format exists to make it less bad. Arrow narrows the cost. It doesn't remove the boundary, so the hierarchy still holds: built-ins beat pandas UDFs beat plain UDFs.

The follow-up: "Is it on by default?"

In the current documentation, yes for pandas conversion — spark.sql.execution.arrow.pyspark.enabled defaults to true. But I'd check rather than assert, because these defaults have moved between versions, and Arrow optimisation for regular Python UDFs is a separate and more recent setting. On AWS especially, Glue and EMR pin Spark versions behind the current release, so the first thing I do in a new environment is print spark.version and read the docs for that version.


Q. A job has 200 tasks. 199 finish in 20 seconds, one runs for 40 minutes. What's happening?

That's skew — one partition holds far more data than the others. And the important consequence is that scaling the cluster doesn't help, because the extra executors just sit idle watching one task.

To diagnose it I'd look at two things: the distribution of the join or group key, and the actual row count per partition using spark_partition_id. In the Spark UI the equivalent tell is a stage where max task duration is orders of magnitude above the median.

Most often the cause is a dominant key — nulls, zeros, "UNKNOWN", or a genuine whale customer. Nulls are the usual one, because every unmatched row gets the same key.

The follow-up: "How would you fix it?"

Cheapest first. If it's nulls, filter them before the join — they aren't joinable anyway, so there's no reason to drag them through a shuffle. If one side is small enough, broadcast it and there's no shuffle at all. Only if neither works would I salt the hot key, because salting means replicating the other side by the salt factor, which is real cost. And I'd salt only the keys that are actually hot, not everything.

And the one after that: "What if the input is a single 5 GB gzip file?"

Then none of that helps, because gzip isn't splittable — it has to be read by one task no matter what you configure. You'd fix it upstream by writing splittable formats, or by re-encoding it before processing.


Q. Why does calling count() before a write slow the job down?

Because a DataFrame is a plan, not a result. count() executes the whole plan from source, then write executes it again. You've doubled the work to answer a yes/no question. If it's genuinely a "do we have anything" check, take(1) can short-circuit — though even that can be expensive with a very selective filter. Often the best answer is to write unconditionally and let the output be empty.

The follow-up: "When would you cache instead?"

When more than one action consumes the same DataFrame. Cache it, run the actions, then unpersist, because cached data is occupying memory the shuffle wants. What I wouldn't do is cache everything defensively — that evicts the thing that actually needed caching and starves execution memory.


Q. What's the difference between a Spark partition and a partitioned table?

Spark partitions are units of parallelism in memory — one task each. Table partitions are the directory layout in storage, the dt=2026-08-12 folders you get from partitionBy. They interact, because reading a partitioned dataset with a matching filter prunes files and changes how many Spark partitions you start with. The way I keep them straight is: storage partitions decide what you read, Spark partitions decide how you process it.


Design

Q. Design a daily PySpark job processing 500 GB of events into aggregates. Talk me through it.

First I'd ask three things: what's the freshness requirement, is the input partitioned by date, and is it late-arriving. Those decide almost everything.

Assuming daily and partitioned by date: read only the partitions I need so the filter prunes at the scan rather than after it. I'd check the input file layout early — if it's thousands of tiny files or a few unsplittable gzips, that's my parallelism problem before I write any logic.

Structurally I'd split it into pure DataFrame-in, DataFrame-out transformation functions and a thin I/O shell, so the logic is unit-testable with five rows and no cluster.

For the aggregation I'd check key distribution before assuming a plain groupBy is fine — if there's a whale key I'd rather know at design time than at 3am. Dimension joins get broadcast where the dimension fits.

On write, overwrite partitioned by date so a retry replaces rather than doubles — and I'd verify the partition-overwrite semantics for our Spark version, because that setting differs and getting it wrong deletes data. I'd repartition before writing to avoid emitting thousands of small files.

And I'd put assertions in the job — output unique by key, output non-empty — that fail the run. Failing loudly beats writing wrong data, because a failure gets attention in an hour and bad data gets noticed in a quarter by someone who then doesn't trust anything we produce.


Debug

Q. A PySpark job that took 20 minutes now takes 3 hours. Nothing in the code changed. Walk me through it.

Ordered checklist.

One — did input volume or shape change? Not just size: a change in file layout, like an upstream job switching to gzip or emitting many small files, changes parallelism without changing volume.

Two — is it skew or is it uniformly slow? I'd look at the task duration distribution for the slow stage. Max far above median is skew, and that's a completely different fix from "everything is slower".

Three — did key distribution change? A new dominant key, or a spike in nulls from an upstream change, produces exactly this.

Four — read the plan and compare it to what it used to be if I have it. A join that used to be broadcast and is now shuffled, because the "small" side grew past the threshold, is a classic cliff-edge regression.

Five — check the environment actually changed nothing. On Glue or EMR a version bump changes defaults, and I've seen that produce exactly this shape.

I'd do one and two before touching any configuration, because they tell me which class of problem I'm in.

Q. Your UDF's print() statements aren't appearing in CloudWatch.

They're probably appearing, just not where I'm looking. A print from the driver goes to the driver log; a print inside a UDF runs on an executor and goes to executor logs, which is a different log group. So before concluding the UDF isn't running, I'd check the executor logs. If it genuinely isn't running, I'd check whether the action that triggers it is ever called — lazy evaluation means a transformation with no action never executes at all.


The "it depends" question

Q. How many partitions should your job have?

It depends, and here's what I'd want to know: how much data, how many cores are available, and what the job does next.

The reasoning is that you want enough partitions that every core has work and no single task holds more data than an executor can handle — but not so many that scheduling overhead and, on write, small files dominate.

So I'd start from data size and core count rather than a magic number, look at actual task durations and the partition size distribution, and adjust. And I'd treat the write side separately from the compute side — high parallelism during compute and a repartition down before writing is often the right shape, because those two want different numbers.

Anyone who gives you a single number without asking about the data is guessing.


Red flags

❌ Saying this Why it costs you
"UDFs are slow because Python is slow." Misses serialisation and, crucially, the optimiser going blind. Python's speed is the least of it.
"Just increase the number of executors." For skew, that changes nothing. Shows you haven't distinguished distribution from capacity.
"coalesce is always cheaper than repartition." Ignores upstream-parallelism narrowing — the single best discriminator on this topic.
Reciting config keys when asked why something is slow. Names a knob without a mechanism. Give the mechanism first, the key second.
"cache() everything to be safe." Evicts what mattered and starves execution memory.
"A DataFrame holds the data." It holds a plan. This misunderstanding explains a whole class of bugs.
Quoting a UDF speedup multiplier confidently. Those vary hugely by data shape and version. "I'd measure it on our data" is the stronger answer.
Stating a default without saying which version. On AWS the version is pinned behind current. "I'd check spark.version first" scores better than being right by luck.

Practising this

Say each answer to a wall, timed. Then have someone ask only the follow-ups — "so how would you avoid one?", "is it on by default?", "what if the input is a single gzip file?" Those are where the differentiation happens, and they're the ones people haven't rehearsed.