Apache Spark — The answer to a slow job is in the plan and the event log
SQL and DataFrame code become the same plan, and window functions don't reduce rows
In one line
A SQL string and the DataFrame API are two entrances that the same engine turns into the same plan. You choose between them not by which is faster but by which is easier to read. An aggregation groups rows and reduces them, while a window function leaves the rows as they are and looks at neighboring rows, and the result of a window depends on how you set the frame.
Why there are two entrances
A data team mixes people who think in SQL with people who think in Python. Before long, remarks like "writing it in SQL is supposedly faster" or "DataFrames are supposedly better optimized" start to circulate. Both are wrong.
The first paragraph of the Spark SQL guide gives the answer. Spark SQL knows more about the structure of the data and of the computation than RDDs do, and it uses that information to optimize further. And when computing a result, it uses the same execution engine, no matter which API or language you used to express it. The documentation notes that this unification lets you switch back and forth and pick the more natural API for each transformation.
"The same engine" is concrete. A SQL string is parsed into a logical plan, and DataFrame method calls also build a logical plan directly. From there on, both pass through the same analyzer, the same optimizer and the same physical planner. So if you write the same question both ways and print explain(), the physical plans are the same. In this lab you compare them yourself.
How it works
To call a DataFrame from SQL, you need a name. The getting started documentation shows the pattern of registering a DataFrame as a temporary view with createOrReplaceTempView and querying it with spark.sql. The same page notes that a temporary view is session-scoped and disappears when the session that created it ends. A view does not copy data. It is just a name attached to a plan.
orders.createOrReplaceTempView("orders")
by_sql = spark.sql("""
SELECT date_trunc('month', order_ts) AS month, sum(qty) AS units
FROM orders WHERE status = 'paid' GROUP BY 1""")
by_api = (orders.where(F.col("status") == "paid")
.groupBy(F.date_trunc("month", "order_ts").alias("month"))
.agg(F.sum("qty").alias("units")))
The physical plan of an aggregation almost always has the same shape. In the example in the EXPLAIN documentation, the order from the bottom is HashAggregate(... partial_sum ...) → Exchange hashpartitioning(k, 200) → HashAggregate(... sum ...). Each partition first reduces its own share to partial sums, shuffles by key, and then adds the collected partial sums again. What is sent through the shuffle is not the original rows but partial sums reduced to one line per key. That is why, for aggregations such as sums and counts, the shuffle is smaller than you would expect even when the data is large. The 200 in the example is the number of shuffle partitions, which is covered separately in the module on shuffles.
Window functions: an aggregation that does not reduce rows
groupBy leaves one line per category. "The top 3 products per category" cannot be solved that way. You have to leave the product rows as they are and attach, in the next column, the rank within the same category. This is what a window function does. The window function part of the SQL reference separates ranking functions (RANK, DENSE_RANK, ROW_NUMBER and so on) from analytic functions (LAG, LEAD, FIRST_VALUE and so on), and defines the syntax for writing the frame with ROWS or RANGE.
The ranking functions differ on ties. As the rank documentation explains, dense_rank leaves no gaps in the ranking after ties, while rank assigns 5th place next if three rows share 2nd place. If you filter "the top 3" with rank <= 3, ties can produce four or more rows, and if you filter with row_number, which of the tied rows gets in is left to the sort order. Which one is right is decided by the business.
A quieter trap is the default frame. The Window documentation says that without an ordering the default is the whole partition (row-based, from start to end), and with an ordering it is a frame that grows range-based from the start to the current row. According to the rangeBetween documentation, a range boundary is based on the ORDER BY value, not on the position of the row. So if you sort by date and compute a running total, rows with the same date include each other in "up to the current row" and all receive the same running value at once. If you want a running total that grows one row at a time, you have to write a row-based frame with rowsBetween(Window.unboundedPreceding, Window.currentRow).
Exact counts and approximate counts
To count unique customers exactly, you have to gather all customer IDs and remove duplicates, so the values themselves must be shuffled. approx_count_distinct takes an allowed relative standard deviation rsd and estimates, with a default of 0.05. The documentation says that if the rsd must be smaller than 0.01, count_distinct is more efficient instead. The list of built-in aggregate functions states that this estimate is made with HyperLogLog++. Use it where a 5% error is acceptable, such as daily visitors on a dashboard, and not where a single wrong person is unacceptable, such as settlement.
What it looks like in the field
Errors in a SQL string appear at run time. If you get a column name wrong, the Python editor says nothing, and an analysis error occurs the moment you call spark.sql. The DataFrame API is the same if you write columns as strings. Either way, it is caught in the analysis stage before any action, so the cheapest test is to run it once end to end on a small dataset.
Mixing the two styles is normal. For example, you write complex windows and joins in SQL, and work that builds columns repeatedly with a Python loop and the DataFrame API. Since the plans are the same, mixing costs you nothing.
A temporary view is not visible outside the session. If you call the same name from another application or another session, the table is not found. To pass data between applications, you have to write it to a file or a table.
What really matters in practice
- SQL and DataFrames are the same engine and the same plan. Choose by readability, not performance, and if in doubt, compare with explain.
- An aggregation is partial aggregation → shuffle → final aggregation. What gets shuffled is the reduced partial result, not the original rows.
- A window does not reduce rows. Rank the top N with a ranking function and choose the tie rule to suit the business.
- State the frame explicitly for a running total. The default for a window with an ordering is a range frame, so rows with the same value are grouped together at once.
- Use approximate distinct counts only where the error is acceptable. The default rsd is 0.05.
What you will do in the next lab
You register the order data as a temporary view and compute monthly revenue with SQL, rewrite the same result with the DataFrame API, and confirm with explain that the plans on both sides are the same. You pick the top 3 products per category with a window function, then compute a running total with a row-based frame and the change from the previous month with lag. Finally, you count the unique customers with both the approximate function and the exact method, measure the error, and summarize it in a report.