PySpark 5 — pandas API on Spark: mở rộng pandas

14 thg 7, 2026 4 lượt xem
#data-engineering
#migration
#pyspark
#pandas
#pandas-on-spark

Vấn đề: pandas quen tay nhưng chạy một máy

Gần như mọi data analyst và data scientist ở NCB đều biết pandas. Nó là ngôn ngữ chung để khám phá dữ liệu: đọc file, lọc, groupby, merge, vẽ biểu đồ — tất cả gói gọn trong vài dòng đọc rất tự nhiên. Vấn đề chỉ lộ ra khi dữ liệu lớn lên. pandas là thư viện chạy trên một tiến trình, một máy, giữ toàn bộ DataFrame trong RAM. Một notebook chạy ngon trên file mẫu 2 triệu dòng sẽ MemoryError ngay khi bạn thử nạp bảng giao dịch thật hàng trăm triệu dòng. Cách chữa cháy quen thuộc là "lấy mẫu 5%", nhưng mẫu che mất đúng những đuôi phân phối mà rủi ro và gian lận hay nằm ở đó.

Có hai hướng thoát. Nếu dữ liệu vẫn vừa một máy khỏe (vài chục GB, tối ưu bộ nhớ), một engine như Polars nhanh hơn pandas nhiều lần và tiết kiệm RAM — đó là lựa chọn đúng cho quy mô một máy. Nhưng khi dữ liệu vượt quá một máy, bạn cần tính toán phân tán trên cụm. Đây là lúc pandas API on Spark xuất hiện: bạn giữ nguyên phong cách viết pandas, nhưng công việc chạy phân tán trên Spark.

pandas API on Spark là gì

pyspark.pandas là một module đi kèm PySpark, cung cấp lớp API mô phỏng pandas nhưng thực thi bằng Spark bên dưới. Dự án này trước đây tên là Koalas, phát triển độc lập, rồi được hợp nhất chính thức vào Apache Spark từ phiên bản 3.2. Ý tưởng cốt lõi rất gọn: bạn viết psdf.groupby("city").amount.sum() y như pandas, nhưng đối tượng psdf không phải pandas DataFrame — nó là một pandas-on-Spark DataFrame, và phép groupby().sum() đó được dịch thành một kế hoạch Spark chạy song song trên nhiều executor.

Nói cách khác, pyspark.pandaslớp tương thích cú pháp đặt trên DataFrame API. Nó không phải một engine mới; nó dịch các thao tác kiểu pandas thành các phép biến đổi Spark. Điều này mang lại lợi ích thực tế lớn nhất: migrate code pandas sẵn có lên quy mô lớn với công sức tối thiểu — nhiều khi chỉ cần đổi dòng import.

import pyspark.pandas as ps

# Đọc trực tiếp thành pandas-on-Spark DataFrame (psdf)
psdf = ps.read_parquet("s3://ncb-lake/transactions/")

# Cú pháp y hệt pandas — nhưng chạy phân tán trên Spark
result = (psdf[psdf["kind"] == "transfer"]
          .groupby("city")["amount"]
          .agg(["count", "sum", "mean"]))

Tạo và chuyển đổi psdf

Có nhiều đường vào ra với pandas-on-Spark. Nắm ba trục chuyển đổi này là hiểu được toàn bộ vị trí của nó trong hệ sinh thái:

  • Tạo trực tiếp: ps.DataFrame({...}), ps.read_parquet(...), ps.read_csv(...), ps.read_delta(...), ps.range(n) — tương ứng các hàm cùng tên của pandas.
  • Từ pandas: một pandas DataFrame pdf có thể nâng lên bằng ps.from_pandas(pdf). Lưu ý dữ liệu đang nằm ở driver được phân phối ra cụm.
  • Qua lại với Spark DataFrame: psdf.to_spark() trả về một Spark DataFrame thuần để cắm vào pipeline production; ngược lại sdf.pandas_api() (hoặc ps.DataFrame(sdf)) bọc một Spark DataFrame thành psdf. Chuyển đổi hai chiều này không tốn chi phí gom dữ liệu vì cả hai đều phân tán và lazy — nó chỉ đổi "vỏ" API.
  • Về pandas thật: psdf.to_pandas() gom toàn bộ dữ liệu về driver thành một pandas DataFrame. Đây là điểm nguy hiểm nhất: gọi trên bảng lớn sẽ làm driver hết RAM (OutOfMemoryError). Chỉ dùng khi kết quả đã đủ nhỏ (đã groupby/agg/head).

Sơ đồ cho thấy psdf là cây cầu giữa thế giới pandas một máy và Spark phân tán. Cạnh to_pandas() là cạnh duy nhất kéo dữ liệu về một máy — hãy coi nó như cửa thoát hiểm, không phải cửa chính.

