Hướng dẫn thực tế về Window Functions trong PySpark

02 tháng 9, 2026·8 phút đọc

Window Functions trong PySpark cho phép tính toán trên các nhóm dữ liệu liên quan mà không làm mất đi thông tin chi tiết của từng bản ghi, vượt qua giới hạn của groupBy(). Bài viết giới thiệu cách sử dụng partitionBy, orderBy và các hàm xếp hạng cùng ví dụ thực tế về phân tích bán hàng, giúp lập trình viên xử lý dữ liệu lớn hiệu quả hơn.

Hướng dẫn thực tế về Window Functions trong PySpark

Khám phá Window Functions trong PySpark: Vũ khí xử lý dữ liệu mà groupBy() thiếu sót

Khi cần tổng hợp dữ liệu, hàm groupBy() tiêu chuẩn của PySpark thường là lựa chọn đầu tiên, nhưng nó có một hạn chế cơ bản: chỉ trả về một dòng kết quả cho mỗi nhóm dữ liệu. Window Functions ra đời để giải quyết vấn đề này, cho phép bạn giữ nguyên từng bản ghi gốc đồng thời bổ sung các phép tính tổng hợp trên toàn nhóm, mở ra khả năng phân tích linh hoạt như xếp hạng, tổng lũy tiến và so sánh giữa các bản ghi.

Vì sao Window Functions vượt trội hơn groupBy()?

Hành vi mặc định của groupBy() chỉ phù hợp khi bạn muốn một kết quả duy nhất cho mỗi nhóm. Ví dụ, nếu bạn tính tổng doanh số của một cửa hàng với 1.000 giao dịch, bạn sẽ nhận về đúng một con số. Nhưng trong nhiều tình huống thực tế — như phân tích nhật ký sự kiện, dữ liệu tài chính, hay hành vi khách hàng — bạn cần biết cả tổng nhóm và chi tiết từng dòng.

Window Functions hoạt động bằng cách định nghĩa một “cửa sổ” các bản ghi liên quan, sau đó tính toán giá trị dựa trên cửa sổ đó mà không nén dữ liệu thành một dòng duy nhất. Đây là công cụ lý tưởng cho các tác vụ như xếp hạng trong nhóm, tính tổng lũy tiến, so sánh giá trị trước/sau, và xác định tỷ trọng của từng phần tử trong tổng thể.

Cài đặt và khởi tạo PySpark

Để bắt đầu, bạn cần cài đặt PySpark trong môi trường ảo. Nếu dùng uv, các lệnh sau sẽ giúp bạn thiết lập nhanh chóng:

mkdir pyspark-windows
cd pyspark-windows
uv init
uv venv
.venv\Scripts\activate
uv pip install pyspark

Tiếp theo, tạo một Spark session chạy cục bộ — không cần cluster phức tạp:

import os
import sys
os.environ["PYSPARK_PYTHON"] = sys.executable
os.environ["PYSPARK_DRIVER_PYTHON"] = sys.executable

from pyspark.sql import SparkSession

spark = (SparkSession.builder
    .master("local[*]")
    .appName("sales-analysis")
    .config("spark.pyspark.python", sys.executable)
    .config("spark.pyspark.driver.python", sys.executable)
    .getOrCreate())

Cấu hình local[*] cho phép Spark chạy trên máy tính cá nhân và tận dụng toàn bộ nhân CPU.

Xây dựng bộ dữ liệu mẫu

Bài viết sử dụng dữ liệu bán hàng gồm 12 giao dịch từ ba cửa hàng tại London, Manchester và Bristol trong tháng 1 năm 2026:

from datetime import date
from pyspark.sql import functions as F
from pyspark.sql import types as T
from pyspark.sql.window import Window

sales_data = [
    (1, "London_Store", date(2026, 1, 2), "Laptop", 1200.00),
    (2, "London_Store", date(2026, 1, 3), "Monitor", 350.00),
    (3, "London_Store", date(2026, 1, 5), "Keyboard", 90.00),
    (4, "London_Store", date(2026, 1, 8), "Laptop", 1350.00),
    (5, "Manchester_Store", date(2026, 1, 2), "Monitor", 320.00),
    (6, "Manchester_Store", date(2026, 1, 4), "Laptop", 1100.00),
    (7, "Manchester_Store", date(2026, 1, 6), "Mouse", 45.00),
    (8, "Manchester_Store", date(2026, 1, 9), "Laptop", 1250.00),
    (9, "Bristol_Store", date(2026, 1, 3), "Keyboard", 85.00),
    (10, "Bristol_Store", date(2026, 1, 4), "Monitor", 300.00),
    (11, "Bristol_Store", date(2026, 1, 7), "Laptop", 1050.00),
    (12, "Bristol_Store", date(2026, 1, 10), "Monitor", 330.00),
]

