Optimize a Spark SQL plan
SkillDatabases & dataOptimize Apache Spark SQL and DataFrame queries using the final Adaptive Query Execution plan and runtime statistics rather than source code alone. Use to reduce runtime, shuffle, spill, scan cost, skew, join amplification, Python UDF overhead, poor partitioning, or unnecessary work while preserving query semantics.
Available today. Use it from your connected AI after setup.
No other account needed.
Connect ahel once, and every AI you use reads what you have installed.
Then ask your AI: use the Optimize a Spark SQL plan skill
What this skill tells your AI
The instructions your AI receives, as published by embrasureai/spark-observability-skills in skills/optimize-spark-sql-plan/SKILL.md and read by ahel’s review.
Find the highest-impact improvement supported by the executed plan and its runtime metrics. Work from the final adaptive plan (isFinalPlan=true), not source code or the initial plan: AQE may already have broadcast, coalesced, or split what you were about to recommend.
Get the evidence
export SPARK_HISTORY_URL="https://history.example.com"
python3 scripts/spark_history_api.py sql-list --app-id <application-id>
python3 scripts/spark_history_api.py sql --app-id <application-id> --execution-id <execution-id> > /tmp/spark-sql.json
Run from this skill directory. If SPARK_HISTORY_URL is unset, find the server before asking the user: try http://localhost:18080, a running application's UI on http://localhost:4040, and the history-server or eventLog settings in the local Spark config; ask only when nothing responds. Authentication comes from SPARK_HISTORY_AUTHORIZATION, SPARK_HISTORY_COOKIE, or SPARK_HISTORY_HEADERS_JSON; never ask for credentials in chat and never disable TLS verification (--ca-file for a private CA).
The snapshot's layout: sqlExecution.planDescription is the executed plan text, sqlExecution.nodes[].metrics carry per-operator counters with edges giving parent-child links, jobs[].stageIds tie the execution to stages, longestStages[] holds each stage with its taskSummary quantiles (arrays align with quantiles; middle is median, last is max), and environment.sparkProperties is the effective configuration. For data beyond the profile, other subcommands (slow, failure, applications) and the raw $SPARK_HISTORY_URL/api/v1 endpoints (/stages/{stage}/{attempt}/taskSummary, /allexecutors, /environment) are available with the same auth. If the API omits the plan, recover it from the driver log or generate explain(mode="formatted") from the deployed code yourself. Get the deployed query code too: the plan tells you what ran, the code tells you what was intended.
Highest-impact opportunities, in order, and how to check each
-
Cardinality blowup. Join output far exceeding both inputs is usually an unintended many-to-many:
jq '.sqlExecution.nodes[] | select(.nodeName | test("Join")) | {nodeId, nodeName, rows: [.metrics[] | select(.name | test("output rows"))]}' /tmp/spark-sql.jsonCompare each join's rows against its children's (follow
edges). Fix keys or deduplicate first; it dominates every downstream metric. -
Pruning. Scans reading far more than the query uses:
jq '.sqlExecution.nodes[] | select(.nodeName | startswith("Scan")) | {nodeName, metrics}' /tmp/spark-sql.jsonCompare files and bytes read against rows output, and check the pushed/partition filters on the scan in
planDescription. -
Broadcast. A
SortMergeJoinwhose smaller side is observed (not estimated) to be broadcastable:jq -r '.sqlExecution.planDescription' /tmp/spark-sql.json | grep -n 'SortMergeJoin\|BroadcastHashJoin'Read the smaller side's "data size" metric from its node, and
spark.sql.autoBroadcastJoinThresholdfromenvironment.sparkProperties. Stale statistics are the usual reason AQE missed it; confirm executor memory headroom before recommending. -
Unnecessary shuffles. Exchanges whose partitioning an upstream operation already satisfies, or repartitions the query does not need:
jq -r '.sqlExecution.planDescription' /tmp/spark-sql.json | grep -n 'Exchange'Every exchange should map to a semantic requirement; look for back-to-back exchange/sort pairs and matching partitioning expressions above and below.
-
Skewed joins and aggregations. One partition dominating a join or aggregate stage:
jq '{jobStages: [.jobs[] | {jobId, stageIds}], stages: [.longestStages[] | {id: .stage.stageId, run: .taskSummary.executorRunTime, shuffleRead: .taskSummary.shuffleReadMetrics.readBytes}]}' /tmp/spark-sql.jsonMax many times median on the join's stage is skew; check
planDescriptionforskewed=trueand, if AQE did not handle it, why. Re-run with a higher--stage-limitif the stage is missing. -
Partition count. Uniformly oversized shuffle partitions with spill:
grep -n 'AQEShuffleRead' <(jq -r '.sqlExecution.planDescription' /tmp/spark-sql.json)for coalescing, stagenumTasksfromlongestStages[].stage, and spill from the check below. Raisespark.sql.shuffle.partitionsor AQE targets before adding explicit repartitions, which can insert redundant exchanges. -
Excess spill in stateful operators. Aggregates, sorts, and windows spilling:
jq '.sqlExecution.nodes[] | select(.nodeName | test("Aggregate|Sort|Window")) | {nodeName, spill: [.metrics[] | select(.name | test("spill"))]}' /tmp/spark-sql.jsonReduce state (partial aggregation, narrower rows, fewer window columns) before adding memory.
-
Python UDF boundaries. Serialization often costs more than the function itself:
jq '.sqlExecution.nodes[] | select(.nodeName | test("Python")) | {nodeName, metrics}' /tmp/spark-sql.jsonWeigh node time against rows processed; a native expression removes the boundary.
-
Caching.
jq -r '.sqlExecution.planDescription' /tmp/spark-sql.json | grep -c 'InMemoryTableScan': a cached result scanned once wastes memory, and a subtree recomputed under several plan branches wants a cache.
The commands above are starting points, not limits: compose your own jq, call the REST API directly, or pull logs, code, and platform state when a question needs it. Rule each candidate in or out with evidence, and when a signal is suggestive but not conclusive, go a level deeper (more task samples, the exact log lines, the executed plan) until you are confident it is or is not the cause. Do not settle for the first plausible explanation.
An Exchange is not automatically waste, a sort-merge join is not automatically wrong, a broadcast is not automatically safe, and a Python UDF is not automatically material: tie every recommendation to the node's runtime metrics, and check whether Catalyst or AQE already applied it.
Changes must preserve semantics: row counts, join cardinality, null behavior, duplicates, and ordering assumptions. Changing partitioning can alter low-order bits of floating-point aggregations.
Report
State the highest-impact change with its plan node, runtime evidence, and expected effect; why Spark did not already do it; correctness risks; and any remaining smaller findings. Do not apply hints, code, or configuration changes to production without approval.
Signals
- GitHub stars
- 52
- Forks
- 13
- Last commit
- Aug 2026
Advanced
- Catalog kind
- skill
- Gateway key
optimize-spark-sql-plan- Source
- github.com/embrasureai/spark-observability-skills