API gần giống pandas — nhưng không phải pandas

Phần lớn thao tác quen thuộc đều có: chỉ mục theo cột psdf["col"], lọc boolean psdf[psdf.amount > 0], groupby, merge/join, sort_values, fillna, astype, describe, value_counts, apply, cả plot (vẽ trực tiếp, mặc định lấy mẫu để không kéo cả tỉ dòng về vẽ). Nhờ vậy một analyst đọc code pandas-on-Spark gần như không thấy lạ.

Nhưng "gần giống" không phải "y hệt". Có một số khác biệt bản chất bắt nguồn từ việc dữ liệu bị phân tán và tính toán lazy. Đây là phần quan trọng nhất của bài — hiểu sai chỗ này là nguồn gốc của những bug âm thầm khi migrate.

1. Thứ tự dòng không được đảm bảo

Trong pandas, DataFrame có thứ tự dòng cố định và một index rõ ràng. Trong pandas-on-Spark, dữ liệu nằm rải trên nhiều partition ở nhiều máy, nên thứ tự dòng mặc định không được đảm bảo giữa các lần chạy. Hệ quả: những thao tác ngầm dựa vào vị trí — ví dụ giả định "dòng thứ nhất là bản ghi cũ nhất" mà không sort — có thể cho kết quả khác nhau. Nếu cần thứ tự, phải nêu tường minh bằng sort_values(...).

2. Index có chi phí

pandas gắn chặt với khái niệm index (nhãn dòng). pandas-on-Spark phải mô phỏng index này trên nền Spark vốn không có khái niệm index. Việc mô phỏng đó tốn công: để đánh index tuần tự toàn cục, Spark có thể phải gom dữ liệu hoặc chạy một phép window đắt. Option compute.default_index_type điều khiển cách sinh index mặc định, với các lựa chọn:

  • "sequence" — index liên tục 0,1,2,… nhưng đắt vì cần một phép tính toàn cục (không song song hóa tốt).
  • "distributed-sequence" — vẫn liên tục và song song hơn, là mặc định ở các bản Spark gần đây; cần thêm một lượt quét để đếm.
  • "distributed" — rẻ nhất, song song hoàn toàn, nhưng index không liên tục (chỉ đảm bảo duy nhất, không đảm bảo 0,1,2,…).
import pyspark.pandas as ps
ps.set_option("compute.default_index_type", "distributed")

Bài học thực chiến: đừng phụ thuộc vào index như trong pandas. Nếu logic của bạn dựa vào giá trị index cụ thể, hãy đưa nó thành một cột dữ liệu thật.

3. apply/transform có thể chậm

Khi bạn gọi psdf.apply(fn) với một hàm Python tùy biến theo dòng, pandas-on-Spark phải chạy hàm Python đó trên các tiến trình Python của executor — đúng cơ chế và đúng chi phí của một UDF. Toàn bộ phân tích ở bài UDF, pandas_udf & Arrow áp dụng nguyên vẹn: dữ liệu phải vượt biên JVM↔Python, mất tối ưu Catalyst, và chậm hơn hàm dựng sẵn nhiều lần. Quy tắc vàng vẫn là: ưu tiên thao tác vector hóa (phép trên cả cột) thay vì apply từng dòng. Ví dụ, dùng psdf["amount"] * psdf["rate"] thay cho psdf.apply(lambda r: r.amount * r.rate, axis=1).

4. Lazy vs eager

pandas là eager: mỗi lệnh chạy ngay và trả kết quả tức thì. Spark là lazy: các phép biến đổi chỉ dựng kế hoạch, chỉ thực thi khi có một action. pandas-on-Spark kế thừa tính lazy này ở nhiều thao tác, nên đôi khi lỗi trong logic của bạn không phát nổ tại dòng viết ra nó mà tại dòng to_pandas(), count(), hay head() phía sau — làm việc debug bối rối hơn pandas. Bù lại, sự lazy cho phép Spark tối ưu cả chuỗi biến đổi trước khi chạy.

5. Một số API chưa hỗ trợ

pandas cực kỳ đồ sộ; pandas-on-Spark phủ phần lớn nhưng không phải 100%. Một số hàm, tham số, hoặc thao tác đặc thù (ví dụ vài kiểu reshape phức tạp, một số phương thức chuỗi/thời gian ít dùng) chưa có, hoặc có nhưng ngữ nghĩa hơi khác. Khi gặp, thường có ba lối: tìm thao tác tương đương, hạ xuống to_spark() để dùng DataFrame API cho bước đó, hoặc (nếu dữ liệu đủ nhỏ) to_pandas().

Bảng đối chiếu nhanh những kỳ vọng dễ sai:

