# Spark cho FDE: đo dữ liệu của khách trước khi chọn cỡ cluster

> Khi khách hỏi "cần mấy node", bạn chỉ cần một ngày dữ liệu thật và vài lệnh đo là biết họ cần compaction, cần cluster hay một máy là đủ.

Bản gốc: https://fdetimes.net/vi/bach-khoa/spark-khi-du-lieu-khach-qua-lon/

Buổi họp kỹ thuật thứ hai với khách thường có một câu hỏi như thế này: "Bên mình nên dựng cluster Spark mấy node?" Lúc đó bạn chưa mở một file dữ liệu nào của họ. Trả lời bằng cảm tính thì nguy hiểm, vì một con số đã nói ra trong phòng họp thường sẽ bị ghi thẳng vào ngân sách.

Cách trả lời an toàn hơn là xin một ngày dữ liệu thật và vài giờ để đo. Đây không phải cách né câu hỏi. Các cấu hình quan trọng của Spark đều là ngưỡng tính bằng byte, số partition hay số core, nên khi chưa có số đo thì con số nào bạn đưa ra cũng chỉ là đoán.

Ở vị trí FDE, bạn có thể phải trả lời câu này trên một hệ thống lạ, trong ít thời gian, trước những người không quan tâm partition là gì. Vì thế, bạn cần một quy trình đo đủ nhanh để mang con số tới cuộc họp tiếp theo.

## Spark đã có sẵn ngưỡng, việc của bạn là đo dữ liệu để so

Spark SQL có sẵn vài ngưỡng mặc định để bạn đem số đo ra so. Khi đọc Parquet, JSON hay ORC, mỗi partition chứa tối đa 128 MB (`spark.sql.files.maxPartitionBytes`), vì vậy dung lượng đầu vào quyết định có bao nhiêu task đọc. Khi join hoặc aggregation, Spark chia dữ liệu thành 200 shuffle partition theo mặc định (`spark.sql.shuffle.partitions`).

Ngưỡng thứ ba là broadcast join: bảng nhỏ hơn 10 MB sẽ được gửi nguyên vẹn tới mọi worker, nhờ đó bảng lớn không phải shuffle. Ngưỡng thứ tư nằm trong tài liệu tuning: Spark khuyến nghị 2–3 task cho mỗi CPU core, tức số partition phải tính theo tổng số core của cluster.

Từ Spark 3.2.0, Adaptive Query Execution (AQE) được bật mặc định. AQE dùng thống kê lúc chạy để tối ưu lại plan, gộp các shuffle partition nhỏ theo kích thước mục tiêu mặc định 64 MB, và có thể tách partition bị lệch trong sort-merge join.

Theo tài liệu Spark, nhờ AQE bạn không cần tự chọn số shuffle partition cho vừa dataset. Dù vậy, AQE chỉ sửa được plan lúc chạy. Dung lượng thật, cách xếp file và độ lệch key thì bạn vẫn phải tự đo.

**Điểm mấu chốt:** Cỡ cluster được tính ra từ số đo dữ liệu, đừng chọn nó trước rồi mới đo.

## Thử đo một job chạy đêm

Lấy một ví dụ giả định. Một chuỗi bán lẻ có job chạy mỗi đêm: join bảng `orders` với bảng `stores` rồi tổng hợp doanh thu theo cửa hàng. Bạn xin một ngày dữ liệu và chạy ba lệnh đo, không cần dựng gì cả.

Lệnh đầu tiên đếm dung lượng và số file: `aws s3 ls --summarize --recursive s3://khach/orders/date=2026-10-01/`. Lệnh thứ hai dùng DuckDB để xem cách xếp file Parquet:

```sql
SELECT file_name,
count(DISTINCT row_group_id) AS so_row_group,
max(row_group_num_rows)      AS dong_moi_group_max
FROM parquet_metadata('orders/*.parquet')
GROUP BY file_name
ORDER BY dong_moi_group_max DESC
LIMIT 10;
```

Lệnh thứ ba là một notebook PySpark chạy local:

