In Spark, does writing SQL instead of DataFrame code change how fast a query runs?
answer
- two front doors, one hallway
- both forms build the same tree
- the analyzer never learns which you typed
- compare the optimized logical plans
- the exceptions are RDDs and Python UDFs
basics
~20 sNo. A SQL string and the equivalent DataFrame code build the same logical plan, and Catalyst analyzes, optimizes and compiles both identically. Pick whichever reads better. The real exceptions are the RDD API and Python UDFs, which Catalyst cannot optimize through.
solid answer
~40 sBoth front doors lead to the same engine. `spark.sql("SELECT ...")` is parsed into an unresolved logical plan; DataFrame calls such as `df.filter(...).groupBy(...)` build that same plan directly without parsing. From that point on there is one path: the Analyzer resolves names against the catalog, the Catalyst optimizer rewrites the logical plan, the planner picks physical operators, and Tungsten generates code. You can prove it by comparing `df.explain(True)` for the two forms — the optimized logical plans are identical. Choose by readability and testability, not performance: SQL is easier for analysts to review, the DataFrame API composes better in code. Two things genuinely do change performance: dropping to the RDD API, which bypasses Catalyst entirely, and Python UDFs, which Catalyst treats as an opaque box in either dialect.
code
python · 15 lines# same query, two dialects
sql_form = spark.sql("""
SELECT dept, sum(salary) AS total
FROM emp
WHERE salary > 100000
GROUP BY dept
""")
api_form = (spark.table("emp")
.filter("salary > 100000")
.groupBy("dept")
.sum("salary"))
sql_form.explain(True)
api_form.explain(True)go deeper
Know that a SQL string and DataFrame code produce the same plan and the same performance, and be able to run df.explain() to show it. Pick whichever form is clearer for the task.
Explain the pipeline both dialects share — parse or build, analyze, optimize, plan, generate code — and name the two things that genuinely break it: the RDD API and Python UDFs.
Be ready to settle this argument on a real team with plan diffs rather than benchmarks, and to spot the mixed-in opaque expression that is the actual cause when one formulation really is slower.
Own the convention: which dialect the platform standardizes on, how queries are reviewed and tested, and a rule that opaque UDFs need justification because they cost the optimizer its visibility.
## Two front doors, one engine Spark SQL accepts a query in several shapes — a SQL string through `spark.sql(...)`, a DataFrame chain in Python or Scala, a typed `Dataset` in Scala — and the very first thing it does with any of them is turn them into the *same* internal data structure: an unresolved logical plan, a tree of operators such as `Project`, `Filter`, `Aggregate` and `UnresolvedRelation`. A SQL string reaches that tree through a parser (Spark uses an ANTLR grammar). DataFrame code skips the parser: each API call appends a node to the tree directly. That is the *only* difference between the two, and it costs microseconds of planning time on a job that runs for minutes. From the unresolved plan onward there is a single pipeline, called Catalyst: 1. **Analysis** — the Analyzer looks up table and column names in the catalog, attaches types, resolves function calls, and applies type coercion. An unknown column raises `AnalysisException` here, before a single task is scheduled. 2. **Logical optimization** — rule batches rewrite the tree: pushing filters closer to the scan, pruning unused columns, folding constants, simplifying boolean expressions. 3. **Physical planning** — strategies turn logical operators into executable ones and choose join implementations; preparation rules insert `Exchange` (shuffle) nodes where an operator needs its input redistributed. 4. **Code generation** — Tungsten's whole-stage code generation fuses adjacent operators into a single compiled JVM loop. Because steps 1–4 have no idea which dialect you typed, identical queries perform identically. ## Proving it The reliable check is the plan, not a stopwatch: ```python a = spark.sql("SELECT dept, sum(salary) FROM emp WHERE salary > 100000 GROUP BY dept") b = (spark.table("emp") .filter("salary > 100000") .groupBy("dept").sum("salary")) a.explain(True) b.explain(True) ``` Compare the `== Optimized Logical Plan ==` section of each. If the two are the same tree, the runtime work is the same too. This is a good habit generally: when someone claims one formulation is faster, ask for the two plans. ## PySpark is not slower — for DataFrame operations A persistent myth holds that PySpark is inherently slower than Scala Spark. For DataFrame and SQL operations it is not: the Python API is a thin wrapper that builds a plan on the JVM side through a gateway (or, with Spark Connect, sends an unresolved plan to the server). No row of data crosses into Python at all. Every row is read, filtered, joined and aggregated inside the JVM by generated code. The myth is true for exactly one thing, and it is worth knowing precisely why. ## The two real exceptions **Python UDFs.** A function you wrap with `udf(...)` is opaque to Catalyst — it cannot be pushed into the data source, folded, or compiled into the fused loop. At runtime the plan grows a Python-evaluation node, and rows are serialized out of the JVM into a Python worker process and back. That is a genuine, sometimes order-of-magnitude cost, and it exists whether you invoke the UDF from a SQL string or from DataFrame code. The fix is to express the logic with built-in functions from `pyspark.sql.functions`, or to use a vectorized `pandas_udf` when a Python implementation is unavoidable. **The RDD API.** An RDD is a distributed collection of opaque JVM objects with no schema. Calling `df.rdd.map(...)` leaves Spark SQL entirely: no analysis, no predicate pushdown, no column pruning, no code generation. Spark still schedules stages and tasks, but every optimization described above is gone, and converting back with `toDF()` does not recover what was lost. In Scala, a typed `Dataset` lambda such as `.filter(p => p.age > 18)` has the same problem in miniature — Catalyst cannot see inside the closure, so it cannot push that predicate down, and rows must be deserialized into JVM objects to run it. Written as a column expression, `.filter($"age" > 18)`, it is transparent and pushes down. ## What to actually choose on Since performance is a tie, decide on other grounds. SQL is portable, reviewable by analysts, and natural for long joins and window expressions. The DataFrame API composes: you can build a query from functions, parameterize it, and unit-test the pieces. Mixing is fine and common — register a temporary view and query it in SQL, then keep chaining in code. What matters for speed is the shape of the query and whether every expression in it is something Catalyst can see through.
- Where does the RDD API sit in this picture?Outside it. An RDD holds opaque JVM objects with no schema, so there is nothing for the Analyzer to resolve and nothing for the optimizer to rewrite — no predicate pushdown, no column pruning, no whole-stage code generation. Spark still builds a DAG and runs tasks, but every Catalyst optimization is forfeited, and calling toDF() afterwards does not recover it.
- Is a Scala Dataset with typed lambdas as optimizable as a DataFrame?Only partly. Column expressions are transparent to Catalyst, but a typed lambda is a compiled closure the optimizer cannot inspect, so a predicate inside it cannot be pushed into the scan and rows must be deserialized into JVM objects to evaluate it. Expressing the same filter as a column expression keeps it optimizable while losing compile-time type safety.
- How would you prove to a colleague that two formulations are equivalent?Run explain(True) on both and diff the optimized logical plan section. Identical trees mean identical runtime work. If they differ, the difference is visible right there — a missing PushedFilters entry, an extra Exchange, a different join operator — which is far more informative than timing two runs on a shared cluster.
Two people order the same dish, one by pointing at the menu and one by naming it. The order ticket that reaches the kitchen is identical, so the food arrives at the same speed.
saying these in an interview costs you the question
- Says SQL strings are interpreted while DataFrame code is compiled
- Claims the DataFrame API bypasses the optimizer and runs directly
- Believes PySpark DataFrame operations are slower than Scala ones
- Thinks a Python UDF is optimized like a built-in function
- Assumes df.rdd.map still gets predicate pushdown