Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. How PySpark talks to the JVM

PySpark · Execution Model

How PySpark talks to the JVM

Mediumpyspark-53
py4jpython-workerudfarrow

Question

How does Python code actually run in PySpark, and where does the cost come from?

Solution

PySpark is a thin Python layer on top of the JVM. Your Python code mostly sends instructions to the JVM, and the JVM does the work. The cost shows up only when Python code has to process the actual rows.

Two very different paths

Path one: DataFrame operations.

df.filter(F.col("amount") > 0).groupBy("region").sum("amount")

Your Python process talks to the JVM driver through Py4J, a bridge that lets Python call JVM objects. Python builds a logical plan. Catalyst optimizes it, and executors run it in the JVM. No row is ever handed to Python. The speed is about the same as Scala.

Path two: Python code that touches rows. This means Python UDFs, rdd.map(lambda ...), and foreachPartition.

JVM executor --serialize rows--> Python worker --run your function--> serialize result --> JVM

For each task Spark starts (or reuses) a Python worker process next to the executor and streams rows to it. Every row is serialized, sent, deserialized, processed by interpreted Python, serialized again and sent back. That is slow compared to the JVM, and it breaks Catalyst optimizations around the UDF.

Reducing the cost

  • Use built-in functions wherever possible, so rows never leave the JVM.
  • If you need Python logic, use a Pandas UDF. It transfers data in columnar batches through Apache Arrow, and you work on whole pandas Series, so the per-row overhead drops a lot.

Memory consequence

Python workers run outside the JVM heap, so their memory is not part of executor-memory. It comes from spark.executor.memoryOverhead (and spark.executor.pyspark.memory if set). A job with heavy Python UDFs that has a small overhead gets its container killed by YARN or Kubernetes with an out-of-memory message from the resource manager, not a Java error. Raising overhead is the fix.

🎯 Put this concept into practice

Solidify this answer with real hands-on interview drills in the browser studio.

Open related drill →
PreviousNext