PySpark 2 — DataFrame API bằng Python thực chiến
DataFrame API — nơi bạn viết 90% code hằng ngày
Ở bài mở đầu PySpark bạn đã dựng được SparkSession và chạy job đầu tiên. Bài này đi vào công cụ bạn dùng nhiều nhất: DataFrame API bằng Python. Nếu muốn hiểu vì sao engine chạy được như vậy — Catalyst tối ưu ra sao, sự khác nhau RDD và DataFrame — hãy đọc RDD vs DataFrame và Spark SQL. Ở đây ta tập trung vào cách viết: những phương thức nào, kết hợp thế nào, và các bẫy thực tế.
Một DataFrame là bảng phân tán có schema (tên cột + kiểu dữ liệu). Bạn không thao tác trực tiếp trên dữ liệu, mà mô tả phép biến đổi; Spark ghi lại thành một kế hoạch (plan) và chỉ thực thi khi có action. Toàn bộ API xoay quanh hai ý niệm: column expression (biểu thức cột) và transformation (phép biến đổi trả về DataFrame mới).
Column expression: nền tảng của mọi thứ
Điểm khác biệt lớn nhất so với pandas: trong PySpark, df["amount"] * 2 không tính toán gì cả. Nó tạo ra một đối tượng Column — một biểu thức mô tả "lấy cột amount rồi nhân 2". Spark chỉ tính khi bạn gọi action. Đây là lý do phải phân biệt rõ Column (biểu thức chạy trên cluster) và giá trị Python (chạy trên driver, cố định lúc build plan).
Quy ước gần như bắt buộc trong dự án thật:
# (minh hoạ)
from pyspark.sql import functions as F
from pyspark.sql import Window
# 3 cách trỏ tới cột — tương đương nhau
df.select(F.col("amount"))
df.select(df["amount"])
df.select(df.amount)
F.col("amount") là cách an toàn nhất vì tên cột dạng chuỗi, tránh nhầm lẫn khi cột chứa ký tự đặc biệt hoặc trùng tên sau join.
Các thao tác cột hay dùng
# (minh hoạ)
df2 = df.select(
F.col("account_id"),
(F.col("amount") * 1.1).alias("amount_vat"), # toán tử số học
F.when(F.col("amount") >= 100_000_000, "large") # rẽ nhánh điều kiện
.when(F.col("amount") >= 10_000_000, "medium")
.otherwise("small").alias("bucket"),
F.col("amount").cast("decimal(18,2)").alias("amount_dec"), # ép kiểu
F.lit("VND").alias("currency"), # hằng số -> cột
)
Bốn công cụ cần nhớ:
F.when(cond, val).otherwise(val)— tương đươngCASE WHENtrong SQL. Nối nhiều.when()để có nhiều nhánh..cast("type")— ép kiểu:"int","double","string","date","decimal(18,2)","timestamp"…F.lit(value)— biến một giá trị Python thành cột hằng. Bắt buộc khi so sánh cột với hằng trong vài hàm, hoặc khi thêm cột cố định..alias("name")— đặt tên cột kết quả (tránh tên xấu kiểuCASE WHEN...).
Toán tử logic dùng & (và), | (hoặc), ~ (phủ định) — không phải and/or của Python — và phải bọc ngoặc từng vế: (F.col("a") > 0) & (F.col("b") < 10).
Các phép biến đổi cốt lõi
# (minh hoạ)
(df
.select("account_id", "amount", "kind") # chọn cột
.filter(F.col("amount") > 0) # lọc dòng (where là bí danh)
.where("kind = 'transfer'") # where chấp nhận cả chuỗi SQL
.withColumn("fee", F.col("amount") * 0.001) # thêm/ghi đè cột
.withColumnRenamed("kind", "txn_kind") # đổi tên cột
.drop("amount") # bỏ cột
.distinct() # loại dòng trùng hoàn toàn
)
| Phương thức | Ý nghĩa | Ghi chú thực tế |
|---|---|---|
select(...) | chọn/tính cột | truyền tên chuỗi hoặc Column |
filter / where | lọc dòng | nhận Column hoặc chuỗi SQL |
withColumn(name, expr) | thêm/ghi đè 1 cột | gọi lặp nhiều lần thì chậm plan, ưu tiên select khi thêm nhiều cột |
withColumnRenamed(old, new) | đổi tên | — |
drop("c1","c2") | bỏ cột | không lỗi nếu cột không tồn tại |
distinct() | bỏ dòng trùng toàn phần | gây shuffle |
dropDuplicates(["k"]) | bỏ trùng theo tập cột | giữ 1 dòng/nhóm, không đảm bảo dòng nào nếu không orderBy |
orderBy(...) | sắp xếp | F.col("x").desc() để giảm dần |
limit(n) | lấy n dòng đầu | — |
Gộp nhóm với groupBy().agg()
Đây là trái tim của việc dựng bảng tổng hợp. Ưu tiên .agg() với các hàm trong F để tính nhiều chỉ tiêu cùng lúc:
# (minh hoạ)
summary = (df
.groupBy("branch_id", "txn_month")
.agg(
F.count("*").alias("txn_count"),
F.sum("amount").alias("total_amount"),
F.avg("amount").alias("avg_amount"),
F.countDistinct("account_id").alias("active_accounts"),
F.max("amount").alias("max_amount"),
F.collect_list("kind").alias("kinds"), # gom giá trị thành mảng
)
)
countDistinct đếm số phần tử phân biệt; collect_list/collect_set gom giá trị của nhóm thành mảng (cẩn thận nhóm quá lớn gây OOM trên executor).
JOIN: nối các bảng
Cú pháp chung: df_left.join(df_right, on=<điều kiện>, how=<kiểu>).
# (minh hoạ)
# Join theo tên cột chung -> không bị nhân đôi cột khoá
txn.join(acct, on="account_id", how="inner")
# Join theo điều kiện tường minh -> giữ cả hai cột khoá
txn.join(acct, txn["account_id"] == acct["id"], how="left")
Các kiểu join (how):
how | Giữ lại gì | Dùng khi |
|---|---|---|
inner | chỉ dòng khớp cả hai | mặc định, ghép dữ liệu chắc chắn có |
left (left_outer) | mọi dòng trái + khớp phải | làm giàu bảng trái, giữ nguyên số dòng trái |
right | mọi dòng phải + khớp trái | ít dùng, thường viết lại thành left |
outer (full) | tất cả hai bên | đối soát, tìm chênh lệch |
left_semi | dòng trái có khớp (không lấy cột phải) | lọc "tồn tại trong bảng B" |
left_anti | dòng trái không khớp | tìm "có ở A mà thiếu ở B" |
Broadcast join là kỹ thuật tối ưu quan trọng nhất khi một bảng nhỏ (vài chục MB trở xuống). Thay vì shuffle cả hai bảng lớn, Spark gửi bản sao bảng nhỏ tới mọi executor để join tại chỗ, tránh shuffle bảng lớn:
# (minh hoạ) — dim_branch nhỏ, fact_txn hàng tỉ dòng
fact_txn.join(F.broadcast(dim_branch), on="branch_id", how="left")
Khi join hai bảng lớn không cân, dữ liệu dễ bị skew (một khoá chiếm phần lớn dòng) khiến một task chạy mãi không xong. Cách xử lý skew (salting, AQE skew join) và cách chọn số partition xem shuffle & partitions; một pipeline ETL ngân hàng hoàn chỉnh dùng các phép này ở ETL ngân hàng với PySpark.
Window function: phân tích trong từng nhóm
Window cho phép tính toán trên một cửa sổ dòng liên quan tới dòng hiện tại — mà không gộp dữ liệu lại như groupBy. Đây là công cụ để xếp hạng, tính lũy kế, so với dòng trước/sau.
# (minh hoạ)
w_rank = Window.partitionBy("branch_id").orderBy(F.col("total_amount").desc())
w_run = Window.partitionBy("account_id").orderBy("created_at")
ranked = (summary
.withColumn("rank_in_branch", F.row_number().over(w_rank))
.withColumn("dense_rank", F.dense_rank().over(w_rank)))
flows = (txn
.withColumn("prev_amount", F.lag("amount", 1).over(w_run)) # giá trị dòng trước
.withColumn("next_amount", F.lead("amount", 1).over(w_run)) # giá trị dòng sau
.withColumn("running_total",
F.sum("amount").over(w_run))) # lũy kế theo thời gian
Các hàm hay dùng qua .over(window):
row_number()— số thứ tự 1,2,3… (không trùng). Dùng để lấy "dòng mới nhất mỗi khách": lọcrank_in_group == 1.rank()/dense_rank()— xếp hạng có/không nhảy bậc khi đồng hạng.lag(col, n)/lead(col, n)— lấy giá trị n dòng trước/sau, để tính chênh lệch, phát hiện thay đổi.sum()/avg()/count().over(w)— tổng/trung bình lũy kế hoặc theo nhóm mà vẫn giữ nguyên từng dòng.
Khái niệm window function chung (frame ROWS/RANGE, running total) trình bày sâu ở Spark SQL; cú pháp PySpark ở trên chỉ là một mặt của cùng cơ chế.
Hàm dựng sẵn hữu ích
Module F (pyspark.sql.functions) có hàng trăm hàm. Nhóm cần thuộc:
# (minh hoạ)
# Chuỗi
F.concat_ws("-", F.col("branch_id"), F.col("txn_month")) # nối có ký tự phân cách
F.split(F.col("full_path"), "/") # tách thành mảng
F.regexp_replace(F.col("phone"), r"\D", "") # bỏ ký tự không phải số
F.upper(F.col("city")); F.trim(F.col("name"))
# Ngày & thời gian
F.to_date(F.col("created_at")) # timestamp -> date
F.date_trunc("month", F.col("created_at")) # về đầu tháng
F.datediff(F.current_date(), F.col("hired_at")) # số ngày chênh
F.date_format(F.col("created_at"), "yyyy-MM") # định dạng chuỗi tháng
# Mảng & bảng chéo
F.explode(F.col("kinds")) # 1 dòng mảng -> nhiều dòng
# pivot: xoay giá trị cột thành cột
df.groupBy("branch_id").pivot("kind").agg(F.sum("amount"))
explode "bung" một cột mảng thành nhiều dòng (mỗi phần tử một dòng) — dùng khi dữ liệu nguồn dạng JSON lồng. pivot biến các giá trị phân biệt của một cột thành các cột riêng (ví dụ mỗi loại giao dịch một cột tổng tiền), rất tiện cho báo cáo.
Lazy evaluation & actions
Đây là điểm dễ sai nhất với người mới. Mọi transformation đều lười — chúng chỉ nối thêm vào plan, không chạy. Spark chỉ thực sự tính khi gặp một action.
| Loại | Ví dụ | Trả về |
|---|---|---|
| Transformation (lười) | select, filter, join, groupBy, withColumn | DataFrame mới |
| Action (kích hoạt job) | show(), count(), collect(), write..., toPandas() | kết quả / ghi ra |
Hệ quả thực tế:
collect()vàtoPandas()kéo toàn bộ dữ liệu về driver — chỉ dùng cho kết quả nhỏ, nếu không sẽ OOM driver.show(n)inndòng đầu, mặc định 20 — dùng để kiểm tra nhanh.count()phải quét toàn bộ, đừng gọi bừa trong vòng lặp.
explain() — đọc kế hoạch thực thi
# (minh hoạ)
result.explain() # xem physical plan
result.explain(True) # xem cả logical + optimized + physical
Đọc plan để kiểm tra: join có được broadcast không (tìm BroadcastHashJoin), có bao nhiêu Exchange (mỗi cái là một lần shuffle tốn kém), filter có bị đẩy xuống nguồn (PushedFilters) không.
cache() / persist() — và cảnh báo
Khi tái dùng một DataFrame nhiều lần (ví dụ tính vài chỉ tiêu khác nhau từ cùng bảng đã lọc), mỗi action sẽ tính lại từ đầu. cache() giữ kết quả trong bộ nhớ để lần sau dùng lại:
# (minh hoạ)
base = df.filter(F.col("amount") > 0)
base.cache() # đánh dấu; thực sự lưu khi action đầu chạy
base.count() # kích hoạt việc cache
# ... nhiều phép dùng lại base ...
base.unpersist() # nhớ giải phóng khi xong
Cảnh báo: đừng cache vô tội vạ. Cache chiếm bộ nhớ executor (có thể đẩy dữ liệu khác ra đĩa), và nếu chỉ dùng một lần thì cache chậm hơn. Chỉ cache khi (1) tái dùng ≥2 lần, và (2) phần tính trước đó đắt. Chi tiết tuning bộ nhớ, storage level xem Spark tuning.
So sánh nhanh với pandas và SQL
Người quen pandas/SQL chuyển sang PySpark cần định vị lại một vài phản xạ:
| Ý định | pandas | SQL | PySpark |
|---|---|---|---|
| Lọc dòng | df[df.amount>0] | WHERE amount>0 | df.filter(F.col("amount")>0) |
| Thêm cột | df["f"]=df.a*0.1 | a*0.1 AS f | df.withColumn("f", F.col("a")*0.1) |
| Gộp nhóm | groupby().agg() | GROUP BY | groupBy().agg() |
| Đổi tên | rename() | AS | withColumnRenamed() |
| Xếp hạng nhóm | groupby().rank() | RANK() OVER | F.rank().over(Window...) |
| Thực thi | ngay lập tức (eager) | ngay | lười (chờ action) |
| Dữ liệu | 1 máy, trong RAM | trong DB | phân tán, nhiều máy |
Điểm khác biệt cốt lõi: pandas tính ngay và giữ toàn bộ trong RAM một máy; PySpark lười và phân tán. Vì vậy đừng dùng vòng lặp Python duyệt từng dòng DataFrame — hãy diễn đạt bằng column expression để Spark xử lý song song.
Ghép lại: một pipeline tổng hợp giao dịch
# (minh hoạ) — tổng hợp giao dịch theo chi nhánh/tháng, xếp hạng, gắn tên khách
from pyspark.sql import functions as F
from pyspark.sql import Window
txn = spark.read.parquet("/lake/raw/transactions") # tỉ dòng
acct = spark.read.parquet("/lake/dim/accounts")
cust = spark.read.parquet("/lake/dim/customers") # bảng nhỏ
monthly = (txn
.filter(F.col("amount") > 0)
.withColumn("txn_month", F.date_format(F.col("created_at"), "yyyy-MM"))
.join(acct.select("id", "customer_id", "branch_id"),
txn["account_id"] == acct["id"], "left")
.groupBy("branch_id", "txn_month")
.agg(
F.count("*").alias("txn_count"),
F.sum("amount").alias("total_amount"),
F.countDistinct("customer_id").alias("active_customers"),
))
w = Window.partitionBy("txn_month").orderBy(F.col("total_amount").desc())
ranked = monthly.withColumn("branch_rank", F.row_number().over(w))
report = (ranked
.join(F.broadcast(cust.select(F.col("id").alias("customer_id"),
"full_name")),
on="customer_id", how="left") # ví dụ làm giàu; cust nhỏ -> broadcast
.filter(F.col("branch_rank") <= 10))
report.write.mode("overwrite").parquet("/lake/mart/branch_monthly_top")
Chuỗi trên là lười cho tới write — chỉ khi đó Spark mới tối ưu và chạy job.
Use case thực tế
Bối cảnh NCB — bảng tổng hợp giao dịch tỉ dòng. Đội Data Engineering cần một bảng mart phục vụ báo cáo lãnh đạo: mỗi chi nhánh, mỗi tháng, tổng doanh số giao dịch, số khách hoạt động, và thứ hạng chi nhánh trong tháng. Nguồn là bảng transactions khoảng 1,2 tỉ dòng/năm trên data lake (Parquet), cùng hai bảng chiều accounts (~8 triệu dòng) và branches (~250 dòng).
Cách làm với DataFrame API (các con số dưới đây là ước lượng minh hoạ):
- Đọc & lọc sớm:
read.parquetrồifilter(amount > 0)và giới hạn theo dải ngày cần chạy. Nhờ predicate pushdown, Spark chỉ đọc các file partition liên quan, giảm dữ liệu quét từ 1,2 tỉ xuống ~100 triệu dòng/tháng. - Chuẩn hóa tháng:
withColumn("txn_month", date_format(created_at, "yyyy-MM")). - Join lấy chi nhánh:
transactionsjoinaccountstheoaccount_idđể lấybranch_id. Đây là join hai bảng lớn — cần chọn số shuffle partition hợp lý và theo dõi skew (một chi nhánh lớn có thể ôm 15–20% giao dịch). - Gộp:
groupBy("branch_id","txn_month").agg(sum, count, countDistinct). - Broadcast bảng
branches(250 dòng, vài KB) bằngF.broadcastđể gắn tên chi nhánh — tránh shuffle không cần thiết. - Window xếp hạng:
row_number().over(partitionBy("txn_month").orderBy(total_amount.desc))để có top chi nhánh mỗi tháng, rồi lọcbranch_rank <= 20. - Ghi mart:
write.mode("overwrite").parquet(...)phân vùng theotxn_month.
Kết quả ước lượng: job chạy trên cluster 20 executor, mỗi cái 4 core / 16 GB, hoàn thành trong khoảng 8–12 phút cho một tháng, tạo ra bảng mart vài chục nghìn dòng. Toàn bộ chỉ dùng DataFrame API — không viết một dòng vòng lặp nào — và explain() xác nhận branches được BroadcastHashJoin, còn join lớn dùng SortMergeJoin với AQE bật để tự chia lại partition khi phát hiện skew.
Ghi nhớ
- Column expression là lười:
df["a"]*2chỉ mô tả phép tính, không chạy; phân biệt rõ Column (trên cluster) và giá trị Python (trên driver). - Bốn công cụ cột phải thuộc:
F.when().otherwise(),.cast(),F.lit(),.alias(). Toán tử logic dùng& | ~và bọc ngoặc từng vế. - Ưu tiên
groupBy().agg(...)với nhiều hàmFđể tính nhiều chỉ tiêu một lần;countDistinct/collect_listkhi cần. - Broadcast join (
F.broadcast) cho bảng nhỏ để tránh shuffle; cẩn thận skew khi join hai bảng lớn. - Window cho phân tích trong nhóm mà không gộp dòng:
row_numberlấy "dòng mới nhất",lag/leadso dòng kề,sum().overđể lũy kế. - Transformation lười, action mới chạy.
collect()/toPandas()kéo dữ liệu về driver — chỉ cho kết quả nhỏ. - Đọc
explain()để soi số lầnExchange(shuffle) và kiểu join; dùngcache()chỉ khi tái dùng ≥2 lần, và nhớunpersist(). - So với pandas: PySpark phân tán + lười — diễn đạt bằng column expression, đừng lặp Python trên từng dòng.
Nguồn tham khảo
- Apache Spark — Spark SQL, DataFrames and Datasets Guide
- Apache Spark — PySpark API Reference:
pyspark.sql.functions - Apache Spark — PySpark API Reference:
pyspark.sql.DataFrame - Apache Spark — PySpark Window Functions
- Damji, Wenig, Das, Lee — "Learning Spark: Lightning-Fast Data Analytics", 2nd Edition (O'Reilly, 2020)
- Chambers & Zaharia — "Spark: The Definitive Guide" (O'Reilly, 2018)
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ẻ!