Apache Spark — The answer to a slow job is in the plan and the event log
Monthly revenue, per-category ranking and running totals in SQL and with DataFrames
Goal
You solve the same question with Spark SQL and with the DataFrame API, confirm that the same physical plan comes out, and compute rank, running total and change from the previous month with window functions. Finally, you compare the exact distinct count with the approximate one.
Why it matters
In Spark, SQL and DataFrames are not two engines. Both are turned into the same logical plan, pass through the same optimizer (Catalyst), and run as the same physical plan. That is why "SQL is faster" or "DataFrames are faster" is usually the wrong question. Whatever the team writes in, you can tell right away from the plan whether they are the same or different. After aggregation, the most frequently used tool is the window function. Because it can look at neighboring rows without reducing rows, report numbers such as rank, running total and change from the previous month all come from here. If you set the frame of a window (partition, ordering, range) wrongly, you get wrong numbers without any error, so you must be able to explain the frame in words. Counting distinct values is an operation with a large shuffle. If you accept an error of a few percent, you can count much more cheaply with HyperLogLog++. Approximate for a dashboard, exact for an invoice: which one to use is decided by what the number is used for.
Steps
- In /root/spk/sql/common.py, create
load(spark), which reads orders and products with a schema and registers them as the temporary viewsordersandproducts, and in /root/spk/sql/sql_monthly.py (applicationspk-sql-monthly) write monthly revenue (columnsmonth,revenue) computed with SQL as a CSV with a header row to /root/spk/sql/out/monthly_sql. - In /root/spk/sql/df_monthly.py (application
spk-sql-df), write the same result to /root/spk/sql/out/monthly_df using the DataFrame API withoutspark.sql. - In /root/spk/sql/plans.py (application
spk-sql-plans), save theexplain(mode="formatted")output of the two approaches to /root/spk/sql/out/plan_sql.txt and /root/spk/sql/out/plan_df.txt, and then actually run both DataFrames withcollect(). - In /root/spk/sql/top3.py (application
spk-sql-top3), write the top 3 products by revenue for each category (category) to /root/spk/sql/out/top3 with the columnscategory,product_id,revenue,rank. For a tie, the one with the smallerproduct_idranks higher. - In /root/spk/sql/running.py (application
spk-sql-running), write the daily revenue per channel for January 2026 and the running total within each channel to /root/spk/sql/out/running with the columnschannel,day,revenue,running. - In /root/spk/sql/mom.py (application
spk-sql-mom), write the monthly revenue per channel, the previous month's revenue and the growth rate (percent, rounded to two decimal places) to /root/spk/sql/out/mom with the columnschannel,month,revenue,prev_revenue,growth_pct. Leave the previous-month value of the first month empty. - In /root/spk/sql/distinct.py (application
spk-sql-distinct), count the customers who made a paid order both exactly (countDistinct) and withapprox_count_distinct(rsd=0.05), and write them to /root/spk/sql/out/distinct.json as{"exact": 정수, "approx": 정수}(both values are integers). - In /root/spk/sql/report.md, write three sections:
## SQL 과 DataFrame,## 윈도 함수and## 근사 집계(use exactly these Korean headings in this order; they mean "SQL and DataFrame", "Window functions" and "Approximate aggregation"). Put the two numbers from step 7 in the third section.
Notes
- Original data:
/data/shop/orders.csv(order_id, customer_id, product_id, qty, order_ts, status, channel) and/data/shop/products.csv(product_id, category, price). Revenue isqty × priceof the paid (status='paid') orders. - If you put the scripts in the same folder as
common.py(/root/spk/sql),from common import loadworks (the folder of the script you ran is placed on the Python module path). - The result CSVs are small, so gathering them into one file with
coalesce(1)makes them easy to read by eye (the grader reads them even if there are several files). explainprints to standard output. Capture it withcontextlib.redirect_stdoutand write it to a file.- Common mistakes: using
rankso that two rows get the same rank on a tie; leaving the running-total window at its default (RANGE … CURRENT ROW when there is an ordering) so that rows of the same day overlap; and including canceled or refunded orders in revenue. - Official documentation: Spark SQL Guide · Window Functions · EXPLAIN · Built-in Functions
Monthly revenue with a temporary view and SQL
In /root/spk/sql/common.py, create load(spark), which reads orders and products with a schema and registers them as the temporary views orders and products. Then create /root/spk/sql/sql_monthly.py under the application name spk-sql-monthly and, using spark.sql, write the monthly revenue of paid orders (the column month is yyyy-MM, and revenue is the sum of qty×price) as a CSV with a header row to /root/spk/sql/out/monthly_sql.
A temporary view is a name visible only inside this SparkSession. Registering a view reads nothing; it only attaches a name to a plan. Build the month with date_format(order_ts, 'yyyy-MM'), then join with products and multiply by the price.
The same question with the DataFrame API
Create /root/spk/sql/df_monthly.py under the application name spk-sql-df and, without using spark.sql, write the same result as in step 1 using the DataFrame API (where, join, groupBy, agg) as a CSV with a header row (month,revenue) to /root/spk/sql/out/monthly_df.
Group with F.date_format("order_ts", "yyyy-MM").alias("month") and sum with F.sum(F.col("qty") * F.col("price")). If the join keys have the same name, pass them as a string like join(p, "product_id") so that only one column remains.
See whether the physical plans of the two approaches are the same
Create /root/spk/sql/plans.py under the application name spk-sql-plans, build the two queries from steps 1 and 2, save the explain(mode="formatted") output to /root/spk/sql/out/plan_sql.txt and /root/spk/sql/out/plan_df.txt, and then actually run both DataFrames with collect().
The formatted mode prints an operator tree at the top and numbered operator descriptions below it. The column numbers (such as #12) differ from query to query, but the names and order of the operators must be the same. When you actually run them, the identical plan document is left in the SQL execution record of the event log; the grader checks whether your files match that record and whether the operators of the two plans are the same.
Top 3 products per category: window ranking
Create /root/spk/sql/top3.py under the application name spk-sql-top3, compute the paid revenue per category (category) and product, and write ranks 1 to 3 within each category, sorted by revenue descending (for a tie, product_id ascending), as a CSV with a header row (category,product_id,revenue,rank) to /root/spk/sql/out/top3.
The window is Window.partitionBy("category").orderBy(...). rank() gives ties the same number and skips the next number, so if you want exactly 3 rows, you must use row_number() and put the tie rule into the ordering.
Running total per channel: setting the range of the window
Create /root/spk/sql/running.py under the application name spk-sql-running, and write the daily revenue per channel of paid orders in January 2026 (day is a date) and the running total running accumulated by date within each channel as a CSV with a header row (channel,day,revenue,running) to /root/spk/sql/out/running.
The running total is a sum over Window.partitionBy("channel").orderBy("day").rowsBetween(Window.unboundedPreceding, Window.currentRow). First collect it into one row per date, and then apply the window: if you apply it before collecting, the running total differs for each order of the same day.
Change from the previous month: looking at the neighboring row with lag
Create /root/spk/sql/mom.py under the application name spk-sql-mom, and write the monthly revenue per channel, the previous month's revenue prev_revenue fetched with lag, and the growth rate growth_pct (= (this month − previous month) / previous month × 100, rounded to two decimal places) as a CSV with a header row (channel,month,revenue,prev_revenue,growth_pct) to /root/spk/sql/out/mom. For the first month of a channel, the previous month and the growth rate must be empty.
F.lag("revenue").over(Window.partitionBy("channel").orderBy("month")) fetches the value of the immediately preceding month of the same channel. The first month has no previous row, so it is null, and a division involving null is also null, so no special handling is needed.
Exact distinct count and approximate distinct count
Create /root/spk/sql/distinct.py under the application name spk-sql-distinct, count the customers who made a paid order with both countDistinct and approx_count_distinct(rsd=0.05) in one pass, and write them to /root/spk/sql/out/distinct.json as {"exact": 정수, "approx": 정수} (both values are integers).
The approximate side merges HyperLogLog++ sketches. Distinct values do not have to be gathered by a shuffle, so it gets cheaper as the data grows. The rsd is a target for the relative standard error, so calculate yourself by what percent the result differs from the exact value. The grader also checks whether the approximate function is actually in the plan in the event log.
Record the same plan, the frame of the window and the error of the approximation
In /root/spk/sql/report.md, write three sections: ## SQL 과 DataFrame, ## 윈도 함수 and ## 근사 집계 (use exactly these Korean headings in this order; they mean "SQL and DataFrame", "Window functions" and "Approximate aggregation"). Put the exact value and the approximate value from step 7 in the third section as numbers.
In the first section, list a few operator names that were the same in the two plans, and in the second section, write one line each on how you set the window for the ranking, the running total and the change from the previous month. In the third section, it is enough to write the two numbers and the percentage of their difference.