Data skew là khi dữ liệu sau shuffle chia không đều theo key: vài partition chứa phần lớn số dòng. Stage chỉ xong khi task chậm nhất xong, nên 199 task chạy 1 phút còn 1 task chạy 40 phút.
Ví dụ: join events với users theo user_id, nhưng 30% event có user_id là NULL hoặc thuộc vài tài khoản bot → tất cả rơi vào cùng một partition.
Phát hiện: Spark UI → Stages → Summary Metrics. Max duration và Shuffle Read Size lớn gấp nhiều lần median là dấu hiệu rõ nhất. Sau đó đếm theo key để tìm key nóng:
SELECT user_id, COUNT(*) AS c FROM events GROUP BY user_id ORDER BY c DESC LIMIT 20;Xử lý, từ rẻ đến tốn công:
- AQE skew join (bật mặc định từ Spark 3.2): spark.sql.adaptive.skewJoin.enabled tự chẻ partition lệch (lớn hơn 5 lần median và trên 256MB) thành nhiều task. Chỉ áp dụng cho sort-merge join.
- Tách key rác: lọc NULL hoặc key mặc định ra xử lý riêng rồi union lại.
- Broadcast bảng bên kia nếu đủ nhỏ, khi đó không còn shuffle theo key.
- Salting: thêm số ngẫu nhiên vào key ở bảng lớn, nhân bản bảng nhỏ theo số salt, join theo key đã salt.
from pyspark.sql import functions as F
N = 16
big = events.withColumn("salt", (F.rand() * N).cast("int"))
small = users.crossJoin(spark.range(N).withColumnRenamed("id", "salt"))
joined = big.join(small, ["user_id", "salt"])- Với aggregation bị lệch: aggregate hai tầng — gom theo (key, salt) trước, gom theo key sau.
Lưu ý: tăng executor memory hay tăng shuffle.partitions không chữa được skew, vì key nóng vẫn dồn về một partition. Phải biết key nào gây lệch trước khi chọn cách xử lý.