Bài 6

Tối ưu & Spark UI

Hiểu Spark chạy thế nào bên dưới để viết job nhanh hơn và đỡ tốn tài nguyên. Bài này tập trung vào cách Spark chia việc, đọc Spark UI để soi điểm nghẽn, và những kỹ thuật tối ưu thực dụng nhất: giảm shuffle, đọc query plan, cache, repartition, broadcast join và AQE.

Job → Stage → Task
Spark chia nhỏ công việc như thế nào

Mỗi khi bạn gọi một action (như count(), show(), write(), collect()), Spark tạo ra 1 job. Job đó được chia thành nhiều stage, và mỗi stage lại chia thành nhiều task.

Ranh giới giữa các stage chính là shuffle. Mỗi lần data cần xáo trộn qua mạng (groupBy, join, distinct...), Spark phải cắt sang một stage mới. Trong cùng một stage, các phép biến đổi được "xâu chuỗi" (pipeline) chạy liền trên từng partition mà không cần trao đổi data.

Mỗi task xử lý đúng 1 partition của data. Vì vậy số task trong một stage = số partition của data ở stage đó. Nhiều partition → nhiều task chạy song song (tốt cho data lớn), nhưng quá nhiều partition nhỏ lại tốn chi phí quản lý.

Phân cấp Job - Stage - Task Scroll / zoom · Mở draw.io ↗
Đọc Spark UI
Công cụ soi job mạnh nhất, miễn phí

Khi có job đang chạy, mở trình duyệt tới http://localhost:4040 để xem Spark UI của ứng dụng. Đây là nơi bạn nhìn thấy chính xác Spark đang làm gì và tốn thời gian ở đâu.

TabCho bạn xem gì
JobsDanh sách job + tiến độ, thời gian chạy mỗi job, job nào đang chạy / đã xong / lỗi
StagesChi tiết từng stage: thời gian, số task, và quan trọng nhất là cột Shuffle Read / Shuffle Write
StorageCác DataFrame/RDD đã cache() — kích thước, nằm ở RAM hay đĩa, đã materialize bao nhiêu
SQL / DataFrameSơ đồ query plan trực quan của từng query — thấy rõ shuffle (Exchange), join, scan
ExecutorsRAM / CPU mỗi executor, số task đã chạy, lượng data xử lý — soi executor nào bị quá tải
Mẹo soi nhanh: vào tab Stages, tìm stage tốn thời gian nhất; nếu cột Shuffle Read / Write to thì nghi ngay shuffle nặng. Dấu hiệu data skew: trong một stage có 1 task chạy lâu hơn hẳn các task còn lại (các task khác xong từ lâu mà cả stage vẫn chờ một task).
Shuffle — kẻ thù số 1
Vì sao job chậm, gần như luôn là shuffle

Shuffle là khi data phải di chuyển qua mạng giữa các partition để gom các dòng cùng key về chung một chỗ. Nó xảy ra ở: groupBy, join, distinct, repartition, orderBy (sắp xếp toàn cục).

Shuffle tốn vì Spark phải ghi data ra đĩa ở stage trước, rồi truyền qua mạng, rồi đọc lại ở stage sau. Đó là thao tác đắt nhất trong một job.

Narrow vs Wide Transformation — Shuffle & Stage Scroll / zoom · Mở draw.io ↗

Cách giảm shuffle:

Data skew: nếu 1 key có quá nhiều dòng (vd cột country mà 80% là "VN"), thì khi groupBy/join, một task phải ôm trọn key đó → chạy mãi không xong dù các task khác đã rảnh. Cách giảm phổ biến là salting: thêm một hậu tố ngẫu nhiên vào key để tách key "nóng" thành nhiều phần, xử lý song song rồi gộp lại. (Spark 3 có AQE tự xử skew phần nào — xem mục cuối.)
explain() — đọc query plan
Xem Spark định làm gì trước khi nó chạy

Trước khi chạy, Spark tối ưu query của bạn thành một plan. Dùng explain() để xem plan đó và phát hiện shuffle thừa hay join chọn sai kiểu.

# chỉ physical plan (ngắn gọn)
df.explain()

# đầy đủ: parsed → analyzed → optimized → physical plan
df.explain(True)

Đọc plan từ dưới lên (dưới cùng là bước chạy đầu tiên: đọc file). Vài thứ cần để ý:

Thấy trong planNghĩa là
ExchangeMột lần shuffle — data bị xáo trộn qua mạng (đây là chỗ tốn)
BroadcastHashJoinJoin kiểu broadcast — bảng nhỏ được gửi đi, không shuffle bảng lớn (nhanh)
SortMergeJoinJoin kiểu sort-merge — cả hai bảng đều phải shuffle + sort (đắt hơn)
FileScan + PushedFiltersĐọc file có đẩy điều kiện lọc xuống tận tầng đọc (predicate pushdown) — đọc ít data hơn
Quy tắc nhanh: thấy nhiều Exchange trong plan = nhiều shuffle = nhiều khả năng chậm. Mục tiêu tối ưu thường là giảm số Exchange và đổi SortMergeJoin sang BroadcastHashJoin khi có thể.
cache / persist
Đừng bắt Spark tính lại từ đầu

