Spark Mastery: Bản đồ kiến thức Apache Spark từ A đến Z
Spark là chủ đề có nhiều bài viết nhất trong site này — nhưng đọc rời rạc thì khó thấy bức tranh chung. Bài này làm trục kiến thức: sắp mọi chủ đề thành lộ trình có thứ tự, trả lời gọn từng câu hỏi phỏng vấn kinh điển, và trỏ đến bài phân tích sâu tương ứng. Đọc xong trục này, bạn biết mình còn thiếu mảnh nào.
1. Bản đồ tổng thể
flowchart TB
subgraph L1["Tầng 1: Nền tảng"]
A["Kiến trúc & Execution Model<br/>(Driver, Executor, DAG)"]
B["Cluster Managers<br/>(Standalone / YARN / K8s)"]
C["RDD / DataFrame / Dataset"]
end
subgraph L2["Tầng 2: Cơ chế thực thi"]
D["Jobs → Stages → Tasks"]
E["Narrow vs Wide<br/>Transformations & Shuffle"]
F["Partitioning & Coalesce"]
end
subgraph L3["Tầng 3: Tối ưu"]
G["Persist / Storage Levels"]
H["Joins & Broadcast"]
I["Data Skew & Salting"]
J["AQE, Catalyst, Tungsten"]
K["Serialization & File Formats"]
end
subgraph L4["Tầng 4: Mở rộng"]
L["Structured Streaming"]
M["MLlib"]
N["Tích hợp HDFS / Cassandra / S3"]
end
L1 --> L2 --> L3 --> L4
2. Lộ trình đọc theo thứ tự
| # | Chủ đề | Bài trong site |
|---|---|---|
| 1 | Spark là gì, kiến trúc tổng quan | Apache Spark |
| 2 | Driver/Executor, lazy evaluation, DAG | Spark Execution Model |
| 3 | Chạy ở đâu: Standalone, YARN, K8s, deploy mode | Spark Cluster Managers |
| 4 | Ba tầng API + transformation/action cơ bản | RDD vs DataFrame vs Dataset |
| 5 | Job/Stage/Task, đọc Spark UI | Spark Jobs, Stages, Tasks |
| 6 | Shuffle — thao tác đắt nhất | Shuffle |
| 7 | Partition: sizing, repartition vs coalesce | Spark Partition |
| 8 | Cache/persist và storage levels | Persist & Storage Levels |
| 9 | Chiến lược join, broadcast join | Spark Joins |
| 10 | Broadcast variables, accumulators, Kryo, định dạng file | Broadcast & Serialization |
| 11 | Data skew: nhận diện và chữa | Data Skew + Salting |
| 12 | Bộ ba tối ưu tự động | Catalyst + Tungsten + AQE |
| 13 | Sự cố bộ nhớ | Troubleshooting OOM + Spill to Disk |
| 14 | SQL trên Spark | Spark SQL |
| 15 | Streaming | Streaming Processing, Watermark, Exactly-once |
| 16 | Luyện phỏng vấn tổng hợp | Spark Optimization Interview |
3. Trả lời nhanh các câu hỏi kinh điển
Narrow vs wide transformation khác nhau thế nào? Narrow: mỗi partition đầu ra chỉ cần dữ liệu từ một partition đầu vào (filter, map, select, union) — chạy nối tiếp trong cùng stage, không tốn mạng. Wide: partition đầu ra cần dữ liệu từ nhiều partition đầu vào (groupBy, join, distinct, orderBy, repartition) — buộc shuffle: ghi đĩa, truyền mạng, đọc lại, và tạo ranh giới stage mới. Kỹ năng tối ưu Spark, gói gọn, là nghệ thuật giảm số lượng và kích thước wide transformation.
Persist dữ liệu thế nào, có những storage level nào? cache()/persist(level) giữ kết quả sau action đầu tiên để các action sau khỏi tính lại. Các level: MEMORY_ONLY, MEMORY_AND_DISK (mặc định DataFrame), MEMORY_ONLY_SER, MEMORY_AND_DISK_SER, DISK_ONLY, biến thể _2 (replica). Quy tắc: chỉ cache thứ dùng ≥2 lần và đắt để tính lại, cache tại điểm hẹp nhất của pipeline, luôn unpersist() — chi tiết và bảng so sánh đầy đủ trong bài persist.
Xử lý data skew thế nào? Nhận diện qua Spark UI (max task duration >> median trong cùng stage). Chữa theo thứ tự chi phí: bật AQE (spark.sql.adaptive.skewJoin.enabled=true — Spark 3 tự chẻ partition lệch); broadcast join để né shuffle bảng lớn; tách riêng hot key xử lý riêng; cuối cùng mới đến salting — thêm hậu tố ngẫu nhiên vào key để rải đều, đổi lấy code phức tạp hơn.
Tối ưu bằng partitioning và coalescing? repartition(n) shuffle toàn bộ để chia đều — dùng khi cần tăng song song hoặc partition theo cột; coalesce(n) gộp partition không shuffle — dùng khi giảm số file đầu ra. Đích nhắm: partition ~128 MB, tránh cả nghìn file bé lẫn vài partition khổng lồ. Ghi ra lake thì partitionBy("dt") theo cột lọc phổ biến. Chi tiết trong Spark Partition.
Broadcast variable là gì? Biến read-only gửi một lần cho mỗi executor thay vì kèm theo từng task — nền tảng của broadcast join. Xem bài broadcast.
Spark tương tác với định dạng serialize nào? Ở tầng file: Parquet/ORC (cột — mặc định cho analytics nhờ pruning + pushdown), Avro (dòng, schema evolution tốt), JSON/CSV (nguồn thô, nên khai schema tay). Ở tầng runtime: Kryo nhanh gọn hơn Java serializer cho RDD/closure; DataFrame dùng Tungsten binary riêng. Xem bài serialization.
4. Spark Structured Streaming trong xử lý realtime
Mô hình của Spark: stream = bảng không có đáy, query streaming là query batch chạy lặp trên phần dữ liệu mới (micro-batch, độ trễ ~trăm ms đến giây). Cùng một API DataFrame:
events = (spark.readStream.format("kafka") .option("kafka.bootstrap.servers", "broker:9092") .option("subscribe", "orders").load())
agg = (events.selectExpr("CAST(value AS STRING) AS json") .select(F.from_json("json", schema).alias("e")).select("e.*") .withWatermark("event_time", "10 minutes") # chấp nhận trễ 10 phút .groupBy(F.window("event_time", "5 minutes"), "country") .agg(F.sum("amount").alias("revenue")))
(agg.writeStream.outputMode("update") .option("checkpointLocation", "s3a://chk/orders-agg/") # bắt buộc cho recovery .format("delta").start("s3a://gold/revenue_5m/"))Ba khái niệm quyết định độ đúng của kết quả: event-time vs processing-time, watermark (chờ dữ liệu trễ bao lâu trước khi chốt cửa sổ), và checkpoint để đạt exactly-once với sink hỗ trợ idempotent/transactional. Khi nào chọn Spark Streaming thay vì Flink: team đã dùng Spark, chấp nhận micro-batch, cần tích hợp lakehouse chặt; Flink thắng ở độ trễ thấp thật sự và state phức tạp.
5. Spark cho Machine Learning (MLlib)
Giá trị của MLlib không nằm ở thuật toán tân tiến nhất (deep learning thì dùng framework chuyên dụng) mà ở chỗ train trên dữ liệu không vừa một máy và cùng một pipeline chạy từ feature engineering đến scoring:
from pyspark.ml import Pipelinefrom pyspark.ml.feature import StringIndexer, VectorAssemblerfrom pyspark.ml.classification import GBTClassifier
pipeline = Pipeline(stages=[ StringIndexer(inputCol="country", outputCol="country_idx"), VectorAssembler(inputCols=["amount", "country_idx", "n_orders"], outputCol="features"), GBTClassifier(labelCol="churned", featuresCol="features"),])model = pipeline.fit(train_df) # train phân tánmodel.write().overwrite().save("s3a://models/churn-gbt/")scored = model.transform(new_df) # batch scoring hàng tỷ dòngĐiểm kiến trúc đáng nhớ: Pipeline serialize cả bước tiền xử lý lẫn model thành một artifact — loại bỏ training-serving skew (đúng mẫu dự án EcomLake dùng MLflow + SparkML). Dùng pyspark.ml (DataFrame-based); pyspark.mllib (RDD-based) đã ở chế độ bảo trì. Trong hệ sinh thái rộng hơn, Spark thường giữ vai feature engineering phân tán + batch scoring, còn training model phức tạp giao cho GPU framework — cầu nối là Feature Store và MLflow.
6. Tích hợp storage ngoài: HDFS, S3, Cassandra
Spark không có storage riêng — đó là ưu điểm kiến trúc (storage-compute decoupling): cùng một engine đọc được nhiều hệ lưu trữ, mỗi hệ một vai.
HDFS / S3 / MinIO (file system & object storage): nguồn và đích mặc định của batch analytics. Với HDFS, Spark còn tận dụng data locality — scheduler cố đặt task lên node đang giữ block dữ liệu, giảm đọc qua mạng; với S3/object storage thì locality không tồn tại, đổi lại co giãn và rẻ (lưu ý committer — xem dự án lakehouse).
Cassandra (NoSQL): qua Spark Cassandra Connector, mở ra mẫu kiến trúc “OLTP serving + OLAP analytics trên cùng dữ liệu”:
df = (spark.read.format("org.apache.spark.sql.cassandra") .options(table="user_events", keyspace="prod").load())# Connector đẩy filter theo partition key xuống Cassandra (predicate pushdown)recent = df.filter(F.col("user_id") == "u123") # → CQL WHERE, không full scanLợi thế của các tích hợp này trong pipeline: (1) pushdown — filter/column pruning được đẩy xuống tầng lưu trữ, Spark chỉ nhận phần cần; (2) không cần bước copy dữ liệu trung gian — đọc thẳng, transform, ghi thẳng sang hệ khác (Cassandra → Parquet trên S3 trong một job); (3) mỗi hệ giữ đúng vai — Cassandra phục vụ low-latency lookup, lake phục vụ scan analytics, Spark là cầu nối. Cái giá cần quản: job Spark quét mạnh có thể đè chết cluster Cassandra đang phục vụ production — giới hạn spark.cassandra.input.readsPerSec hoặc đọc từ replica/datacenter tách riêng.
Nguồn Tham Khảo
- Spark Documentation Overview - Apache Spark.
- Structured Streaming Programming Guide - Apache Spark.
- MLlib: Main Guide (DataFrame-based API) - Apache Spark.
- Spark SQL Data Sources - Apache Spark.
- Spark Cassandra Connector - DataStax.
- Hadoop-AWS: Integration with Amazon Web Services - Apache Hadoop.
🔗 Bài viết liên quan
Các nội dung khác có nhắc đến hoặc liên quan mật thiết với chủ đề này:
Bình luận & Thảo luận