Unbalanced partitions leave a few tasks running long after the others finish. In Summary Metrics on the Spark UI, Max duration 50% or more higher than the 75th percentile points to it.
Also called data skew.
Read more: Microsoft Learn
In the Ultra Transcenders books
Each book explains Skew in context, with comparison tables and the common traps.
Terms in this definition
- Spark UI
Classic compute's Apache Spark web console. Databricks advises checking the jobs timeline first, then the slowest stage, skew or spill, I/O and lastly the SQL DAG. Serverless and SQL warehouses offer a query profile in its place.
- MAX
Gives back the biggest value found in a column holding numbers, dates or text, or compares two expressions and returns the bigger. For TRUE/FALSE values, which this DAX aggregation doesn't accept, turn to MAXA.
Related terms
- AQE
Adaptive query execution: Spark re-plans a query while it runs, for instance turning a sort-merge join into a broadcast join, merging shuffle partitions, dealing with skew and propagating empty relations. Databricks advises leaving it enabled.
- Spark advisor
A Fabric feature that examines Spark code and runs as it executes, giving error, warning and information messages both under notebook cells and in an application's Diagnostics panel. Its checks include spotting data skew, explaining the root cause of errors and flagging when work falls back from the native engine.
- Spark History Server
After a Spark application has ended (whether it Completed, Failed, Stopped or was Canceled), Fabric offers this enhanced web UI for it. Event logs are replayed to show how the run went, and an extra Diagnosis tab points out skew in data or time.