Sau khi định nghĩa schema phù hợp, bạn có thể tạo DataFrame và xem dữ liệu được sắp xếp theo cửa hàng và ngày.

Hiểu về “cửa sổ” (Window) trong PySpark

Một window specification xác định tập hợp các dòng mà PySpark sẽ xem xét khi tính giá trị cho dòng hiện tại. Có ba thành phần chính:

  • partitionBy() — Chia dữ liệu thành các nhóm độc lập.
  • orderBy() — Xác định thứ tự các dòng bên trong mỗi nhóm.
  • rowsBetween() hoặc rangeBetween() — Định nghĩa khung cửa sổ tương đối với dòng hiện tại.

Ví dụ, muốn tính tổng doanh số của từng cửa hàng mà vẫn giữ từng giao dịch, ta dùng:

store_window = Window.partitionBy("store")

sales_with_store_total = sales.withColumn(
    "store_total",
    F.sum("amount").over(store_window),
)

Kết quả trả về vẫn đủ 12 dòng, nhưng mỗi dòng đều có thêm cột store_total với tổng doanh số của cửa hàng tương ứng — điều mà groupBy() không thể làm được.

Xếp hạng dòng trong từng nhóm

Một trong những ứng dụng phổ biến nhất là xếp hạng các giao dịch theo giá trị giảm dần trong mỗi cửa hàng:

sales_rank_window = (Window.partitionBy("store")
    .orderBy(F.col("amount").desc()))

ranked_sales = sales.withColumn(
    "sale_rank",
    F.row_number().over(sales_rank_window),
)

Lúc này, giao dịch lớn nhất ở mỗi cửa hàng được gán rank 1, kế tiếp là rank 2, v.v.

Phân biệt row_number(), rank()dense_rank()

Ba hàm xếp hạng này khác nhau khi có giá trị trùng nhau:

  • row_number() — Luôn gán số thứ tự duy nhất, không xử lý tình huống hòa.
  • rank() — Các dòng bằng điểm nhận cùng hạng, nhưng để lại khoảng trống phía sau.
  • dense_rank() — Giống rank() nhưng không tạo khoảng trống.

Ví dụ với các giá trị 100, 100, 80:

+--------+------------+------+------------+
| amount | row_number | rank | dense_rank |
+--------+------------+------+------------+
| 100    | 1          | 1    | 1          |
| 100    | 2          | 1    | 1          |
| 80     | 3          | 3    | 2          |
+--------+------------+------+------------+

Hãy chọn row_number() khi bạn cần chính xác một dòng cho mỗi thứ hạng, và dùng rank()/dense_rank() khi các giá trị bằng nhau cần được đối xử công bằng.

Chọn lọc bản ghi hàng đầu từ mỗi nhóm

Kết hợp xếp hạng với bộ lọc, bạn có thể dễ dàng lấy hai giao dịch lớn nhất từ mỗi cửa hàng:

top_two_sales_per_store = (sales
    .withColumn("sale_rank", F.row_number().over(sales_rank_window))
    .filter(F.col("sale_rank") <= 2))

Mỗi cửa hàng sẽ chỉ còn lại hai dòng có giá trị cao nhất.

Tính tổng lũy tiến và so sánh liên tục

Window Functions cũng hữu ích để tính tổng tích lũy theo thời gian. Ví dụ, muốn biết tổng doanh số của mỗi cửa hàng tính đến từng ngày:

store_date_window = (Window.partitionBy("store")
    .orderBy("sale_date", "transaction_id"))

running_total_window = (store_date_window.rowsBetween(
    Window.unboundedPreceding,
    Window.currentRow,
))

sales_with_running_total = sales.withColumn(
    "running_store_total",
    F.sum("amount").over(running_total_window),
)

Để so sánh giá trị hiện tại với giao dịch trước đó, dùng hàm lag():