```python
from pyspark.sql import functions as F

orders = spark.read.parquet("orders/")
stores = spark.read.parquet("stores/")

print("read partitions:", orders.rdd.getNumPartitions())

total = orders.count()
(orders.groupBy("store_id").count()
.withColumn("ty_le", F.col("count") / total)
.orderBy(F.desc("count")).show(10))

orders.cache().count()   # rồi mở tab Storage trên web UI
```

Giả sử kết quả như sau: `orders` nặng 40 GB, chia thành 12.000 file, trung bình khoảng 3,4 MB mỗi file. Bảng `stores` nặng 6 MB. Một `store_id` duy nhất, kênh bán online, chiếm 35% số dòng. Mỗi kết quả này dẫn tới một quyết định khác nhau.

Nhưng file của khách nhỏ hơn rất nhiều so với khoảng 100 MB–10 GB mà hướng dẫn layout Parquet của DuckDB khuyến nghị.

Vì thế việc nên đề xuất đầu tiên là compaction, chưa cần thêm node. Nếu truy vấn metadata cho thấy file nào chỉ có một row group khổng lồ, bạn cũng biết file đó chỉ xử lý được bằng một thread. Row group hợp lý nằm trong khoảng 100K–1M dòng.

Bảng `stores` 6 MB, nằm dưới ngưỡng 10 MB, nên Spark sẽ broadcast nó và bảng `orders` không phải shuffle cho phép join. Hãy ghi lại con số này, vì đến khi khách mở thêm cửa hàng và bảng vượt 10 MB thì plan sẽ đổi mà không ai báo cho bạn.

Key chiếm 35% là con số cần theo dõi kỹ nhất. Tài liệu Spark nói thẳng rằng data skew có thể làm join chậm đi nghiêm trọng, nhưng tính năng tự tách partition lệch của AQE chỉ áp dụng cho sort-merge join.

Trên web UI, mở stage của aggregation và xem event timeline: một task chạy lâu hơn hẳn các task còn lại, hoặc có shuffle read lớn gấp nhiều lần, chính là dấu hiệu cần báo với khách từ trước.

Cuối cùng là bộ nhớ. Tài liệu Spark cho rằng cách ước lượng tốt nhất là cache dataset rồi xem trang Storage trên web UI. Nếu cache một ngày dữ liệu, bạn có số đo thật để nhân lên theo khung thời gian cần giữ. Con số này đáng tin hơn mọi quy tắc nhân hệ số.

## Có cần dựng cluster cho 40 GB không?

Trước khi tính đến cluster, nên có một phép đối chứng. Một bài viết của MotherDuck nhận xét rằng với dữ liệu nhỏ và vừa, Spark rốt cuộc chỉ chạy ở cấu hình tối thiểu và tốn nhiều overhead.

Bài viết lấy ví dụ AWS Glue: dịch vụ này yêu cầu tối thiểu 2 DPU, tức bạn trả tiền cho 32 GB RAM và 8 vCPU dù dataset mẫu chỉ khoảng 1 GB.

DuckDB chạy trên một máy và vẫn xử lý được dữ liệu lớn hơn RAM nhờ spill ra đĩa. Vì vậy bạn chạy đúng phép join và aggregation của job đêm trên DuckDB, đo thời gian, rồi so với khung giờ batch của khách. Buổi họp sau, bạn mang tới một con số đo được thay vì một ý kiến.

Ví dụ khung batch là 3 giờ và DuckDB chạy xong 40 GB trong 25 phút: job này chưa cần cluster.

Chỉ ba tín hiệu có thể đổi kết luận đó: thời gian chạy trên một máy tiến gần khung 3 giờ khi tính theo tốc độ tăng dữ liệu của khách; bảng dimension lớn tới mức hai bảng to phải join với nhau; hoặc lượng spill ra đĩa chiếm phần lớn thời gian chạy.

