PySpark 1 — Tổng quan, kiến trúc Python↔JVM & Setup
PySpark là gì — và bài này bổ sung gì
PySpark là API Python chính thức của Apache Spark. Nó không phải một engine riêng: engine vẫn là Spark chạy trên JVM (Java Virtual Machine). PySpark chỉ là lớp vỏ Python để bạn viết chương trình phân tán bằng cú pháp Python quen thuộc, còn phần tính toán nặng thực chất diễn ra trong tiến trình Java/Scala.
Bài này cố tình không lặp lại phần lý thuyết engine đã trình bày kỹ ở Spark 1 — Tổng quan & kiến trúc thực thi (mô hình in-memory, Driver–Executor, DAG, vì sao thắng MapReduce). Nếu bạn chưa nắm những khái niệm đó, hãy đọc bài kia trước. Ở đây ta tập trung vào góc nhìn Python thực hành: chuyện gì xảy ra giữa Python và JVM, vì sao có đoạn code nhanh và có đoạn chậm bất ngờ, cách setup, và khi nào nên (và không nên) chọn PySpark.
Đây là điểm mấu chốt phân biệt người dùng PySpark "biết dùng" với người "dùng được": hiểu ranh giới Python↔JVM. Không hiểu ranh giới này là nguyên nhân số một khiến một job PySpark chạy chậm gấp hàng chục lần mà không rõ tại sao.
Kiến trúc Python↔JVM qua Py4J
Khi bạn gõ from pyspark.sql import SparkSession rồi tạo session, có hai tiến trình được sinh ra ở phía driver:
- Tiến trình Python — nơi script
.pyhoặc notebook của bạn chạy. - Tiến trình JVM (Spark Driver thật) — nơi Spark xây dựng kế hoạch truy vấn, điều phối executor, quản lý DAG.
Hai tiến trình này nói chuyện với nhau qua Py4J — một thư viện cầu nối cho phép code Python gọi trực tiếp các đối tượng và phương thức Java trong JVM. Mỗi khi bạn viết df.filter(...).groupBy(...), Python không tự tính toán gì cả: nó chỉ gửi lệnh qua Py4J để JVM dựng cây thao tác (logical plan). Toàn bộ tối ưu hoá (Catalyst) và thực thi thật đều nằm trong JVM.
Vì sao DataFrame nhanh còn UDF Python chậm
Đây là hệ quả trực tiếp của kiến trúc trên, và là điều cần khắc cốt ghi tâm:
- Thao tác DataFrame/Spark SQL dựng sẵn (
filter,join,groupBy,withColumnvới hàm built-in nhưF.upper,F.sum,F.when...): code chỉ mô tả ý định, JVM thực thi hoàn toàn bằng mã Java/Scala đã tối ưu. Dữ liệu không bao giờ rời khỏi JVM. Tốc độ ngang với Scala Spark — vì thực chất nó là Scala Spark chạy bên dưới. - UDF Python (
udf): khi bạn viết một hàm Python thuần rồi đăng ký làm UDF, mỗi executor phải fork một tiến trình Python worker riêng. Với mỗi dòng (hoặc mỗi batch), dữ liệu bị serialize (pickle) từ JVM sang Python worker, hàm Python chạy, rồi kết quả serialize ngược về JVM. Chi phí đi-về này (kèm việc mất khả năng tối ưu của Catalyst) khiến UDF Python thường chậm hơn hàm built-in nhiều lần.
Quy tắc thực hành: ưu tiên hàm built-in pyspark.sql.functions tối đa; chỉ dùng UDF khi thật sự không có cách nào khác; và khi buộc phải dùng, chọn pandas_udf (vectorized UDF) dùng Apache Arrow để truyền dữ liệu theo batch cột thay vì từng dòng — nhanh hơn UDF thường rất nhiều. Chủ đề này được đào sâu ở PySpark 4 — UDF & pandas_udf.
SparkSession — điểm vào của mọi chương trình
Từ Spark 2.0, SparkSession là điểm vào thống nhất, thay cho việc phải tạo SQLContext/HiveContext riêng như đời cũ. Bên trong nó vẫn bọc một SparkContext — đối tượng đại diện kết nối tới cluster, quản lý cấu hình và tạo RDD ở tầng thấp. Trong công việc hằng ngày với DataFrame, bạn gần như chỉ chạm tới SparkSession; SparkContext (spark.sparkContext) chỉ dùng khi cần thao tác RDD cấp thấp hoặc đọc cấu hình cluster.
from pyspark.sql import SparkSession
spark = (
SparkSession.builder
.appName("ncb-txn-etl")
.master("local[*]") # chạy local, dùng hết CPU core
.config("spark.sql.shuffle.partitions", "200")
.getOrCreate()
)
# minh hoạ (Python) — không phải SQL sandbox
getOrCreate() trả về session đang tồn tại nếu có, nếu không thì tạo mới — nên gọi nhiều lần vẫn an toàn. Trong Databricks hay notebook được cấu hình sẵn, biến spark thường đã tồn tại; bạn không cần (và không nên) tự tạo lại.
Các môi trường chạy PySpark
master quyết định code chạy ở đâu:
master | Ý nghĩa | Khi nào dùng |
|---|---|---|
local[*] | Chạy trên 1 máy, dùng tất cả core | Dev, test, dữ liệu nhỏ |
local[4] | 1 máy, giới hạn 4 core | Dev có kiểm soát tài nguyên |
yarn | Cluster Hadoop/YARN | On-prem, hệ sinh thái Hadoop |
k8s://... | Cluster Kubernetes | Cloud-native, containerized |
spark://host:7077 | Standalone cluster | Cluster Spark tự dựng |
| (Databricks) | Managed, không set master thủ công | Databricks quản lý cluster |
- Local: cả driver lẫn executor nằm trong một JVM trên máy bạn. Rất tiện để học và test, nhưng không phải xử lý phân tán thật.
- Cluster (YARN/K8s/Standalone): driver điều phối, executor chạy rải trên nhiều node. Đây là nơi PySpark phát huy giá trị.
- Databricks / notebook (Jupyter, Colab): môi trường tương tác. Databricks là nền tảng thương mại quản lý sẵn cluster; Colab/Jupyter thường chạy
local[*]để học.
Khi nào chọn PySpark — và khi nào KHÔNG
Đây là quyết định kiến trúc quan trọng và hay bị làm sai theo hướng "dùng búa tạ đập hạt dẻ". Spark sinh ra cho dữ liệu lớn, phân tán. Với dữ liệu vừa một máy, nó thường chậm hơn và phức tạp hơn các công cụ single-node.
PySpark vs pandas / Polars / DuckDB (single-node)
- pandas: xử lý trong RAM một máy, API phong phú nhưng chậm và tốn bộ nhớ với dữ liệu lớn. Phù hợp vài MB đến vài GB.
- Polars: DataFrame viết bằng Rust, đa luống, lazy — nhanh hơn pandas nhiều lần trên một máy. Xử lý tốt vài chục GB nếu đủ RAM.
- DuckDB: "SQLite cho phân tích" — engine OLAP nhúng, chạy SQL trên file Parquet/CSV cực nhanh, không cần server.
- Xem tổng quan lựa chọn thư viện Python ở Python cho dữ liệu — Tổng quan.
Điểm chung của Polars/DuckDB: chúng thắng PySpark khi dữ liệu vừa một máy, vì không có chi phí điều phối cluster, không có shuffle qua mạng, không có overhead JVM khởi động. Nhiều pipeline "big data" thực ra chỉ vài GB — dùng Polars/DuckDB gọn và nhanh hơn hẳn.
PySpark thắng khi dữ liệu thật sự lớn — vượt sức một máy (hàng trăm GB đến hàng chục TB), cần chia ra hàng chục/hàng trăm node để xử lý song song, hoặc khi bạn đã có sẵn hạ tầng cluster/lakehouse và cần tích hợp với Delta Lake, Iceberg, streaming.
PySpark vs Scala Spark
Cả hai chạy trên cùng engine JVM. Khác biệt nằm ở lớp API:
| Tiêu chí | PySpark | Scala Spark |
|---|---|---|
| Hiệu năng DataFrame/SQL | Ngang nhau (đều xuống JVM) | Ngang nhau |
| Hiệu năng UDF | Chậm hơn (qua Py4J/pickle) | Nhanh (native JVM) |
| Hệ sinh thái ML/DS | Rất mạnh (pandas, scikit-learn, ML) | Hạn chế |
| Đường cong học | Thấp (Python phổ biến) | Cao hơn |
| Tính năng mới | Đôi khi ra sau vài phiên bản | Ra trước (API gốc) |
Vì sao đa số team dữ liệu ngân hàng chọn Python cho Spark? Vì hệ sinh thái ML/Data Science của Python áp đảo: cùng một ngôn ngữ để làm ETL, huấn luyện mô hình chấm điểm tín dụng, phân tích, và trực quan hoá. Data scientist không phải học Scala. Đánh đổi là: khi phải viết logic phức tạp bằng UDF, PySpark chậm hơn Scala; và một vài API mới đôi khi lên Scala trước. Nhưng chừng nào bạn giữ logic trong DataFrame/SQL built-in, khác biệt hiệu năng gần như bằng không — đây là lý do lời khuyên "tránh UDF Python" quan trọng đến vậy.
Lazy evaluation & Catalyst — nhắc lại ngắn
PySpark kế thừa nguyên vẹn hai đặc tính cốt lõi của engine (chi tiết ở Spark 2 — RDD & DataFrame):
- Lazy evaluation: các phép transformation (
select,filter,join,withColumn...) chỉ ghi lại ý định, chưa chạy gì. Chỉ khi gặp một action (show,count,collect,write) thì Spark mới thật sự thực thi. Nhờ đó Catalyst nhìn được toàn bộ chuỗi và tối ưu tổng thể (đẩy filter xuống sớm, gộp bước, cắt cột thừa). - Catalyst optimizer: bộ tối ưu truy vấn biến logical plan thành physical plan hiệu quả. Vì PySpark DataFrame cũng đi qua Catalyst, nó nhận được toàn bộ lợi ích tối ưu như Scala.
Thực hành quan trọng: DataFrame API là API chính. RDD (API cấp thấp) hiếm khi cần trong công việc hằng ngày và không được Catalyst tối ưu — chỉ dùng khi thật đặc thù. Với người mới, hãy coi như "luôn dùng DataFrame".
Ví dụ code PySpark tối thiểu
Dưới đây là một chương trình PySpark hoàn chỉnh nhỏ nhất: tạo session, đọc dữ liệu, một phép biến đổi, và một action. Đây là minh hoạ Python — không phải SQL sandbox nên không đánh dấu chạy được.
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = (
SparkSession.builder
.appName("txn-summary")
.master("local[*]")
.getOrCreate()
)
# Đọc dữ liệu giao dịch từ Parquet (lazy — chưa đọc thật)
txn = spark.read.parquet("/data/transactions/")
# Transformation: lọc giao dịch rút tiền, tính tổng theo tài khoản
summary = (
txn
.filter(F.col("kind") == "withdraw")
.groupBy("account_id")
.agg(F.sum("amount").alias("total_withdraw"))
.orderBy(F.desc("total_withdraw"))
)
# Action: đến đây Spark mới thực thi toàn bộ chuỗi
summary.show(10, truncate=False)
spark.stop()
Lưu ý: read.parquet, filter, groupBy, agg đều lazy — không có byte dữ liệu nào được xử lý. Chỉ tại show(10) (action) Spark mới dựng physical plan, phân task ra executor, và trả 10 dòng đầu về driver. Toàn bộ tính toán này chạy trong JVM; Python chỉ điều phối.
Lộ trình series "PySpark thực chiến"
| Bài | Nội dung |
|---|---|
| 01 — Tổng quan (bài này) | Kiến trúc Python↔JVM, SparkSession, setup, khi nào dùng |
| 02 — DataFrame API | select/filter/join/groupBy, cột, biểu thức, window |
| 03 — I/O & nguồn dữ liệu | Đọc/ghi Parquet, CSV, JDBC, Delta; partition, schema |
| 04 — UDF & pandas_udf | UDF thường vs vectorized, Arrow, khi nào dùng |
| 05 — pandas API on Spark | Viết cú pháp pandas, chạy phân tán trên Spark |
| 06 — MLlib | Pipeline ML phân tán, feature, huấn luyện quy mô lớn |
| 07 — Streaming & Testing | Structured Streaming trong Python, kiểm thử job |
| 08 — ETL ngân hàng | Dự án ETL giao dịch end-to-end tại NCB |
Nền tảng engine chung nằm ở series Spark chuyên sâu; PySpark series này bám vào thực hành Python.
Use case thực tế
Bối cảnh NCB. Khối dữ liệu cần xử lý bảng giao dịch lõi (core banking) để dựng báo cáo hành vi khách hàng và phục vụ mô hình phát hiện gian lận. Quy mô ước lượng: khoảng 8 triệu tài khoản hoạt động, trung bình mỗi tài khoản ~2–3 giao dịch/ngày → tầm 20 triệu giao dịch/ngày, tích luỹ ~200 triệu giao dịch/tháng. Lưu dạng Parquet, mỗi tháng nén lại vẫn cỡ vài trăm GB; nếu cần đọc lại 12–24 tháng lịch sử để tính đặc trưng (feature) thì chạm ngưỡng hàng TB.
Vì sao không xử lý trên một máy? Nhóm từng thử một job pandas đọc 3 tháng dữ liệu (~600 triệu dòng) trên máy 64 GB RAM — tiến trình bị OOM (out of memory) vì pandas nạp toàn bộ vào RAM và tốn bộ nhớ gấp nhiều lần kích thước file. Chuyển sang Polars/DuckDB cải thiện với 1–2 tháng, nhưng khi cần quét toàn bộ lịch sử nhiều TB để backfill feature thì một máy không đủ dung lượng lẫn thời gian.
Giải pháp PySpark. Job feature engineering chạy trên cluster YARN 20 executor (mỗi executor ~5 core, 20 GB RAM). Vì dữ liệu Parquet đã partition theo ngày, Spark chỉ đọc đúng partition cần (partition pruning), và toàn bộ phép join + aggregate giữ trong DataFrame built-in nên xuống thẳng JVM. Kết quả ước lượng:
- Backfill 12 tháng (~2,4 tỷ dòng) từ không chạy nổi trên 1 máy → hoàn tất trong khoảng 35–50 phút trên cluster.
- Job hằng ngày (20 triệu giao dịch) chạy tăng trưởng (incremental) chỉ vài phút.
- Bài học vận hành: một phiên bản đầu dùng UDF Python để chuẩn hoá mã giao dịch làm job chậm gấp ~4 lần; viết lại bằng
F.when/F.regexp_replacebuilt-in đưa thời gian về mức bình thường — đúng như lý thuyết ranh giới Python↔JVM ở trên.
Chi tiết pipeline này được dựng đầy đủ ở PySpark 8 — ETL ngân hàng.
Ghi nhớ
- PySpark là API Python, engine vẫn là JVM Spark. Python driver điều khiển JVM driver qua Py4J; code DataFrame chỉ mô tả ý định, JVM mới tính toán thật.
- Ranh giới Python↔JVM quyết định hiệu năng. Thao tác DataFrame/SQL built-in chạy trọn trong JVM → nhanh ngang Scala. UDF Python phải serialize dữ liệu qua lại Python worker → chậm. Ưu tiên built-in; nếu cần UDF thì dùng
pandas_udf(Arrow, vectorized). SparkSessionlà điểm vào duy nhất (bọcSparkContext).masterchọn nơi chạy:local[*]cho dev,yarn/k8s/standalone cho cluster; Databricks/notebook đã cấu hình sẵnspark.- Chọn công cụ theo quy mô, không theo trend. Dữ liệu vừa một máy → pandas / Polars / DuckDB nhanh gọn hơn. Dữ liệu lớn phân tán (hàng trăm GB–TB, cần nhiều node) → PySpark.
- Chọn Python cho Spark vì hệ sinh thái ML/DS; đánh đổi là UDF chậm hơn và vài API ra sau Scala — nhưng không đáng kể nếu giữ logic trong DataFrame built-in.
- DataFrame API là API chính; RDD hiếm khi cần và không được Catalyst tối ưu. Lazy evaluation: chỉ action (
show/count/write) mới kích hoạt thực thi.
Nguồn tham khảo
- Apache Spark Documentation — PySpark — tài liệu chính thức của API Python.
- Apache Spark Documentation — Spark SQL, DataFrames and Datasets Guide — DataFrame/SQL API và Catalyst.
- Apache Spark Documentation — Cluster Mode Overview — mô hình Driver–Executor và các loại
master(YARN, Kubernetes, Standalone). - Py4J Documentation — thư viện cầu nối Python↔JVM mà PySpark dùng.
- "Learning Spark, 2nd Edition" — Jules S. Damji, Brooke Wenig, Tathagata Das, Denny Lee (O'Reilly, 2020).
- "Spark: The Definitive Guide" — Bill Chambers & Matei Zaharia (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ẻ!