PySpark 2 — DataFrame API bằng Python thực chiến

14 thg 7, 2026 4 lượt xem
#data-engineering
#pyspark
#dataframe
#python
#transformations

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 DataFrameSpark 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 đương CASE WHEN trong 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ểu CASE 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ĩaGhi chú thực tế
select(...)chọn/tính cộttruyền tên chuỗi hoặc Column
filter / wherelọc dòngnhận Column hoặc chuỗi SQL
withColumn(name, expr)thêm/ghi đè 1 cộtgọ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ộtkhông lỗi nếu cột không tồn tại
distinct()bỏ dòng trùng toàn phầngây shuffle
dropDuplicates(["k"])bỏ trùng theo tập cộtgiữ 1 dòng/nhóm, không đảm bảo dòng nào nếu không orderBy
orderBy(...)sắp xếpF.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):

howGiữ lại gìDùng khi
innerchỉ dòng khớp cả haimặc định, ghép dữ liệu chắc chắn có
left (left_outer)mọi dòng trái + khớp phảilàm giàu bảng trái, giữ nguyên số dòng trái
rightmọ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_semidòng trái khớp (không lấy cột phải)lọc "tồn tại trong bảng B"
left_antidòng trái không khớptì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ọc rank_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ạiVí dụTrả về
Transformation (lười)select, filter, join, groupBy, withColumnDataFrame mới
Action (kích hoạt job)show(), count(), collect(), write..., toPandas()kết quả / ghi ra

Hệ quả thực tế:

  • collect()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) in n dò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ạ:

Ý địnhpandasSQLPySpark
Lọc dòngdf[df.amount>0]WHERE amount>0df.filter(F.col("amount")>0)
Thêm cộtdf["f"]=df.a*0.1a*0.1 AS fdf.withColumn("f", F.col("a")*0.1)
Gộp nhómgroupby().agg()GROUP BYgroupBy().agg()
Đổi tênrename()ASwithColumnRenamed()
Xếp hạng nhómgroupby().rank()RANK() OVERF.rank().over(Window...)
Thực thingay lập tức (eager)ngaylười (chờ action)
Dữ liệu1 máy, trong RAMtrong DBphâ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ườiphâ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ạ):

  1. Đọc & lọc sớm: read.parquet rồi filter(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.
  2. Chuẩn hóa tháng: withColumn("txn_month", date_format(created_at, "yyyy-MM")).
  3. Join lấy chi nhánh: transactions join accounts theo account_id để lấy branch_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).
  4. Gộp: groupBy("branch_id","txn_month").agg(sum, count, countDistinct).
  5. Broadcast bảng branches (250 dòng, vài KB) bằng F.broadcast để gắn tên chi nhánh — tránh shuffle không cần thiết.
  6. 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ọc branch_rank <= 20.
  7. Ghi mart: write.mode("overwrite").parquet(...) phân vùng theo txn_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"]*2 chỉ 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 & | ~bọc ngoặc từng vế.
  • Ưu tiên groupBy().agg(...) với nhiều hàm F để tính nhiều chỉ tiêu một lần; countDistinct/collect_list khi 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_number lấy "dòng mới nhất", lag/lead so 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ần Exchange (shuffle) và kiểu join; dùng cache() 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

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.

13 thg 7, 2026 10

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.

13 thg 7, 2026 9

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.

13 thg 7, 2026 8

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.

13 thg 7, 2026 7

Cảm nhận của bạn

Bình luận

Bạn cần để viết bình luận.

Chưa có bình luận. Hãy là người đầu tiên chia sẻ!