Bài này trả lời thẳng những câu hỏi hay gây bối rối nhất khi mới học: Spark là engine hay là ngôn ngữ, hay là database? Spark SQL nằm ở đâu trong bức tranh và chạy nó ở chỗ nào? Và DataFrame của Spark khác gì DataFrame của pandas mà bạn đã quen? Hiểu xong bài này, 6 bài sau chỉ còn là thực hành.
Câu trả lời ngắn gọn: Spark là một compute engine (động cơ tính toán) phân tán. Hãy nghĩ tới động cơ của một chiếc xe — bản thân nó không phải chiếc xe, không phải xăng, không phải con đường. Nó là bộ phận biến nhiên liệu thành chuyển động. Spark cũng vậy: nó là bộ phận nhận dữ liệu vào, thực hiện phép tính, trả kết quả ra — và làm điều đó bằng cách chia việc cho nhiều máy chạy song song.
Chỗ gây nhầm lẫn là chữ "engine". Nó không có nghĩa Spark là một kỹ sư (engineer) hay một phần mềm bạn "mở lên bấm nút". Trong ngành phần mềm, "engine" chỉ một lõi xử lý mà thứ khác gọi vào — ví dụ "game engine", "search engine", "database engine". Spark là data processing engine: bạn viết code (hoặc SQL), đưa cho Spark, Spark lo phần tính toán nặng.
Để định vị cho rõ, đây là những thứ Spark hay bị nhầm — và Spark thực chất khác chúng ra sao:
| Thứ hay bị nhầm với Spark | Nó là gì | Spark khác ở chỗ |
|---|---|---|
| pandas | Thư viện xử lý dữ liệu chạy trên 1 máy, trong 1 tiến trình Python. | Spark chia dữ liệu ra nhiều máy, xử lý song song. Xem so sánh chi tiết ở cuối bài. |
| MySQL / Oracle / PostgreSQL | Database — vừa lưu trữ dữ liệu lâu dài vừa truy vấn. | Spark không lưu trữ gì cả. Nó đọc dữ liệu từ nơi khác (file, S3, Kafka, DB), tính toán trong RAM, rồi ghi ra. Tắt Spark là dữ liệu biến mất. |
| SQL | Một ngôn ngữ để mô tả truy vấn. | Spark là engine thực thi SQL đó. SQL là "yêu cầu", Spark là "người làm". |
| Hadoop / HDFS | HDFS là hệ thống lưu file phân tán; MapReduce là engine đời cũ. | Spark là engine thay thế MapReduce (nhanh hơn nhờ tính trong RAM), thường đọc dữ liệu từ HDFS/S3 nhưng không cần Hadoop. |
| Airflow | Bộ lập lịch — quyết định job nào chạy lúc mấy giờ, theo thứ tự nào. | Airflow ra lệnh "chạy job Spark lúc 2h sáng"; Spark là thứ thực sự làm việc tính toán đó. |
Nhiều người tưởng "Spark" và "Spark SQL" là hai phần mềm khác nhau. Thực ra Spark SQL là một thư viện nằm bên trong Spark. Toàn bộ Spark được xây trên một lõi chung (Spark Core), và các thư viện chuyên biệt cắm lên trên lõi đó:
Điểm quan trọng: DataFrame và Spark SQL đều thuộc thư viện Spark SQL. Khi bạn viết df.groupBy(...) hay viết một câu SELECT ... GROUP BY ..., cả hai đều đi qua cùng một bộ tối ưu tên là Catalyst và được dịch ra cùng một kế hoạch thực thi. Đó là lý do trong Spark, viết bằng DataFrame API hay bằng SQL cho ra kết quả và tốc độ y hệt nhau — chỉ khác cú pháp.
Spark viết bằng Scala (chạy trên JVM), nhưng nó cho bạn nhiều "tay lái" để ra lệnh:
SELECT. Đây là "ngôn ngữ" đưa bạn tới phần tiếp theo.Đây là chỗ hay bí. Với MySQL bạn có Workbench, với Oracle có SQL Developer — "mở lên là gõ SQL". Spark thì không có một chỗ cố định như vậy, vì Spark là engine — nó chạy ở bất cứ nơi nào bạn khởi động một SparkSession. Dưới đây là các "cửa" phổ biến để gõ Spark SQL, từ dễ tới production:
spark.sql("..."). Đây là cách bạn dùng trong lab này (JupyterLab). Đăng ký DataFrame thành view rồi query:
sales.createOrReplaceTempView("sales") spark.sql(""" SELECT category, SUM(qty*unit_price) AS revenue FROM sales GROUP BY category """).show()Kết quả trả về vẫn là một DataFrame — dùng tiếp bình thường. Đây là Bài 3 của khóa.
spark-sql — CLI gõ SQL trực tiếp. Trong terminal có Spark, gõ spark-sql sẽ mở một dấu nhắc giống hệt SQL client cổ điển: bạn gõ thẳng SELECT ...; và nhấn Enter, không cần viết một dòng Python nào. Gần nhất với trải nghiệm "MySQL Workbench" mà bạn quen.pyspark / spark-shell — REPL tương tác. Gõ pyspark mở một shell Python đã có sẵn biến spark, thử lệnh từng dòng. spark-shell là bản Scala tương đương.spark-submit file.py — chạy production. Khi code đã xong, đóng gói thành file .py rồi spark-submit lên cụm để chạy hàng loạt (thường được Airflow gọi theo lịch). Không tương tác, chạy một mạch.spark.sql().Cả hai đều là "bảng dữ liệu có cột và dòng", cùng gọi là DataFrame — nên rất dễ tưởng chúng như nhau. Nhưng chúng khác nhau ở bản chất vận hành. Đây là bảng so sánh đầy đủ:
| Tiêu chí | pandas DataFrame | Spark DataFrame |
|---|---|---|
| Chạy ở đâu | 1 máy, 1 tiến trình Python | Nhiều máy (cụm), nhiều executor song song |
| Giới hạn dữ liệu | Phải vừa RAM của 1 máy (vài GB) | Terabyte — chia partition, tràn ra đĩa nếu cần |
| Thực thi | Eager — chạy ngay từng dòng lệnh | Lazy — chỉ chạy khi gặp action (show/count/write) |
| Tối ưu tự động | Không — chạy đúng thứ tự bạn viết | Có — Catalyst tối ưu toàn chuỗi trước khi chạy |
| Mutable? | Có — sửa tại chỗ (df["x"]=...) | Bất biến — mỗi phép tạo DataFrame mới |
| Chỉ số (index) | Có index dòng | Không có index |
| Thứ tự dòng | Được giữ nguyên | Không đảm bảo (dữ liệu rải nhiều partition) |
| Tốc độ với data nhỏ | Rất nhanh (không có overhead) | Chậm hơn vì chi phí khởi động/điều phối |
| Tốc độ với data lớn | Chậm dần rồi hết RAM → crash | Co giãn — thêm máy là chạy tiếp |
| Chạy SQL | Không trực tiếp (cần thư viện ngoài) | Có sẵn spark.sql() |
Nhìn code cạnh nhau sẽ rõ hơn. Cùng một việc "thêm cột thành tiền rồi lọc" — cú pháp na ná, nhưng cơ chế bên dưới rất khác:
import pandas as pd df = pd.read_csv("sales.csv") df["amount"] = df["qty"] * df["unit_price"] df = df[df["qty"] > 0] # chạy NGAY, kết quả có liền trong RAM print(df.head())
from pyspark.sql import functions as F df = spark.read.csv("sales.csv", header=True) df = (df .withColumn("amount", F.col("qty")*F.col("unit_price")) .filter(F.col("qty") > 0)) # CHƯA chạy gì — tới show() mới thực thi df.show()
Khác biệt cốt lõi nằm ở dòng cuối: pandas đã tính xong ngay khi bạn gán df["amount"]. PySpark thì mọi thứ còn là "kế hoạch" cho tới khi gọi show() — lúc đó Catalyst mới nhìn cả chuỗi, tối ưu, rồi rải task xuống executor.
toPandas(): khi dữ liệu Spark đã gom nhỏ (ví dụ kết quả tổng hợp còn vài trăm dòng), bạn có thể result.toPandas() để đưa về pandas mà vẽ biểu đồ (matplotlib/seaborn) hay xử lý tiếp bằng thư viện Python quen thuộc.
summary = spark.sql("SELECT category, SUM(amount) rev FROM sales GROUP BY category") pdf = summary.toPandas() # giờ là pandas, nhỏ gọn, vẽ chart thoải mái pdf.plot.bar(x="category", y="rev")
toPandas() (hay collect()) trên DataFrame lớn hàng triệu/tỉ dòng. Toàn bộ dữ liệu sẽ bị kéo về RAM của một máy driver — đúng cái giới hạn mà Spark sinh ra để tránh. Chỉ toPandas() sau khi đã lọc/tổng hợp cho nhỏ.toPandas() phần nhỏ để vẽ.spark.sql() trong notebook (cách của khóa này), CLI spark-sql, Thrift Server cho BI, hay nền tảng như Databricks. Nó luôn chạy trong cụm Spark của bạn.toPandas() để bắc cầu — nhưng chỉ sau khi đã thu nhỏ dữ liệu.