Khi một trong ba tín hiệu đó xuất hiện, khuyến nghị 2–3 task mỗi core giúp bạn kiểm tra cấu hình bằng phép tính. Giả sử cluster dự kiến có tổng 16 core: mức khuyến nghị tương ứng khoảng 32–48 partition.

Mỗi lúc vẫn chỉ có 16 task chạy, mỗi core lần lượt xử lý hai đến ba task, nên core nào xong sớm sẽ nhận việc tiếp thay vì ngồi chờ. Shuffle partition thì để AQE tự gộp theo mục tiêu 64 MB, đừng chỉnh tay con số 200.

## Những lỗi hay gặp ở hiện trường

Lỗi phổ biến nhất là copy cấu hình từ dự án trước sang, kèm theo một `shuffle.partitions` được chỉnh tay cho một dataset khác hẳn. Lỗi thứ hai là đo trên dữ liệu mẫu đã được làm sạch. Mẫu sạch không có key 35%, nên lúc đo mọi thứ đều chạy tốt.

Lỗi thứ ba là chỉ nhìn tổng dung lượng mà không xem cách xếp file. 40 GB trong 12.000 file nhỏ và 40 GB trong 40 file 1 GB là hai bài toán khác nhau. Lỗi thứ tư là quên rằng bảng dimension sẽ lớn dần. Nên đặt cảnh báo khi bảng gần chạm 10 MB, vì vượt ngưỡng đó thì plan sẽ đổi.

## Trả lời câu hỏi phỏng vấn bằng số đo

Một câu đáng tập trả lời trước khi đi phỏng vấn FDE: "Khách có 40 GB dữ liệu mới mỗi đêm và hỏi có nên dựng Spark cluster không. Bạn làm gì?" Câu trả lời yếu là kể tên công cụ. Câu trả lời mạnh hơn là trình bày một kế hoạch đo.

Có thể trả lời theo bốn bước. Đầu tiên, xin một ngày dữ liệu thật và đo dung lượng, số file, row group, kích thước bảng dimension và độ lệch key. Tiếp theo, cache mẫu để đo bộ nhớ thật.

Sau đó, chạy baseline trên một máy bằng DuckDB để so với khung batch. Cuối cùng, nói rõ con số nào sẽ khiến bạn đổi ý. Hãy dành thời gian cho bước cuối này, vì nó cho người nghe thấy bạn biết khi nào kết luận của chính mình hết đúng.

Trong CV, hãy viết thành kết quả cụ thể, ví dụ: "Đo profile dữ liệu 40 GB/đêm, đề xuất compaction và chạy trên một máy thay vì dựng cluster." Nếu JD nhắc tới tối ưu chi phí hạ tầng hay chọn cỡ hệ thống, đó là chỗ nên kể ví dụ này, kèm số liệu.

Lần tới khách hỏi "cần mấy node", hãy xin họ một ngày dữ liệu và hẹn hôm sau mang con số tới.

**Thử ngay tuần này:**

- Lấy một dataset Parquet công khai, chạy truy vấn parquet_metadata để đếm row group của từng file, rồi so kích thước file với khoảng 100 MB–10 GB.
- Chạy cùng một phép join và aggregation trên PySpark local và trên DuckDB, ghi lại thời gian và bộ nhớ của cả hai vào một bảng.
- Viết sẵn câu trả lời dài 2 phút cho câu hỏi "Khách có 40 GB mỗi đêm, có nên dựng cluster không?", trong đó có ít nhất ba con số bạn sẽ đo.

## Nguồn

- [Tuning - Spark 4.2.0 Documentation](https://spark.apache.org/docs/latest/tuning.html)

- [Performance Tuning - Spark 4.2.0 Documentation](https://spark.apache.org/docs/latest/sql-performance-tuning.html)

- [Making PySpark Code Faster with DuckDB](https://motherduck.com/blog/making-pyspark-code-faster-with-duckdb)

- [Tuning Workloads](https://duckdb.org/docs/current/guides/performance/how_to_tune_workloads.html)

- [File Formats](https://duckdb.org/docs/current/guides/performance/file_formats.html)
