PySpark 10 — Tối ưu & Gỡ lỗi từ góc PySpark
Vì sao bài này viết từ "góc PySpark"
Cơ chế engine — shuffle vận hành thế nào, cost-based optimizer, memory model của executor — đã được mổ xẻ ở tuning Spark và shuffle & partitions. Bài này không lặp lại lý thuyết đó. Nó trả lời một câu hỏi hẹp hơn nhưng thực chiến hơn: khi bạn ngồi trước một job PySpark chạy chậm, bạn gõ gì, đọc gì, sửa gì?
Vòng lặp tối ưu luôn giống nhau: đo → tìm nút thắt → sửa → đo lại. Sai lầm phổ biến nhất là bỏ qua bước "đo" và nhảy thẳng vào "sửa" — tăng bừa spark.sql.shuffle.partitions, rải .cache() khắp nơi, hoặc thêm executor. Đó là đoán mò. Hai công cụ để đo: df.explain() (đọc kế hoạch trước khi chạy) và Spark UI (đọc thực tế sau khi chạy).
Đọc physical plan bằng df.explain(mode="formatted")
Spark biến DataFrame của bạn thành một physical plan — cây các toán tử thực thi. Đọc được cây này là kỹ năng nền tảng. Gọi:
df.explain(mode="formatted")
Ngoài "formatted" (dễ đọc nhất, có đánh số node + phần chi tiết bên dưới) còn có "simple", "extended", "cost" (kèm ước lượng thống kê) và "codegen". Đọc plan từ dưới lên: node dưới cùng chạy trước.
Những gì cần soi trong plan:
Exchange— đây chính là shuffle. MỗiExchangelà một lần dữ liệu bị băm lại và chuyển qua mạng giữa các executor. Đây là toán tử đắt nhất. Đếm sốExchange; mỗi cái là một điểm nghi vấn "có thực sự cần không?".- Kiểu join.
BroadcastHashJoinnghĩa là Spark gửi (broadcast) bảng nhỏ tới mọi executor — rất nhanh, không shuffle bảng lớn.SortMergeJoinnghĩa là cả hai phía đều bị shuffle rồi sort — đắt. Nếu bạn biết một phía nhỏ mà plan vẫn raSortMergeJoin, đó là dấu hiệu broadcast threshold quá thấp hoặc thống kê sai. Scanvà pushdown. Ở node scan, đọcPushedFiltersvàReadSchema. Nếu filter được đẩy xuống (PushedFilters: [GreaterThan(amount,0)]) thì storage lọc hộ bạn — tốt. Nếu bộ lọc không xuất hiện ở đây mà nằm ở mộtFilterphía trên scan, Spark đang đọc thừa dữ liệu.ReadSchemacho biết có column pruning không: nếu bạn chỉ cần 3 cột mà schema đọc ra 40 cột, bạn đang lãng phí I/O.AQEShuffleRead,AdaptiveSparkPlan— xuất hiện khi AQE bật (xem bên dưới).
Nhận biết plan xấu qua vài mẫu quen thuộc: chuỗi nhiều Exchange liên tiếp (shuffle chồng shuffle); SortMergeJoin với một phía rõ ràng nhỏ; scan không có PushedFilters dù code có .filter(); hoặc CartesianProduct (join thiếu điều kiện — gần như luôn là bug).
Đọc Spark UI: nơi sự thật lộ ra
explain() cho biết kế hoạch; Spark UI cho biết chuyện đã thực sự xảy ra. UI truy cập qua cổng 4040 (application đang chạy) hoặc History Server. Bốn tab quan trọng:
Tab SQL — điểm khởi đầu tốt nhất. Mỗi query hiển thị một sơ đồ toán tử kèm số liệu ngay trên từng node: số hàng, dung lượng shuffle read/write, thời gian. Đây là cách nhanh nhất để thấy toán tử nào "ăn" thời gian và dữ liệu.
Tab Jobs → Stages — một job gồm nhiều stage; ranh giới stage chính là chỗ shuffle. Với mỗi stage, mở phần Summary Metrics for Tasks: bảng phân vị (Min, 25th, Median, 75th, Max) của Duration, Shuffle Read Size, Spill. Dấu hiệu bệnh:
- Task lệch (skew): cột
Maxcủa Duration hoặc Shuffle Read lớn hơnMedianhàng chục lần. Ví dụ median 3 giây nhưng max 6 phút → một task ôm phần lớn dữ liệu, cả stage phải chờ nó. Đây là biểu hiện kinh điển của data skew. - Spill:
Spill (Memory)vàSpill (Disk)khác 0 nghĩa là dữ liệu tràn khỏi RAM executor xuống đĩa — chậm đột biến. Do partition quá to hoặc bộ nhớ executor quá nhỏ. - Shuffle lớn: tổng Shuffle Read/Write hàng chục–hàng trăm GB cho một phép tính lẽ ra nhỏ → nghi shuffle thừa.
- GC time: cột GC Time chiếm tỉ lệ lớn so với Duration → executor thiếu bộ nhớ, dành thời gian dọn rác thay vì tính.
Tab Stages → DAG Visualization — sơ đồ khối cho thấy dòng chảy toán tử và chỗ nào tạo shuffle boundary. Hữu ích để hình dung một stage dài bất thường bắt nguồn từ đâu.
Nguyên tắc đọc UI: đừng nhìn tổng thời gian job, hãy tìm stage chiếm phần lớn thời gian, rồi trong stage đó tìm chênh lệch Max/Median và spill. Nút thắt gần như luôn nằm ở một chỗ rất cụ thể.
Bản đồ nguyên nhân chậm → cách sửa
| Nguyên nhân | Dấu hiệu trên UI/plan | Cách sửa từ PySpark |
|---|---|---|
| Data skew | 1 task Max ≫ Median, treo cuối stage | Salting khoá; bật AQE skew join |
| Shuffle thừa | Nhiều Exchange, shuffle GB lớn | Bớt groupBy/join không cần; broadcast bảng nhỏ; repartition đúng khoá |
| Small files / partition sai | Rất nhiều task siêu ngắn; scan chậm | coalesce, đọc đúng partition, nén file nhỏ |
| Python UDF | Stage Python chậm, ít song song | Thay bằng hàm built-in / pandas_udf |
| collect / toPandas | Driver OOM, một mình driver bận | Ghi ra file; chỉ gom kết quả nhỏ |
| Cache sai chỗ | Storage tab đầy, tính lại nhiều lần | cache() DF dùng lại nhiều lần; unpersist khi xong |
Data skew là thủ phạm số một trong dữ liệu ngân hàng: một chi nhánh lớn, một mã sản phẩm phổ biến, hay null chiếm phần lớn khoá join. Hai cách xử lý (chi tiết ở phần ví dụ): salting và AQE skew join.
Shuffle thừa thường do join/groupBy không cần thiết hoặc repartition sai. Mẹo: dùng broadcast() thủ công cho bảng tra cứu nhỏ (danh mục chi nhánh, mã sản phẩm) để tránh shuffle bảng giao dịch lớn.
Small files và partition sai — xem đọc/ghi nguồn dữ liệu. Ngàn file vài KB làm scheduler nghẹt vì mỗi file thành một task; ngược lại partition khổng lồ gây spill.
Python UDF giết hiệu năng vì phá vỡ vectorization và phải serialize dữ liệu qua ranh giới JVM↔Python — xem UDF & pandas_udf. Luôn ưu tiên hàm trong pyspark.sql.functions.
collect() / toPandas() kéo toàn bộ dữ liệu về driver — một máy đơn. Với bảng lớn, driver hết RAM (OOM) ngay. Chỉ dùng khi kết quả đã nhỏ (sau agg, limit).
Adaptive Query Execution (AQE)
AQE để Spark điều chỉnh kế hoạch trong lúc chạy dựa trên thống kê runtime thật, thay vì khoá cứng plan từ đầu. Từ Spark 3.2 nó bật mặc định. Ba lợi ích lớn:
- Coalesce shuffle partitions: sau shuffle, AQE tự gộp các partition nhỏ lại, tránh cảnh đặt
shuffle.partitions=200rồi phần lớn partition rỗng. Bạn bớt phải chỉnh tay con số này. - Skew join: AQE phát hiện partition to bất thường sau shuffle và tự chẻ nó thành nhiều mảnh chạy song song — xử lý skew join mà không cần salting thủ công trong nhiều trường hợp.
- Chuyển sang broadcast: nếu tới lúc chạy mới biết một phía join thật ra nhỏ, AQE đổi
SortMergeJointhànhBroadcastHashJoin.
Cấu hình liên quan (tên đúng):
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.coalescePartitions.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
Khi AQE hoạt động, plan bọc trong AdaptiveSparkPlan và bạn sẽ thấy AQEShuffleRead trên UI.
cache / persist / checkpoint đúng cách
Chỉ cache khi một DataFrame được dùng lại nhiều lần. Nếu DF chỉ đọc một lần, cache() chỉ tốn RAM vô ích. cache() là persist() với mức mặc định MEMORY_AND_DISK. Chọn StorageLevel khi cần: MEMORY_ONLY (nhanh nhất, mất nếu thiếu RAM), MEMORY_AND_DISK (an toàn, tràn xuống đĩa), MEMORY_AND_DISK_SER (nén serialize, tiết kiệm RAM đổi lấy CPU).
from pyspark import StorageLevel
df_ref = df_ref.persist(StorageLevel.MEMORY_AND_DISK)
df_ref.count() # materialize cache (action)
# ... dùng df_ref nhiều lần ...
df_ref.unpersist() # giải phóng khi xong
Lưu ý: cache là lazy — chưa thực sự lưu cho tới khi có action đầu tiên. Và luôn unpersist() khi không cần nữa, nếu không bộ nhớ đầy dần rồi tràn.
Checkpoint khác cache: nó ghi DataFrame ra đĩa và cắt đứt lineage (chuỗi phụ thuộc). Với pipeline lặp (vòng lặp ML, join dây chuyền hàng chục bước), lineage dài đến mức việc lập kế hoạch trở nên chậm hoặc gây StackOverflow. df.checkpoint() "chốt hạ" dữ liệu thành nguồn mới, lineage reset về 0. Cần setCheckpointDir(path) trước.
Cấu hình hiệu năng cốt lõi
| Config | Nghĩa & ảnh hưởng |
|---|---|
spark.sql.shuffle.partitions | Số partition sau shuffle (join/groupBy). Mặc định 200. Quá cao → nhiều task rỗng; quá thấp → partition to, spill. AQE giảm nhu cầu chỉnh tay. |
spark.sql.autoBroadcastJoinThreshold | Ngưỡng dung lượng bảng để tự broadcast (mặc định 10MB). Tăng để broadcast bảng tra cứu lớn hơn; đặt -1 để tắt hẳn broadcast. |
spark.executor.memory | RAM heap mỗi executor. Thiếu → spill, GC nhiều, OOM. |
spark.executor.memoryOverhead | Bộ nhớ off-heap (Python worker, buffer shuffle). PySpark nặng UDF thường cần tăng cái này. |
spark.executor.cores | Số task chạy song song mỗi executor. |
spark.driver.memory | RAM driver. Phải tăng nếu buộc collect() kết quả lớn (nhưng nên tránh collect). |
Đừng nhớ máy móc các con số — nhớ triệu chứng → nút cần vặn: spill/GC → tăng executor.memory; join lẽ ra broadcast mà không broadcast → tăng autoBroadcastJoinThreshold; driver OOM → xem lại collect/toPandas và driver.memory.
Debug các lỗi hay gặp
OOM — phân biệt driver vs executor. Đọc log kỹ: OOM ở driver (thường do collect/toPandas/broadcast quá lớn) khác OOM ở executor (partition quá to, skew, cache quá nhiều). Sửa driver OOM bằng cách không gom dữ liệu về driver; sửa executor OOM bằng repartition, xử skew, hoặc tăng bộ nhớ/overhead.
Serialization errors (PicklingError, Task not serializable) — thường do dùng biến/đối tượng không pickle được bên trong UDF (kết nối DB, client, logger toàn cục). Sửa: khởi tạo tài nguyên bên trong hàm chạy trên executor, đừng bắt (capture) từ driver.
Thư viện thiếu trên executor — code chạy tốt trên driver nhưng executor báo ModuleNotFoundError. Vì executor là các máy khác, không tự có thư viện bạn cài ở driver. Cách xử lý ở config & deployment: đóng gói môi trường (--archives/conda-pack) hoặc cài sẵn trên cluster.
Dữ liệu null / sai kiểu — kết quả bất ngờ (tổng ra null, join hụt hàng) thường do khoá join có null, kiểu không khớp (string vs int), hoặc parse ngày sai. Kiểm tra bằng printSchema(), đếm null trên khoá, và ép kiểu tường minh bằng cast().
Ví dụ: một skew join chậm và cách chữa
Truy vấn thực tế: nối bảng giao dịch khổng lồ với danh mục chi nhánh theo branch_id. Một vài chi nhánh trung tâm chiếm phần lớn giao dịch → skew.
from pyspark.sql import functions as F
# --- Bản gốc: SortMergeJoin, bị skew ---
result = (transactions
.join(branches, "branch_id") # cả hai bị shuffle
.groupBy("branch_id")
.agg(F.sum("amount").alias("total")))
result.explain(mode="formatted")
# Plan cho thấy: SortMergeJoin + Exchange trên transactions.
# Spark UI: stage join có 1 task Max = 6 phút, Median = 4 giây → skew rõ.
Cách 1 — broadcast (bảng branches nhỏ, chỉ vài trăm dòng): ép broadcast để không shuffle bảng giao dịch.
result = (transactions
.join(F.broadcast(branches), "branch_id") # BroadcastHashJoin
.groupBy("branch_id")
.agg(F.sum("amount").alias("total")))
# explain() giờ hiện BroadcastHashJoin, không còn Exchange trên transactions.
Broadcast xử lý skew của join vì phía lớn không bị băm lại. Nhưng nếu skew nằm ở bước groupBy (một chi nhánh gom quá nhiều hàng vào một partition), cần salting hoặc AQE.
Cách 2 — AQE skew join (khi cả hai phía lớn, không broadcast được): để Spark tự chẻ partition to.
spark.conf.set("spark.sql.adaptive.enabled", "true")
spark.conf.set("spark.sql.adaptive.skewJoin.enabled", "true")
# Chạy lại truy vấn gốc: AQE phát hiện partition khổng lồ và chẻ nhỏ.
# UI: các task cân bằng hơn, không còn 1 task treo.
Cách 3 — salting thủ công (kiểm soát tối đa, hoặc khi skew ở aggregation): rải khoá nóng bằng hậu tố ngẫu nhiên, gom hai pha.
N = 16 # số salt
tx_salted = transactions.withColumn("salt", (F.rand() * N).cast("int"))
# nhân bản bảng nhỏ cho mọi salt để join khớp
from pyspark.sql import functions as F
salts = spark.range(N).withColumnRenamed("id", "salt")
br_salted = branches.crossJoin(salts)
partial = (tx_salted
.join(br_salted, ["branch_id", "salt"])
.groupBy("branch_id", "salt")
.agg(F.sum("amount").alias("part"))) # pha 1: rải đều
result = (partial
.groupBy("branch_id")
.agg(F.sum("part").alias("total"))) # pha 2: gộp lại
Salting biến một khoá nóng thành N khoá con chạy song song, rồi cộng dồn. Đánh đổi: code phức tạp hơn và thêm một shuffle nhỏ ở pha 2 — nhưng loại bỏ được task treo. Thứ tự ưu tiên thực dụng: broadcast (nếu một phía nhỏ) → AQE skew join (mặc định, ít công) → salting (khi hai cách trên chưa đủ).
(Toàn bộ đoạn trên là minh hoạ bằng PySpark/Python, không phải SQL sandbox.)
Use case thực tế
Bối cảnh — NCB, job tổng hợp giao dịch cuối ngày. Job gom giao dịch theo chi nhánh để dựng báo cáo doanh số. Gần đây nó kéo dài từ ~25 phút lên hơn 2 giờ, thỉnh thoảng OOM executor. Không ai đổi code — chỉ là một chi nhánh hội sở mới mở kênh online, lượng giao dịch qua khoá branch_id của nó phình gấp nhiều lần phần còn lại.
Đo. Mở Spark UI, tab Stages: stage join+groupBy chiếm gần như toàn bộ thời gian. Bảng Summary Metrics: Duration Median ≈ 5 giây, Max ≈ 1 giờ 50 phút — một task duy nhất treo, kèm Spill (Disk) hàng chục GB ở đúng task đó. explain(mode="formatted") xác nhận SortMergeJoin shuffle cả hai phía. Kết luận chắc chắn: data skew ở khoá chi nhánh, không phải thiếu tài nguyên.
Sửa. Hai thay đổi, không đụng logic nghiệp vụ:
- Bật AQE skew join (
spark.sql.adaptive.enabled+skewJoin.enabled) để Spark tự chẻ partition khổng lồ. - Với bước
groupByvẫn còn nóng ở khoá hội sở, thêm saltingN=16cho khoá đó theo mẫu hai pha ở trên; đồng thờibroadcast()bảng danh mục chi nhánh (nhỏ) để bỏ shuffle bảng giao dịch.
Kết quả (ước lượng). Trên cùng cluster và cùng dữ liệu: thời gian giảm từ ~2 giờ xuống ~18 phút, task Max/Median về gần nhau (không còn task treo), spill biến mất, executor hết OOM. Không thêm một node nào — chỉ đọc đúng chỗ nghẽn và sửa đúng nguyên nhân. Bài học vận hành: khi job "tự nhiên chậm đi", nghi ngờ phân bố dữ liệu thay đổi (skew) trước khi nghĩ tới việc cấp thêm máy.
Ghi nhớ
- Vòng lặp bất biến: đo (
explain+ Spark UI) → tìm nút thắt → sửa → đo lại. Đừng đoán, đừng rải.cache()bừa. df.explain(mode="formatted"): đếmExchange(shuffle), soiBroadcastHashJoinvsSortMergeJoin, kiểmPushedFiltersvàReadSchemaở node scan.- Spark UI: bắt đầu ở tab SQL/Stages, tìm stage chiếm phần lớn thời gian, xem Max vs Median (skew), Spill (thiếu RAM), GC.
- Data skew là thủ phạm số một trong dữ liệu ngân hàng; chữa theo thứ tự broadcast → AQE skew join → salting.
- AQE (mặc định từ Spark 3.2): coalesce partition, skew join, chuyển broadcast — giảm nhu cầu chỉnh
shuffle.partitionstay. - cache/persist chỉ khi tái sử dụng nhiều lần, chọn
StorageLevel, và luônunpersist; checkpoint để cắt lineage dài. - Tránh Python UDF (dùng hàm built-in /
pandas_udf) và tránhcollect/toPandastrên bảng lớn (driver OOM). - Debug: phân biệt OOM driver vs executor; serialization do capture đối tượng không pickle được; thư viện thiếu trên executor; null/sai kiểu ở khoá join.
- Config cốt lõi cần thuộc:
spark.sql.shuffle.partitions,spark.sql.autoBroadcastJoinThreshold,spark.executor.memory/memoryOverhead— nhớ theo triệu chứng → nút cần vặn.
Nền tảng ở tổng quan PySpark và DataFrame API; triển khai môi trường executor ở config & deployment.
Nguồn tham khảo
- Apache Spark — Performance Tuning (SQL) — AQE, partitioning, broadcast, caching
- Apache Spark — Tuning Spark (memory, serialization, GC, data locality)
- Apache Spark — Monitoring and Instrumentation (Web UI: Jobs, Stages, SQL tabs)
- Apache Spark — RDD Programming Guide: RDD Persistence & Storage Levels
- Apache Spark — Configuration (spark.sql.shuffle.partitions, executor.memory, autoBroadcastJoinThreshold…)
- Bill Chambers & Matei Zaharia, Spark: The Definitive Guide (O'Reilly, 2018) — phần Performance Tuning & Debugging
Bài viết liên quan
So sánh các định dạng dữ liệu (CSV, JSON, XML, Avro, Parquet, ORC) và lý do lưu theo cột nhanh hơn cho phân tích. Bài đi sâu vào row vs columnar storage, nén (Snappy/gzip/zstd), schema evolution, OLTP vs OLAP, object storage và partitioning để tối ưu chi phí lẫn tốc độ truy vấn.
Data Engineering là ngành xây dựng và vận hành hệ thống biến dữ liệu thô thành dữ liệu sạch, tin cậy, sẵn sàng cho phân tích và AI. Bài giới thiệu vai trò Data Engineer trong vòng đời dữ liệu (nguồn → ingestion → storage → transformation → serving), phân biệt với Analyst/Scientist/ML Engineer, bức tranh hệ sinh thái công cụ và bài toán đưa dữ liệu core banking sang kho phân tích.
Stream processing là gì, khác biệt batch vs stream (bounded/unbounded), micro-batch (Spark) vs true streaming, Apache Flink là gì và định vị so với Spark Structured Streaming và Kafka Streams. Kiến trúc runtime JobManager/TaskManager, các tầng API, triết lý streaming-first và bối cảnh phát hiện gian lận ngân hàng.
Vì sao một máy không đủ và cần xử lý phân tán: từ MapReduce, Hadoop/HDFS đến Apache Spark in-memory. Kiến trúc driver–executor–cluster manager, các mức trừu tượng RDD/DataFrame/Dataset, cơ chế lazy evaluation với DAG, và vì sao shuffle (wide dependency) là phần tốn kém nhất. Kèm PySpark, Spark SQL và các kỹ thuật tối ưu (partition, broadcast join, cache, chống skew) cùng khi nào KHÔNG nên dùng Spark.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!