0

Làm thế nào để xử lý vấn đề Small File Problem, cách thiết lập job tự động gộp các file Parquet nhỏ thành file lớn trên Data Lake?

"Small File Problem" (Vấn nạn file nhỏ) chính là kẻ thù thầm lặng giết chết hiệu năng của mọi Data Lake [cite: x]. Khi các luồng streaming (như Kafka/Debezium) liên tục xả dữ liệu xuống S3, chúng thường tạo ra hàng ngàn file Parquet bé tí (vài KB đến vài MB) [cite: x].

Hậu quả là khi truy vấn, hệ thống sẽ tốn nhiều thời gian để mở/đóng kết nối mạng (Network I/O) và đọc siêu dữ liệu (Metadata) hơn là thời gian thực sự đọc dữ liệu [cite: x]. Mục tiêu của chúng ta là gộp chúng lại thành các file có kích thước lý tưởng từ 128MB đến 1GB [cite: x].

Quá trình này được gọi là Compaction (Gộp file) [cite: x]. Dưới đây là các chiến lược và cách thiết lập tự động hóa từ cơ bản đến "hạng nặng" [cite: x].


1. Bản chất kỹ thuật của quá trình Compaction

Về mặt logic, Compaction cực kỳ đơn giản [cite: x]:

  1. Đọc toàn bộ các file nhỏ trong một Partition (ví dụ: year=2026/month=08/day=05) [cite: x].
  2. Gộp chúng lại trên RAM [cite: x].
  3. Ghi đè lại xuống S3 dưới dạng 1 hoặc một vài file Parquet lớn [cite: x].
  4. Xóa các file nhỏ cũ đi [cite: x].

Tuy nhiên, bài toán hóc búa ở đây là Tính toàn vẹn (ACID) [cite: x]. Nếu đang ghi file lớn mà server bị sập, hoặc có người đang query dữ liệu đúng lúc bạn đang xóa file cũ thì sao? [cite: x] Đó là lý do chúng ta cần các công cụ chuyên dụng thay vì tự viết script copy-paste thủ công [cite: x].


2. Cách 1: Sử dụng Apache Spark (Giải pháp truyền thống & Phổ biến nhất)

Nếu Data Lake của bạn đang dùng thuần Parquet, công cụ mạnh nhất để xử lý khối lượng dữ liệu này là Apache Spark (chạy trên AWS EMR hoặc AWS Glue) [cite: x].

Kịch bản (PySpark Script):

Bạn viết một Job đọc dữ liệu của "ngày hôm qua", dùng hàm coalesce() hoặc repartition() để ép số lượng file đầu ra, sau đó ghi đè (overwrite) lại vào đúng thư mục đó [cite: x].

from pyspark.sql import SparkSession

spark = SparkSession.builder.appName("ParquetCompaction").getOrCreate()

# 1. Đường dẫn tới thư mục cần gộp (Ví dụ: Dữ liệu ngày hôm qua)
partition_path = "s3://bo-doi-datalake/events/year=2026/month=08/day=05/"

# 2. Đọc toàn bộ các file nhỏ
df = spark.read.parquet(partition_path)

# 3. Ép số lượng file đầu ra (Ví dụ: ép về 1 file duy nhất nếu dung lượng < 1GB)
# Dùng coalesce thay vì repartition để tránh xáo trộn dữ liệu (shuffle) không cần thiết
df_compacted = df.coalesce(1)

# 4. Ghi đè lại dữ liệu (chế độ overwrite)
df_compacted.write.mode("overwrite").parquet(partition_path)

Cách tự động hóa (Scheduling):

  • Dùng AWS Glue + EventBridge: Đóng gói đoạn code trên thành một Glue Job [cite: x]. Setup CloudWatch Events (EventBridge) để trigger Job này chạy vào lúc 2h sáng mỗi ngày (Off-peak hours), dọn dẹp các thư mục của ngày hôm trước [cite: x].
  • Dùng Apache Airflow: Nếu hệ thống có sẵn Airflow, hãy tạo một DAG với SparkSubmitOperator chạy theo lịch cron 0 2 * * * [cite: x].

3. Cách 2: Nâng cấp lên Data Lakehouse (Delta Lake / Apache Iceberg) - Tiêu chuẩn Senior

