Every driver serves a web UI while the application runs, on port 4040 by default (33040 on this workstation), and the same data as JSON under /api/v1. The listing labels each workload with setJobDescription so its jobs are easy to find, then holds the session open for the screenshots in this section.
sc = spark.sparkContext
sc.setJobDescription("genre revenue") # names the jobs in the UI
genre = (lines.where("status = 'delivered'")
.join(books.withColumnRenamed("id", "book_id"), "book_id")
.groupBy("genre").agg(F.sum(F.col("qty") * F.col("unit_price")).alias("revenue")))
genre.collect()
sc.setJobDescription("status join") # 89% of lines: 'delivered'
spark.conf.set("spark.sql.autoBroadcastJoinThreshold", "-1") # force a shuffle join
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "false")
spark.conf.set("spark.sql.adaptive.advisoryPartitionSizeInBytes", "4m")
phase = spark.createDataFrame([("delivered", "closed"), ("returned", "closed"),
("cancelled", "closed"), ("shipped", "open"), ("paid", "open"),
("pending", "open")], "status STRING, phase STRING")
lines.join(phase, "status").write.format("noop").mode("overwrite").save()
sc.setJobDescription("cache order lines")
lines.cache().count()
print("UI:", sc.uiWebUrl, "application:", sc.applicationId)Output
UI: http://10.255.255.254:33040 application: local-1790878201467
| Tab | Answers |
|---|---|
| Jobs | Which actions ran, how long, how many stages were skipped |
| Stages | Per-stage and per-task time, shuffle, spill, GC |
| Storage | What is cached, at which level, how much fits |
| Environment | Every configuration value actually in effect |
| Executors | Memory, cores, task counts and failures per executor |
| SQL / DataFrame | Each query's plan with live metrics on every operator |
Read top-down: find the slow job, open its slowest stage, then compare that stage's tasks with each other. Check Environment first when a setting "doesn't work"; it shows what the driver received, not what you meant.