Diagnosing Data Skew

A shuffle sends every row with the same key to the same task, so one popular key makes one task do most of the work while the others wait. BookNest's catalog is an example: a GROUP BY book_id shows The Quiet Harbor (book 1) on 415,087 order lines, 30.1% of all of them. The listing forces a sort-merge join between order lines and books, as if the catalog were too large to broadcast, and reads each task's shuffle input from the Spark 129 UI's REST API, the numbers the Stages tab shows (Spotting Skew).

Measuring skew in the order-lines joinJavaScript
import json, urllib.request
sc = spark.sparkContext
API = f"{sc.uiWebUrl}/api/v1/applications/{sc.applicationId}/stages"   # the UI's REST API
def join_tasks(label, df):            # run df, then report its busiest shuffle stage's tasks
    sc.setJobGroup(label, label)
    df.write.format("noop").mode("overwrite").save()
    get = lambda path: json.load(urllib.request.urlopen(f"{API}/{path}"))
    reads = [sorted(t["taskMetrics"]["shuffleReadMetrics"]["recordsRead"]
                    for t in get(f"{s}/{a['attemptId']}/taskList?length=1000"))
             for j in sc.statusTracker().getJobIdsForGroup(label)
             for s in sc.statusTracker().getJobInfo(j).stageIds
             for a in get(s) if a["status"] == "COMPLETE"]
    r = max(reads, key=sum)
    print(f"{label:<7} {len(r)} tasks, median {r[len(r) // 2]:>7,} rows, max {r[-1]:>7,}")
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1")      # pretend books is huge
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "4m")   # 64m, scaled down
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "false")
catalog = books.withColumnRenamed("id", "book_id")
join_tasks("skewed", lines.join(catalog, "book_id"))
Output
skewed  5 tasks, median 269,311 rows, max 415,088

After AQE coalesced the 200 shuffle partitions, five tasks did the join and one read 415,088 rows (book 1's lines plus its catalog row). No partition count can spread one key: a hash partitioner sends equal keys to one task by definition. Watch for a task whose time or shuffle read is several times the median, idle cores at the end of a stage, and spills or out-of-memory errors in one task only.