Việc tự viết script ghi đè file Parquet tiềm ẩn rủi ro hỏng dữ liệu nếu Job chết giữa chừng [cite: x]. Trong kiến trúc Data hiện đại, giới kỹ sư đã chuyển sang sử dụng các Table Formats như Delta Lake (của Databricks) hoặc Apache Iceberg thay cho Parquet thuần [cite: x].

Các định dạng này vẫn lưu dữ liệu bằng file Parquet bên dưới, nhưng có thêm một lớp quản lý Transaction Log, cho phép tính năng ACID và tự động hóa Compaction cực kỳ thanh lịch [cite: x].

Nếu bạn chuyển đổi bảng sang định dạng Delta Lake hoặc Iceberg, việc gộp file chỉ tốn đúng 1 dòng SQL duy nhất [cite: x]:

-- Dành cho Delta Lake hoặc Apache Iceberg (chạy qua Athena, Spark, hoặc Trino)
OPTIMIZE user_events;

Cách hoạt động của lệnh OPTIMIZE:

  • Hệ thống sẽ tự động quét ngầm (background) các file nhỏ [cite: x].
  • Tạo ra các file lớn (Target size thường mặc định là 1GB) [cite: x].
  • Đổi con trỏ Metadata sang file mới mà không làm gián đoạn (Zero downtime) các câu lệnh SELECT đang chạy [cite: x].
  • Sau đó, bạn chỉ cần chạy lệnh VACUUM user_events; để xóa hẳn các file nhỏ cũ đi khỏi ổ cứng [cite: x].

4. Cách 3: Compaction "nhà nghèo" bằng DuckDB (Backend Approach)

Nếu hệ thống chưa đủ lớn để phải gánh chi phí khổng lồ của Spark hay AWS Glue, bạn có thể tận dụng chính con server Backend (Node.js/Go) kết hợp với DuckDB để làm một cronjob nhẹ nhàng chạy mỗi đêm [cite: x].

Logic tự động hóa (Ví dụ chạy cronjob Node.js/Python lúc nửa đêm):

-- Dùng DuckDB đọc toàn bộ thư mục S3 và tạo bảng tạm trên RAM
CREATE TABLE temp_compact AS 
SELECT * FROM read_parquet('s3://bo-doi-datalake/events/year=2026/month=08/day=05/*.parquet');

-- Ghi ngược trở lại S3 thành 1 file duy nhất, nén chuẩn Snappy
COPY temp_compact TO 's3://bo-doi-datalake/events/year=2026/month=08/day=05/compacted_data.parquet' 
(FORMAT 'parquet', CODEC 'snappy');

Lưu ý: Sau khi file compacted_data.parquet được tạo thành công, đoạn script của bạn sẽ gọi AWS S3 SDK (hàm deleteObjects) để xóa các file gốc li ti đi, chỉ chừa lại file lớn [cite: x].


5. Những nguyên tắc "Sống Còn" khi thiết kế Compaction Job

  • Idempotency (Tính không thay đổi trạng thái): Đảm bảo rằng nếu Job bị lỗi và chạy lại 10 lần, dữ liệu của bạn không bị nhân bản (duplicate) lên 10 lần [cite: x]. Luôn dùng chế độ overwrite toàn bộ partition, hoặc ghi ra thư mục tạm (temp) rồi dùng API đổi tên (rename) nếu S3 hỗ trợ [cite: x].
  • Luôn xử lý độ trễ (Late-arriving data): Nếu bạn gộp file của ngày hôm qua, nhưng hôm nay vẫn có một vài sự kiện mạng bị trễ rơi rớt vào thư mục ngày hôm qua thì sao? [cite: x] Job của bạn phải đủ thông minh để không xóa nhầm các file mới này, hoặc sử dụng Delta/Iceberg để giải quyết triệt để vấn đề này [cite: x].
  • Giám sát (Monitoring): Đặt cảnh báo (Alert) trên Slack/Telegram [cite: x]. Nếu Job gộp file chết 3 ngày liên tiếp, hàng vạn file nhỏ sẽ tích tụ và làm sập query của team Data [cite: x]. Phải luôn biết khi nào Job thất bại [cite: x].

All rights reserved

Viblo
Hãy đăng ký một tài khoản Viblo để nhận được nhiều bài viết thú vị hơn.
Đăng kí