Khía cạnhpandaspandas-on-Spark
Nơi chạy1 máy, trong RAMPhân tán trên cụm Spark
Thực thiEagerLazy (nhiều thao tác)
Thứ tự dòngCố địnhKhông đảm bảo nếu không sort
IndexMiễn phí, cốt lõiCó chi phí, cần cấu hình
apply theo dòngBình thườngNhư UDF — có thể chậm
Phủ APIToàn bộPhần lớn, không phải tất cả

Chọn công cụ nào cho việc gì

Ba lựa chọn hay bị nhầm lẫn. Nguyên tắc chọn:

  • pandas / Polars: dữ liệu vừa một máy. Nhanh, đơn giản, không cần cụm. Polars khi cần hiệu năng và tiết kiệm RAM; pandas khi cần hệ sinh thái rộng và code hiện có.
  • pandas API on Spark: khi bạn đã có nhiều code/notebook pandas và cần chạy nó trên dữ liệu vượt một máy, nhanh chóng, không muốn viết lại từ đầu. Đây là công cụ để migrate quy mô, không phải để tối ưu tối đa.
  • DataFrame API thuần: cho pipeline production cần hiệu năng và độ ổn định cao nhất. API này bám sát mô hình Spark, dễ tối ưu (predicate pushdown, tránh index thừa, kiểm soát partition), và là lựa chọn chuẩn cho job ETL chạy hằng ngày.

Một chiến lược thực tế rất hay dùng: bắt đầu bằng pandas-on-Spark để chạy được nhanh trên dữ liệu đầy đủ, rồi viết lại các bước "nóng" (tốn tài nguyên, chạy hằng ngày) bằng DataFrame API thuần. pandas-on-Spark đưa bạn từ "chạy được trên mẫu" tới "chạy được trên toàn bộ" trong vài giờ; DataFrame API đưa từ "chạy được" tới "chạy tối ưu ổn định".

Chiến lược migrate pandas → pandas-on-Spark

Quy trình migrate một notebook pandas hiện có, theo từng bước có kiểm soát:

  1. Đổi import. Thay import pandas as pd bằng import pyspark.pandas as ps (và đổi pd.ps. ở các lời gọi tạo/đọc dữ liệu). Với nhiều notebook đơn giản, chừng này đã chạy.
  2. Xử lý điểm gom dữ liệu. Rà mọi chỗ ngầm kéo dữ liệu về một máy: .to_numpy(), .values, .tolist(), vòng lặp for trên dòng, hoặc truyền thẳng psdf vào thư viện chỉ hiểu pandas/NumPy. Đây là nơi hay nổ RAM driver.
  3. Xử lý khác biệt ngữ nghĩa. Thêm sort_values ở chỗ dựa vào thứ tự; bỏ giả định về giá trị index; thay apply theo dòng bằng thao tác vector hóa; chọn compute.default_index_type phù hợp.
  4. Thay apply/vòng lặp bằng built-in. Mọi logic diễn đạt được bằng phép trên cột nên viết lại như vậy để tận dụng Catalyst.
  5. Kiểm thử đối chiếu. Chạy song song phiên bản pandas (trên mẫu) và pandas-on-Spark (trên cùng mẫu), so kết quả tổng hợp phải khớp. Sau đó mới bung ra toàn bộ dữ liệu.

Ba cạm bẫy hiệu năng cần khắc cốt:

  • Đừng gom về pandas giữa chừng. Một to_pandas() vô tình ở giữa pipeline biến job phân tán thành job một máy và thường làm sập driver.
  • Đừng apply row-wise. Nó chậm như UDF; luôn tìm cách vector hóa.
  • Cẩn thận với index đắt. Nếu không cần index tuần tự, chọn distributed để tránh phép tính toàn cục.

Ví dụ: cùng một phân tích, hai cách viết

Giả sử cần tính, theo từng thành phố, tổng và trung bình số tiền của các giao dịch chuyển khoản. Bằng pandas thuần (chạy một máy, minh họa):

import pandas as pd

pdf = pd.read_parquet("transactions_sample.parquet")
out = (pdf[pdf["kind"] == "transfer"]
       .groupby("city")["amount"]
       .agg(["count", "sum", "mean"])
       .sort_values("sum", ascending=False))
print(out.head(10))

Cùng logic, viết bằng pandas-on-Spark (chạy phân tán trên toàn bộ dữ liệu, minh họa):

import pyspark.pandas as ps

psdf = ps.read_parquet("s3://ncb-lake/transactions/")   # cả tỉ dòng
out = (psdf[psdf["kind"] == "transfer"]
       .groupby("city")["amount"]
       .agg(["count", "sum", "mean"])
       .sort_values("sum", ascending=False))

# out đã là bảng nhỏ (mỗi thành phố 1 dòng) → an toàn gom về xem
print(out.head(10).to_pandas())