Spark lazy: nó chỉ thực sự tính khi gặp action, và mặc định mỗi action sẽ tính lại từ đầu cả chuỗi biến đổi. Nếu bạn dùng lại một DataFrame cho nhiều action (đếm, ghi, hiển thị...), hãy cache() để Spark giữ kết quả lại, khỏi đọc + tính lại mỗi lần.

df.cache()          # đánh dấu muốn cache (vẫn lazy)
df.count()          # 1 action để MATERIALIZE — giờ data mới nằm trong cache

# ...dùng lại df nhiều lần ở đây, rất nhanh...

df.unpersist()      # giải phóng khi không cần nữa

cache() tương đương persist(MEMORY_AND_DISK). Bạn có thể chọn StorageLevel khác qua persist():

StorageLevelLưu ở đâu
MEMORY_ONLYChỉ RAM. Nhanh nhất, nhưng không đủ RAM thì phần thừa bị tính lại
MEMORY_AND_DISKƯu tiên RAM, tràn thì xuống đĩa. Mặc định của cache() với DataFrame
DISK_ONLYChỉ đĩa. Chậm hơn RAM nhưng không ngốn bộ nhớ
Đừng cache bừa: cache chiếm RAM và tự nó cũng tốn công. Chỉ cache những DataFrame được dùng lại nhiều lần và đủ nhỏ để vừa RAM. Dùng một lần rồi thôi thì cache chỉ làm chậm thêm.
repartition vs coalesce
Chỉnh số partition cho đúng
repartition(n)coalesce(n)
Tác dụngTăng hoặc giảm số partitionChỉ giảm số partition
Shuffle? shuffle (đắt)KHÔNG shuffle (rẻ)
Cân bằngChia lại data đều các partitionGộp partition sẵn có — có thể lệch
# gộp về ÍT file lớn khi ghi (không shuffle) — vd xuất 1 file
df.coalesce(1).write.mode("overwrite").parquet("out/")

# chia lại ĐỀU thành 10 partition để xử lý song song tốt hơn (có shuffle)
df.repartition(10).write.mode("overwrite").parquet("out/")
Chọn cái nào: muốn ghi ít file lớn (giảm số file output) → dùng coalesce vì nó không tốn shuffle. Muốn chia lại đều để chạy song song nhiều task hơn (hoặc tăng số partition) → dùng repartition.
Broadcast join
Mẹo join nhanh nhất khi 1 bảng nhỏ

Khi join một bảng lớn với một bảng nhỏ (vd bảng giao dịch khổng lồ join với bảng danh mục vài nghìn dòng), thay vì shuffle cả hai bảng, Spark có thể broadcast bảng nhỏ — gửi nguyên bản sao của nó tới mọi executor. Nhờ vậy bảng lớn không cần shuffle, mỗi task join cục bộ với bản copy của bảng nhỏ.

from pyspark.sql import functions as F

# ép broadcast bảng nhỏ (small) khi join với bảng lớn (big)
result = big.join(F.broadcast(small), "id")

Thực ra Spark tự broadcast nếu bảng đủ nhỏ. Ngưỡng do spark.sql.autoBroadcastJoinThreshold quyết định (mặc định ~10MB). Bảng nhỏ hơn ngưỡng → Spark tự dùng broadcast; lớn hơn → nó chuyển sang SortMergeJoin (có shuffle). Dùng F.broadcast() khi bạn chắc bảng nhỏ nhưng Spark chưa tự nhận ra.

Lưu ý: chỉ broadcast bảng thật sự nhỏ. Broadcast bảng lớn sẽ làm tràn RAM executor và còn chậm hơn shuffle bình thường.
AQE & vài config hữu ích
Để Spark tự tối ưu khi chạy

Adaptive Query Execution (AQE) — có từ Spark 3 — cho phép Spark tự tối ưu lúc chạy dựa trên kích thước data thực tế (chứ không chỉ đoán trước). Cụ thể AQE: gộp các partition nhỏ lại sau shuffle (đỡ task vụn), đổi sort-merge join sang broadcast khi phát hiện một bên đủ nhỏ, và xử lý skew bằng cách tự tách task quá lớn.

# bật AQE (Spark 3+; nhiều bản đã bật mặc định)
spark.conf.set("spark.sql.adaptive.enabled", "true")

Vài config hay phải chỉnh:

ConfigÝ nghĩa
spark.sql.shuffle.partitionsSố partition sau shuffle, mặc định 200. Data nhỏ thì 200 task là quá thừa → giảm xuống (vd 8–50) cho nhanh
spark.sql.adaptive.enabledBật/tắt AQE. Nên để true để Spark tự tối ưu lúc chạy
spark.sql.autoBroadcastJoinThresholdNgưỡng tự broadcast (mặc định ~10MB). Tăng nếu muốn Spark broadcast bảng lớn hơn; đặt -1 để tắt
Checklist tối ưu
Những nguyên tắc vàng cần nhớ
← Quay lại
Bài 5: Streaming