ML Engineer MasterClass (October) | 4 seats left

Spark · Data skew
AmazonAmazon Analytics
Amazon · Tune Spark Workloads

Data skew

Measure uneven keys and understand where a two-stage aggregation can help.

Step 1 of 6 · Learn

Measure key frequency

A hot key occurs much more often than others and can concentrate shuffle work. Count rows by key before deciding on a remedy. This small dataset illustrates frequency only: it is not evidence of a slow or skewed production job. Task duration, row size, spill, and shuffle-read metrics provide additional evidence.

Lesson reference: PySpark and Scala

Measure key frequency

A hot key occurs much more often than others and can concentrate shuffle work. Count rows by key before deciding on a remedy. This small dataset illustrates frequency only: it is not evidence of a slow or skewed production job. Task duration, row size, spill, and shuffle-read metrics provide additional evidence.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = (orders.groupBy("status").agg(count("*").alias("row_count"))
    .select(col("status").cast("string").alias("key"), col("row_count"))
    .orderBy(col("row_count").desc(), col("key").asc()))
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = (orders.groupBy("status").agg(count("*").alias("row_count"))
    .select(col("status").cast("string").alias("key"), col("row_count"))
    .orderBy(col("row_count").desc(), col("key").asc()))
result.show(false)

Isolate repeated keys

Filter a frequency report to examine repeated keys. A threshold is a diagnostic choice, not a universal definition of skew. Increasing the number of hash partitions cannot split one identical grouping key across multiple final aggregation groups. First confirm which operator and key are causing expensive work.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = (orders.groupBy("customer_id").agg(count("*").alias("row_count"))
    .filter(col("row_count") >= 1).orderBy("customer_id"))
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = (orders.groupBy("customer_id").agg(count("*").alias("row_count"))
    .filter(col("row_count") >= 1).orderBy("customer_id"))
result.show(false)

Split and recombine an aggregate

Salting adds a secondary key so one logical group can be aggregated in pieces, then recombined. Here order_id modulo 2 defines deterministic salts. sum is mergeable, so summing partial sums preserves the answer. This adds work and is not automatically faster; ordinary sum already supports partial aggregation. Do not average partial averages without their counts, or apply this recipe directly to joins.

PySpark example

from pyspark.sql.functions import col, broadcast, count, sum

result = (orders.withColumn("salt", col("order_id") % 2)
    .groupBy("status", "salt").agg(sum("total").alias("partial_total"))
    .groupBy("status").agg(sum("partial_total").alias("total_value"))
    .orderBy("status"))
result.show(truncate=False)

Scala example

import org.apache.spark.sql.functions._

val result = (orders.withColumn("salt", col("order_id") % 2)
    .groupBy("status", "salt").agg(sum("total").alias("partial_total"))
    .groupBy("status").agg(sum("partial_total").alias("total_value"))
    .orderBy("status"))
result.show(false)
example.pyPySpark
1. Count keysMeasure frequency before choosing a remedy.
2. Compare workFrequent keys can concentrate work; row width and computation also matter.
3. Inspect executionUse task durations and shuffle metrics on representative data; six rows cannot demonstrate production skew.
Source orders6 rows
order_idcustomer_idstatustotalitem_count
1001101Delivered89.52
1002102Shipped1493
1003101Cancelled351
1004103Delivered219.994
1005104Delivered49.991
1006105Shipped1202
Measure key frequency
Result3 rows
keyrow_count
Delivered3
Shipped2
Cancelled1
Delivered is the largest status group with three rows; this alone does not establish a performance problem.

The example is loaded in the editor. Run it as written, then try a small change.

Runs on the Spark backend. First startup may take a moment.

Run your code to see the result.