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.
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ặcrepartition()/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
EXPLAINplan: dùngdf.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
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
-
Apache Spark — Programming GuideTài liệu chính thức Spark SQL/DataFrameDocs
-
Databricks — Performance tuning guideHướng dẫn tối ưu hiệu năng chính thức trên DatabricksDocs
-
Databricks — Adaptive Query ExecutionGiải thích AQE hoạt động ra saoDocs
✅ 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