sales_with_previous = sales.withColumn(
    "previous_amount",
    F.lag("amount").over(store_date_window),
)

Kết quả cho phép bạn tính chênh lệch giữa hai lần mua liên tiếp của cùng một cửa hàng.

Tối ưu hiệu suất khi làm việc với dữ liệu lớn

Trong môi trường sản xuất với hàng triệu dòng, Window Functions có thể gây chậm do yêu cầu repartition, sortshuffle dữ liệu. Dưới đây là một số mẹo quan trọng:

  • Lọc dữ liệu sớm: Hãy giảm lượng dữ liệu đầu vào trước khi áp dụng window, ví dụ chỉ lấy các giao dịch từ một khoảng thời gian cụ thể.
  • Chọn đúng cột cần thiết: Nếu chỉ cần vài cột cho phép tính, hãy loại bỏ các cột khác trước khi thực hiện window operation.
  • Cân nhắc phân vùng lệch: Nếu một partition key có số dòng vượt trội so với các key khác, có thể gây mất cân bằng tải. Ví dụ, phân vùng theo quốc gia thường rất không đồng đều.
  • Cache kết quả trung gian: Nếu cùng một DataFrame (đã qua window) được dùng nhiều lần, hãy cache để tránh tính toán lại.

Kết hợp tất cả trong một bức tranh phân tích hoàn chỉnh

Dưới đây là ví dụ tổng hợp, thêm nhiều chỉ số hữu ích cho từng giao dịch:

store_total_window = Window.partitionBy("store")
store_date_window = (Window.partitionBy("store")
    .orderBy("sale_date", "transaction_id"))
running_total_window = store_date_window.rowsBetween(
    Window.unboundedPreceding, Window.currentRow
)
store_rank_window = (Window.partitionBy("store")
    .orderBy(F.col("amount").desc()))

sales_analysis = (sales
    .withColumn("store_total", F.sum("amount").over(store_total_window))
    .withColumn("share_of_store_total", F.col("amount") / F.col("store_total"))
    .withColumn("running_store_total", F.sum("amount").over(running_total_window))
    .withColumn("previous_amount", F.lag("amount").over(store_date_window))
    .withColumn("change_from_previous", F.col("amount") - F.col("previous_amount"))
    .withColumn("sale_rank", F.row_number().over(store_rank_window))
)

Kết quả hiển thị mỗi giao dịch với tổng doanh số cửa hàng, tỷ trọng đóng góp, tổng lũy tiến theo thời gian, giá trị giao dịch trước đó và mức thay đổi — tất cả thông tin chi tiết vẫn được giữ nguyên.

Những lỗi thường gặp cần tránh

  • Quên partitionBy(): Điều này khiến window áp dụng trên toàn bộ dataset, không đúng ý định nếu bạn muốn tính theo nhóm.
  • Thiếu thứ tự ổn định: Khi chỉ sắp xếp theo ngày, các dòng có cùng ngày có thể xuất hiện không theo thứ tự. Nên thêm cột phụ như transaction_id.
  • Kỳ vọng window làm giảm số dòng: Window Functions không nén dữ liệu — nếu muốn giữ lại chỉ các bản ghi top-ranked, hãy thêm cột rank rồi lọc.
  • Nhầm lẫn giữa row-based và time-based: rowsBetween(-6, 0) lấy 7 dòng, không phải 7 ngày. Dùng rangeBetween nếu bạn muốn định nghĩa theo thời gian.

Kết luận

Window Functions trong PySpark là công cụ đắc lực cho các kỹ sư dữ liệu khi cần phân tích sâu mà không đánh mất chi tiết. Bằng cách hiểu rõ sự kết hợp giữa partitionBy(), orderBy() và các khung cửa sổ, bạn có thể giải quyết hiệu quả nhiều bài toán thực tế: xếp hạng nội bộ, tổng lũy tiến, so sánh tương quan, và xác định các mẫu dữ liệu. Với khả năng mở rộng của Spark, những kỹ thuật này sẽ phát huy mạnh mẽ ngay cả trên những tập dữ liệu khổng lồ đặc trưng của hạ tầng dữ liệu hiện đại.

Chia sẻ:FacebookX
Nội dung tổng hợp bằng AI, mang tính tham khảo. Xem bài gốc ↗