
pandas loads a table into memory and gets to work, and for most tables that is the right tool. The trouble starts when the table is larger than the memory, or the job is slower than one CPU can finish overnight. Spark is the engine for that case: it spreads the data over many machines, runs the same operations on each piece, and hides the fact that a machine may die halfway through. PySpark is its Python face. This post says what the engine adds over pandas, the four mechanisms that make a distributed computation fast and recoverable, and where the price is paid, as of early 2025.
The DataFrame looks like pandas and runs somewhere else
A PySpark session talks to a cluster, which may be one laptop for development. The DataFrame API is close enough to pandas to read at a glance, and the difference is where the rows live: on the cluster’s workers, in partitions, not in the process that wrote the code. The three-row table here is synthetic and only exists to make the mechanics visible; the same lines run unchanged on a table of a billion rows.
from pyspark.sql import SparkSession
from pyspark.sql.functions import avg, col
spark = SparkSession.builder.appName("example").getOrCreate()
df = spark.createDataFrame([("Alice", 30), ("Bob", 25), ("Charlie", 35)], ["name", "age"])
older = df.filter(col("age") > 28)
older.show()filter does not run when the line executes. It records a step, and nothing is computed until an action such as show, count or collect asks for a result. That laziness is the first mechanism and the others depend on it.
Four mechanisms, one idea: plan the whole job before running any of it
- Lazy evaluation and the DAG. The recorded steps form a directed acyclic graph. When an action fires, Spark’s optimiser looks at the whole graph, pushes filters down to the data source, prunes unused columns, and fuses steps that can run in one pass over the data. MapReduce ran one step at a time and wrote to disk in between; the plan is what makes Spark faster than that.
- Partitions and parallelism. A DataFrame is split into partitions and each task processes one, on whichever worker holds it. Ten workers, a hundred partitions, and a filter runs a hundred times in parallel with no code to say so.
- In-memory computation. Intermediate results stay in RAM across the steps of a job, and a DataFrame that will be reused can be cached, so an iterative algorithm does not reread its input from disk each round.
- Lineage for fault tolerance. Spark does not replicate intermediate results. It remembers how each partition was derived, and if a worker dies it recomputes only the partitions that were lost, from the plan. A failed task is retried; a failed machine costs a recompute, not a restart.
The lineage idea is the one that gives the old core structure its name, the resilient distributed dataset, and it is still underneath the DataFrame API. Writing against RDDs directly is rare now; the DataFrame plan lets the optimiser do more.
SQL and DataFrames are the same plan
A DataFrame can be registered as a view and queried in SQL, and the two produce the same plan, so the choice is taste and audience. The end-to-end shape of a job is: read, transform, aggregate, and hand the result on.
df = spark.read.csv("employees.csv", header=True, inferSchema=True)
by_dept = df.filter(col("age") > 30).groupBy("department").agg(avg("salary").alias("avg_salary"))
df.createOrReplaceTempView("employees")
same = spark.sql("""
SELECT department, AVG(salary) AS avg_salary
FROM employees WHERE age > 30 GROUP BY department
""")
by_dept.show()
same.show()Both show calls print the same table, because both are the same DAG. Spark’s machine-learning library and its streaming API sit on the same foundation, which is why one platform covers ETL, model training and stream processing without switching engines.
Where it stops holding
Spark’s overhead is real. Starting a session, planning a job and shipping tasks to workers takes seconds, so on a table that fits in memory pandas finishes before Spark has started, and the rule is to use Spark when the data or the compute does not fit on one machine, not before. The Python side has its own cost: a lambda in a map or a Python UDF runs in a separate Python process on each worker with data serialised across, and is many times slower than the built-in column functions, so the fast path is to stay in the DataFrame API. And a cluster is something someone runs; the managed platforms, Databricks chief among them, exist because that is a job in itself.
Distribute. Data. Plan. First. Cache. Recompute. Losses. Pandas. Until. It. Doesn’t. Fit.
References
- PySpark documentation
- Spark SQL guide
- Zaharia, M. et al. (2012). Resilient Distributed Datasets: a fault-tolerant abstraction for in-memory cluster computing. NSDI 2012.
- Databricks on this blog