Bài 3

Spark SQL

Nếu bạn mạnh SQL thì đây là cách nhanh nhất để dùng Spark. DataFrame API và SQL thực ra chạy trên cùng một engine (Catalyst optimizer), cho ra cùng kết quả và cùng hiệu năng — chọn cái nào hoàn toàn tùy thói quen của bạn.

Tạo View để query SQL
Đăng ký DataFrame thành bảng tạm

Spark không tự biết DataFrame của bạn tên là gì trong câu SQL. Muốn dùng spark.sql(...), trước hết phải đăng ký DataFrame thành một view — coi như đặt cho nó một cái tên bảng để câu lệnh SQL tham chiếu tới.

# đăng ký DataFrame thành view tạm
sales.createOrReplaceTempView("sales")
customers.createOrReplaceTempView("customers")

# giờ có thể viết SQL thuần
spark.sql("""
  SELECT category, SUM(qty*unit_price) AS revenue
  FROM sales
  WHERE qty > 0
  GROUP BY category
  ORDER BY revenue DESC
""").show()

Vài điểm quan trọng về view:

Mẹo: "Replace" trong tên hàm nghĩa là chạy lại nhiều lần không lỗi — view cũ bị ghi đè. Rất tiện khi thử nghiệm trong notebook.
DataFrame API vs SQL — tương đương
Hai cách viết, một engine

Mọi thao tác trong DataFrame API đều có câu SQL tương ứng. Cả hai được Catalyst dịch về cùng một kế hoạch thực thi, nên đừng lo lắng về hiệu năng — chúng giống hệt nhau.

DataFrame APISQL tương ứng
df.filter(df.qty > 0)WHERE qty > 0
df.groupBy("category").agg(sum("amount"))GROUP BY category + SUM(amount)
a.join(b, "customer_id")FROM a JOIN b ON ...
df.orderBy(df.revenue.desc())ORDER BY revenue DESC
df.withColumn("amount", col("qty")*col("unit_price"))SELECT qty*unit_price AS amount
df.select("region").distinct()SELECT DISTINCT region
Nhớ: chọn API hay SQL không ảnh hưởng tốc độ — cùng một Catalyst optimizer, cùng một kế hoạch vật lý, cùng một performance.
JOIN trong SQL
Ghép sales với customers

JOIN trong Spark SQL viết y như SQL truyền thống. Ví dụ tính doanh thu theo vùng bằng cách ghép bảng sales với customers:

spark.sql("""
  SELECT c.region, SUM(s.qty*s.unit_price) AS revenue
  FROM sales s
  JOIN customers c ON s.customer_id = c.customer_id
  GROUP BY c.region
  ORDER BY revenue DESC
""").show()

Các kiểu JOIN bạn sẽ gặp:

Quen thuộc: đặt alias bảng (s, c) giúp câu lệnh ngắn gọn và tránh nhập nhằng khi hai bảng có cột trùng tên.
CTE & Subquery
Chia nhỏ query phức tạp

Khi logic dài, viết tất cả vào một câu SELECT lồng nhau rất khó đọc. CTE (Common Table Expression) — mệnh đề WITH — cho phép đặt tên cho từng bước trung gian, rồi dùng lại như một bảng:

spark.sql("""
  WITH cat_rev AS (
    SELECT category, SUM(qty*unit_price) AS rev
    FROM sales
    WHERE qty > 0
    GROUP BY category
  )
  SELECT *
  FROM cat_rev
  WHERE rev > 1000000
  ORDER BY rev DESC
""").show()

Ở đây cat_rev gom doanh thu theo nhóm hàng trước, rồi câu SELECT chính chỉ việc lọc những nhóm vượt mốc 1 triệu. CTE giúp query dễ đọc, tách bạch từng tầng logic, và bạn có thể nối nhiều CTE liên tiếp bằng dấu phẩy.

Subquery: bạn cũng có thể lồng một SELECT ngay trong FROM hoặc WHERE. CTE thường dễ đọc hơn subquery lồng sâu khi logic phức tạp.
Window function trong SQL
Xếp hạng, running total

Window function tính toán trên một "cửa sổ" các dòng mà không gom chúng lại như GROUP BY. Ví dụ xếp hạng đơn hàng theo doanh thu trong từng nhóm sản phẩm:

spark.sql("""
  SELECT order_id, category, amount,
         ROW_NUMBER() OVER (
           PARTITION BY category
           ORDER BY amount DESC
         ) AS rn
  FROM (
    SELECT *, qty*unit_price AS amount
    FROM sales
    WHERE qty > 0
  )
""").show()

Cấu trúc cốt lõi là OVER (PARTITION BY ... ORDER BY ...):

Ứng dụng kinh điển: lấy top-N mỗi nhóm — bọc query trên thêm một lớp WHERE rn <= 3.
Khi nào dùng SQL vs DataFrame
Không có đáp án sai

Vì hai cách cho cùng hiệu năng, lựa chọn chỉ là về sự thuận tiện và bảo trì:

Mẹo gỡ rối: spark.sql("...").explain() vẫn xem được kế hoạch thực thi y như khi gọi trên DataFrame — tiện để hiểu Spark tối ưu câu SQL của bạn ra sao.
Bài tập
Luyện viết SQL trên dữ liệu lab

Đăng ký view cho salescustomers, rồi giải 3 bài sau bằng spark.sql(...):

Tự kiểm: mỗi bài hãy thử viết lại bằng DataFrame API và so kết quả — chúng phải khớp nhau từng dòng.
← Quay lại
Bài 2: Transform