Databricks for Data & AI Engineers
Trang chủ / Module 4
⚡ Module 4 / 14 · 2-3 tuần

Spark & PySpark trên Databricks

DataFrame API, Spark SQL, kiến trúc thực thi (partition/shuffle/Catalyst/AQE), Photon, tối ưu hiệu năng.

SparkPySparkPerformance
📌 Vì sao module này quan trọng

Databricks được sáng lập bởi chính đội ngũ tạo ra Apache Spark — hiểu Spark hoạt động ra sao "bên dưới" là khác biệt giữa việc viết code chạy được và viết code chạy nhanh, ổn định ở quy mô dữ liệu lớn. Đây là kỹ năng kỹ thuật cốt lõi nhất của một Data Engineer trên nền tảng này.

🎯 Mục tiêu học tập

  • Sử dụng thành thạo DataFrame API và Spark SQL cho các thao tác biến đổi dữ liệu phổ biến
  • Hiểu kiến trúc thực thi: Driver/Executor, partition, shuffle, lazy evaluation
  • Hiểu Catalyst optimizer và Adaptive Query Execution (AQE) tối ưu truy vấn tự động ra sao
  • Áp dụng các kỹ thuật tối ưu hiệu năng thực tế: broadcast join, tránh small files, caching đúng chỗ

DataFrame API & Spark SQL

from pyspark.sql import functions as F

orders = spark.table("bronze.raw_orders")

silver_orders = (orders
    .filter(F.col("status").isNotNull())
    .withColumn("order_date", F.to_date("created_at"))
    .groupBy("customer_id", "order_date")
    .agg(
        F.sum("amount").alias("total_amount"),
        F.count("*").alias("order_count"),
    )
)
silver_orders.write.mode("overwrite").saveAsTable("silver.daily_customer_orders")

# Tương đương bằng Spark SQL
spark.sql(
    "SELECT customer_id, DATE(created_at) AS order_date, "
    "SUM(amount) AS total_amount, COUNT(*) AS order_count "
    "FROM bronze.raw_orders WHERE status IS NOT NULL "
    "GROUP BY customer_id, DATE(created_at)"
)

DataFrame API và Spark SQL sinh ra cùng một execution plan — chọn cách viết nào tự nhiên hơn với bạn (nhiều Data Engineer thích SQL cho biến đổi đơn giản, DataFrame API cho logic phức tạp cần tái sử dụng qua hàm Python).

Kiến trúc thực thi: Driver, Executor, Partition, Shuffle

Driver là process điều phối, giữ execution plan và phân phối task; Executor là các process thực sự chạy task trên từng node. Dữ liệu được chia thành các partition (đơn vị song song hoá nhỏ nhất) — số lượng partition ảnh hưởng trực tiếp đến mức độ song song. Các thao tác cần dữ liệu từ nhiều partition khác nhau gộp lại (groupBy, join, orderBy) gây ra shuffle — di chuyển dữ liệu qua mạng giữa các executor, thường là bước tốn kém nhất trong một Spark job.

Spark hoạt động theo mô hình lazy evaluation: các phép biến đổi (filter, select, withColumn...) chỉ xây dựng execution plan, không chạy thật cho đến khi gặp một action (write, collect, count...) — hiểu điều này giải thích vì sao đo thời gian một dòng .filter() riêng lẻ luôn cho kết quả gần 0.

Catalyst Optimizer & Adaptive Query Execution (AQE)

Catalyst là bộ tối ưu truy vấn tự động phân tích execution plan logic, áp dụng các quy tắc tối ưu (predicate pushdown — lọc dữ liệu sớm nhất có thể; column pruning — chỉ đọc cột cần dùng) trước khi sinh ra physical plan thật sự thực thi. AQE (bật mặc định trên Databricks) đi xa hơn: điều chỉnh lại plan trong lúc job đang chạy dựa trên số liệu thống kê thực tế (tự động gộp partition nhỏ sau shuffle, tự động chuyển sang broadcast join nếu phát hiện 1 bảng nhỏ hơn dự kiến) — giảm đáng kể nhu cầu tinh chỉnh thủ công so với Spark thế hệ cũ.

Tối ưu hiệu năng thực tế

  • Broadcast join: khi join một bảng lớn với một bảng nhỏ (vừa đủ trong bộ nhớ, thường <10MB-100MB), ép Spark gửi toàn bộ bảng nhỏ đến mọi executor thay vì shuffle cả hai bảng — nhanh hơn rất nhiều. AQE thường tự phát hiện, nhưng có thể ép rõ bằng F.broadcast(small_df).
  • Tránh small files: ghi quá nhiều file nhỏ (thường do streaming ghi liên tục, hoặc partition quá nhỏ) làm chậm cả ghi lẫn đọc sau này — dùng OPTIMIZE (Module 3) định kỳ, hoặc repartition()/coalesce() trước khi ghi.
  • Cache đúng chỗ: .cache()/.persist() chỉ có lợi khi cùng một DataFrame được dùng lại nhiều lần trong cùng phiên làm việc — cache một DataFrame chỉ dùng một lần là lãng phí bộ nhớ, không tăng tốc gì.
  • Đọc EXPLAIN plan: dùng df.explain(True) để xem physical plan thật, phát hiện shuffle không cần thiết hoặc full table scan khi lẽ ra có thể partition pruning.

🏋️ Bài tập thực hành

Tối ưu một truy vấn chậm

Lấy 2 bảng mẫu (1 bảng lớn ~1 triệu dòng, 1 bảng nhỏ ~1000 dòng dạng lookup), viết join giữa chúng theo cách "ngây thơ" (không tối ưu), đo thời gian chạy và xem execution plan bằng .explain(). Sau đó ép broadcast join, so sánh lại thời gian và execution plan. Viết 1 pipeline nạp dữ liệu theo từng batch nhỏ (mô phỏng streaming), sau đó chạy OPTIMIZE và đo cải thiện tốc độ truy vấn.

📚 Tài nguyên học tập

✅ Tự đánh giá hoàn thành

  • Viết được biến đổi dữ liệu bằng cả DataFrame API và Spark SQL
  • Giải thích được shuffle là gì và vì sao nó tốn kém
  • Đọc hiểu được execution plan từ .explain()
  • Áp dụng được broadcast join và giải thích khi nào nên dùng