Apache Spark — The answer to a slow job is in the plan and the event log
Use explain to confirm how much less Parquet is read
Goal
After converting the orders to Parquet, you read from the execution plan how much less of the file the same question reads depending on how you write the condition (pushdown, column pruning and partition pruning). You also confirm with the event log how adaptive execution (AQE) changed the plan after it ran.
Why it matters
When Spark is slow, the first thing to look at is not the code but the plan. Of two lines that produce the same result, one reads only the needed columns and rows from the file, and the other reads everything and then throws it away. That difference does not show in the result; it shows only in the plan.
A columnar format such as Parquet can do two things. It reads only the needed columns (column pruning, ReadSchema in the plan), and it skips chunks that do not match the condition using the minimum and maximum statistics of row groups (filter pushdown, PushedFilters in the plan). On top of that, if you select directory-based partitions with a condition, whole directories are skipped (PartitionFilters). All three happen only when the optimizer can understand the condition. The moment you wrap a column in a function, pushdown disappears.
AQE looks at the actual sizes after a shuffle has finished and rewrites the rest of the plan. So the plan before running (isFinalPlan=false) and the plan after running can differ, and what really ran is seen in the plan after the run or in the event log.
Steps
- With /root/spk/plan/convert.py (application
spk-plan-convert), read the orders with a schema, add anorder_date(date) column, and write them as Parquet to /root/spk/plan/orders_pq. - Put the functions
explain,fieldanditemsin the helper /root/spk/plan/planhelp.py, and with /root/spk/plan/explain.py (applicationspk-plan-explain), save theextendedandformattedplans of a query that selectsorder_idandqtyof refunded orders to /root/spk/plan/out/extended.txt and /root/spk/plan/out/formatted.txt, andcollect()that query. - With /root/spk/plan/pushdown.py (application
spk-plan-push), writeorder_id,qty,channelof the orders withstatus == 'refunded'andqty >= 5as a CSV to /root/spk/plan/out/pushdown, and write the list ofPushedFiltersitems and theReadSchemaof that plan to /root/spk/plan/out/pushdown.json. - With /root/spk/plan/nopush.py (application
spk-plan-nopush), change the same question toupper(status) == 'REFUNDED'and write it to /root/spk/plan/out/nopush, and write the list ofPushedFiltersof that plan to /root/spk/plan/out/nopush.json. - With /root/spk/plan/partition.py (application
spk-plan-part), write the data split byorder_dateto /root/spk/plan/orders_by_day, pick the single day 2026-02-14, and write the number of rows counted and thePartitionFiltersof the plan to /root/spk/plan/out/partition.json. - With /root/spk/plan/aqe.py (application
spk-plan-aqe),collect()the number of orders per customer, save the plan to /root/spk/plan/out/aqe_final.txt, count the tasks that read the shuffle from that application's log, and write them to /root/spk/plan/out/aqe.json as{"shuffle_partitions": 200, "reduce_tasks": 정수}(the second value is an integer). - With /root/spk/plan/optimize.py (application
spk-plan-opt), save theextendedplan of a query that chains twowherecalls,qty > 1andqty > 3, and selects1 + 2asthree, to /root/spk/plan/out/optimized.txt, andcollect()it. - In /root/spk/plan/report.md, write three sections:
## 밀어내기,## 가지치기and## AQE(use exactly these Korean headings in this order; they mean "Pushdown", "Pruning" and "AQE"). Put the number of rows from step 5 in the second section and the number of tasks from step 6 in the third section as numbers.
Notes
- Put the scripts in
/root/spk/planand runspark-submitfrom there. You useplanhelp.pyin the same folder withfrom planhelp import explain. - The modes of
explain:simple(physical plan only),extended(parsed, analyzed and optimized logical plans plus the physical plan),codegen,cost,formatted(physical plan tree plus per-operator details). The per-operatorPushedFilters,ReadSchemaandPartitionFiltersappear one per line in formatted. - In the event log, a task that read the shuffle is one whose
Task Metrics→Shuffle Read Metrics→Total Records ReadinSparkListenerTaskEndis greater than 0. - Common mistakes: looking for pushdown in a CSV (a CSV has to be read in full); printing the plan before
collect()and treating it as the AQE final plan; and concluding that less of the file was read just becausePushedFiltersis present (a chunk is skipped only when the row-group statistics can separate the condition). - Official documentation: Performance Tuning · EXPLAIN · Parquet Files · Adaptive Query Execution
CSV to Parquet
Create /root/spk/plan/convert.py under the application name spk-plan-convert, read /data/shop/orders.csv with a schema, add the column order_date = to_date(order_ts), and write it as Parquet to /root/spk/plan/orders_pq (without partitioning).
A columnar file is stored column by column and has minimum and maximum statistics for each row group. These two things are the raw material of "reading less" in all the later steps. The grader checks the row count, the list of columns, and whether the file is really Parquet (the first four bytes are PAR1).
The modes of explain: from logical plan to physical plan
In the helper /root/spk/plan/planhelp.py, put explain(df, mode), which captures the explain output as a string; field(plan, name), which extracts the value of a line of the form 이름: 값 (name: value) in a formatted plan; and items(text), which splits [A(x,1), B(y)] into a list of items. Then create /root/spk/plan/explain.py under the application name spk-plan-explain, save the extended and formatted plans of a query that selects order_id and qty of the rows of orders_pq with status == 'refunded' to /root/spk/plan/out/extended.txt and /root/spk/plan/out/formatted.txt, and then collect() that query.
extended prints four sections in order: the parsed logical plan, the analyzed logical plan, the optimized logical plan and the physical plan. In the analysis stage, column names are bound to actual columns (the #number), and in the optimization stage the condition moves down close to the scan. formatted spells out that physical plan operator by operator. The grader checks whether the formatted file matches the plan that actually ran (the event log).
Filter pushdown and column pruning
Create /root/spk/plan/pushdown.py under the application name spk-plan-push, write order_id, qty and channel of the rows of orders_pq with status == 'refunded' and qty >= 5 as a CSV with a header row to /root/spk/plan/out/pushdown, and write the list of items of PushedFilters and the value of ReadSchema from the formatted plan of that query to /root/spk/plan/out/pushdown.json as {"pushed_filters": [문자열…], "read_schema": 문자열} (a list of strings and a string).
PushedFilters is the condition handed to the Parquet read, and ReadSchema is the columns actually read. Notice that the customer_id and product_id you did not select drop out of ReadSchema, while status, which was used only in the condition, remains. The grader checks whether the list you wrote is identical, character by character, to the one in the plan that actually ran.
Wrap it in a function and pushdown disappears
Create /root/spk/plan/nopush.py under the application name spk-plan-nopush, change the same question as in step 3 to F.upper("status") == "REFUNDED" and qty >= 5, write it as a CSV with a header row to /root/spk/plan/out/nopush, and write the list of items of PushedFilters of that plan to /root/spk/plan/out/nopush.json as {"pushed_filters": [문자열…]} (a list of strings).
The result rows are exactly the same as in step 3. What changes is the plan. Parquet does not know what upper(status) is, so that condition is not handed to the read, and Spark filters after reading the whole file. In practice, people often fall into this trap by comparing date columns wrapped in date_format.
Partition pruning: skipping whole directories
Create /root/spk/plan/partition.py under the application name spk-plan-part, write orders_pq with partitionBy("order_date") to /root/spk/plan/orders_by_day, read it back, and write the number of rows counted with order_date == '2026-02-14' and the PartitionFilters value of that query's plan to /root/spk/plan/out/partition.json as {"rows": 정수, "partition_filters": 문자열} (an integer and a string).
The column you partitioned by is stored not inside the file but in the directory name (order_date=2026-02-14). A condition on that column is used to choose directories, and the files of the dates you did not choose are not even opened. The grader also checks in the event log whether the "number of partitions read" metric is 1.
The plan AQE rewrote after running
Create /root/spk/plan/aqe.py under the application name spk-plan-aqe, and after you collect() the number of orders per customer of orders_pq (groupBy("customer_id").count()), save the explain(mode="simple") output to /root/spk/plan/out/aqe_final.txt. Then count the tasks that read the shuffle in that application's event log and write them to /root/spk/plan/out/aqe.json as {"shuffle_partitions": 200, "reduce_tasks": 정수} (the second value is an integer).
spark.sql.shuffle.partitions is 200, but the shuffle of this data is only a few MB, so AQE merges the partitions. In the plan after running, you see isFinalPlan=true and the merged shuffle read (AQEShuffleRead). The number of tasks that read the shuffle comes from counting the tasks in the log's SparkListenerTaskEnd whose Total Records Read is greater than 0.
What the optimizer merges and folds
Create /root/spk/plan/optimize.py under the application name spk-plan-opt, save the extended plan of orders_pq.where(qty > 1).where(qty > 3).select("order_id", (lit(1) + lit(2)).alias("three")) to /root/spk/plan/out/optimized.txt, and collect() that query.
In the analyzed logical plan there are two Filters and (1 + 2) is still there. In the optimized logical plan, the two conditions are merged into one Filter (filter merging) and 1 + 2 is folded into 3 (constant folding). This is the step where you see that the shape of the code you wrote differs from the shape that actually runs.
Put the amount you read less into numbers
In /root/spk/plan/report.md, write three sections: ## 밀어내기, ## 가지치기 and ## AQE (use exactly these Korean headings in this order; they mean "Pushdown", "Pruning" and "AQE"). Put the number of rows from step 5 in the second section and the number of shuffle-reading tasks from step 6 in the third section as numbers.
In the first section, write how PushedFilters changed between steps 3 and 4; in the second, what the partition pruning skipped; and in the third, how many the 200 was reduced to.