Điểm khác biệt code gần như chỉ nằm ở dòng import và nguồn đọc. Điểm khác biệt vận hành thì lớn: bản trên đọc file mẫu vào RAM một máy, bản dưới quét cả kho lake phân tán, và chỉ gom về máy khi kết quả đã co lại thành bảng nhỏ.

Khi cần cắm kết quả vào pipeline Spark production, chỉ việc chuyển vỏ:

sdf = out.to_spark()          # psdf -> Spark DataFrame để ghi/join tiếp
sdf.write.mode("overwrite").parquet("s3://ncb-lake/marts/transfer_by_city/")

# Chiều ngược lại: bọc một Spark DataFrame thành psdf để dùng lại code pandas
existing_sdf = spark.read.table("core.transactions")
psdf2 = existing_sdf.pandas_api()

Use case thực tế

Bối cảnh. Đội Data Science của NCB có một notebook pandas dài đã dùng ổn định để phân tích hành vi giao dịch: đọc dữ liệu, lọc theo kind, groupby theo thành phố và tháng, tính vài chỉ số, vẽ biểu đồ phân phối. Vì pandas chạy một máy và bảng giao dịch quá lớn để nạp hết, đội buộc phải chạy trên mẫu 5% (ước lượng ~40 triệu trong tổng ~800 triệu dòng/năm). Mẫu này che mất các giao dịch giá trị rất lớn và các thành phố ít giao dịch — đúng những điểm cần soi cho rủi ro.

Cách làm. Thay vì viết lại toàn bộ notebook bằng DataFrame API (ước tính vài ngày công và rủi ro sai lệch logic), đội migrate sang pandas-on-Spark:

  1. Đổi import pandas as pdimport pyspark.pandas as ps, đổi pd.read_parquetps.read_parquet trỏ vào toàn bộ thư mục lake thay vì file mẫu.
  2. Rà và xử lý ba chỗ: một df.values truyền vào scikit-learn (thay bằng lấy mẫu tường minh rồi mới to_pandas), một apply tính nhãn theo dòng (viết lại bằng ps.Series vector hóa), và thêm sort_values cho bước dựa vào thứ tự.
  3. Đặt compute.default_index_type = "distributed" vì notebook không cần index tuần tự.
  4. Kiểm thử đối chiếu: chạy lại trên đúng mẫu 5% cũ, các con số tổng hợp khớp bản pandas → yên tâm bung ra 100% dữ liệu.

Kết quả (ước lượng, minh họa). Notebook chạy được trên toàn bộ ~800 triệu dòng thay vì mẫu 40 triệu, trên cụm Spark có sẵn, với công sức migrate khoảng nửa ngày thay vì vài ngày viết lại. Đội phát hiện một nhóm giao dịch giá trị lớn ở vài tỉnh mà bản lấy mẫu bỏ sót. Sau khi bản phân tích ổn định, hai bước chạy hằng ngày tốn tài nguyên nhất được viết lại bằng DataFrame API thuần để tối ưu — đúng chiến lược "migrate nhanh trước, tối ưu phần nóng sau".

Ghi nhớ

  • pandas API on Spark (import pyspark.pandas as ps, tiền thân Koalas, hợp nhất vào Spark 3.2) cho phép viết code kiểu pandas nhưng chạy phân tán trên Spark; nó là lớp tương thích cú pháp đặt trên DataFrame API, không phải engine mới.
  • Ba trục chuyển đổi: tạo trực tiếp (ps.read_parquet...), qua lại Spark DataFrame (to_spark() / sdf.pandas_api() — không tốn chi phí, chỉ đổi vỏ), và về pandas (to_pandas()gom về driver, chỉ dùng khi kết quả đã nhỏ).
  • Khác biệt phải nhớ so với pandas: thứ tự dòng không đảm bảo (cần sort), index có chi phí (chỉnh compute.default_index_type), apply theo dòng chậm như UDF, thực thi lazy chứ không eager, và một số API chưa hỗ trợ.
  • Chọn công cụ: pandas/Polars cho dữ liệu một máy; pandas-on-Spark để migrate code pandas sẵn có lên quy mô lớn nhanh; DataFrame API thuần cho pipeline production tối ưu nhất.
  • Cạm bẫy hiệu năng lớn nhất: gom về pandas giữa chừng (sập driver) và apply row-wise (chậm) — luôn vector hóa và chỉ to_pandas() khi dữ liệu đã co nhỏ.
  • Chiến lược thực chiến: migrate nhanh bằng pandas-on-Spark để chạy trên toàn bộ dữ liệu thay vì mẫu, kiểm thử đối chiếu trên mẫu cũ, rồi viết lại các bước nóng bằng DataFrame API.

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ẻ!