Khác nhau ở chỗ một partition đầu ra cần dữ liệu từ bao nhiêu partition đầu vào.
- Narrow (
filter,select,map,union): mỗi partition đầu ra chỉ phụ thuộc một partition đầu vào, nên xử lý ngay trên executor đang giữ dữ liệu, không di chuyển qua mạng. - Wide (
groupBy,joinhai bảng lớn,distinct,orderBy,repartition): các dòng cùng key đang rải khắp các partition phải được gom về một chỗ — đó là shuffle.
Shuffle tốn vì: executor ở phía trước ghi dữ liệu đã chia theo key xuống đĩa cục bộ, executor phía sau kéo về qua mạng, kèm serialize và deserialize; thiếu RAM thì còn spill thêm ra đĩa. Đĩa, mạng và CPU cùng bị dùng. Shuffle cũng là ranh giới stage: Spark cắt DAG tại mỗi wide dependency, stage sau phải chờ stage trước ghi xong.
Cách giảm shuffle thường dùng:
python
from pyspark.sql.functions import broadcast
# small dimension table: ship it to every executor instead of shuffling the big table
orders.join(broadcast(countries), "country_code")- Lọc dòng và chọn cột trước
join/groupByđể khối lượng shuffle nhỏ lại. - Broadcast bảng nhỏ; Spark tự broadcast khi bảng dưới
spark.sql.autoBroadcastJoinThreshold(mặc định 10MB). - Chỉnh
spark.sql.shuffle.partitions(mặc định 200) hoặc để AQE tự gộp các partition nhỏ sau shuffle.
Lưu ý: trên Spark UI, cột Shuffle Read/Write của từng stage cho biết nhanh nhất job đang tốn tài nguyên ở đâu.