Batch 6 — Partitioning, File Format & Tối ưu I/O

22 thg 7, 2026 2 lượt xem
#data-engineering
#parquet
#compaction
#partitioning
#batch
#file-format

Mô hình tinh thần: job nhanh hay chậm được quyết định trước khi engine chạy

Có một sự thật hay bị bỏ qua: phần lớn thời gian một job batch không dùng để tính toán, mà để đọc dữ liệu từ đĩa/object store. Bạn có thể tối ưu code, tăng số executor, chỉnh bộ nhớ — nhưng nếu bố cục lưu trữ (storage layout) bắt engine đọc 2 TB để rồi vứt đi 1.98 TB thì mọi tối ưu tính toán đều vô nghĩa. Đọc batch idempotency hay ETL vs ELT cho ta biết khi nàobằng cách nào nạp dữ liệu; bài này trả lời câu hỏi bổ sung: dữ liệu nằm trên đĩa như thế nào để lần đọc sau rẻ nhất có thể.

Mô hình tinh thần gồm ba đòn bẩy, xếp theo thứ tự tác động:

  1. Đọc ít file hơn — partition pruning: chỉ mở những thư mục chứa dữ liệu bạn cần.
  2. Đọc ít cột hơn — column pruning + định dạng cột: chỉ đọc cột query đụng tới.
  3. Đọc ít dòng hơn trong mỗi file — predicate pushdown + thống kê min/max: bỏ qua cả block dữ liệu không khớp filter.

Ba đòn bẩy này đều là I/O — không phải CPU. Chúng cộng dồn theo cấp số nhân: một query quét đúng 1 partition, đúng 3/40 cột, đúng vài row-group khớp filter có thể chỉ đọc dưới 1% kích thước bảng.

Partition: chia dữ liệu theo khoá để prune và chạy song song

Partitioning là chia một bảng logic thành nhiều thư mục vật lý theo giá trị của một (vài) cột khoá. Trên data lake/lakehouse, quy ước phổ biến là Hive-style partitioning: đường dẫn nhúng cả tên và giá trị cột.

s3://bank-dwh/silver/transactions/
  ├─ event_date=2026-07-19/
  │    ├─ part-00000.parquet
  │    └─ part-00001.parquet
  ├─ event_date=2026-07-20/
  └─ event_date=2026-07-21/

Khi query có filter WHERE event_date = '2026-07-21', engine đọc metadata thư mục, thấy chỉ cần mở đúng một thư mục và bỏ qua toàn bộ các thư mục ngày khác — đây là partition pruning. Không đọc, không giải nén, không parse. Đây là đòn bẩy I/O mạnh nhất vì nó loại việc ở mức thư mục, trước cả khi chạm vào nội dung file.

Partition đồng thời là đơn vị song song hoá tự nhiên: mỗi partition (thường tách tiếp thành nhiều file) là một đơn vị việc độc lập, engine giao cho các task khác nhau. Đây cũng chính là ranh giới cho ghi idempotent kiểu INSERT OVERWRITE một partition đã bàn ở batch idempotency: backfill lại ngày 2026-07-20 chỉ cần ghi đè đúng thư mục event_date=2026-07-20/, không đụng ngày khác — xem thêm backfill & reprocessing.

Chọn khoá partition: quy tắc và bẫy

  • Partition theo cột hay xuất hiện trong WHERE, thường là thời gian (ngày nạp / event date). Đây là default an toàn cho ETL batch vì phần lớn job xử lý theo lô ngày.
  • Cardinality vừa phải. Partition theo customer_id (hàng triệu giá trị) tạo hàng triệu thư mục mỗi thư mục vài KB → thảm hoạ small files (mục dưới). Partition theo ngày cho ~365 thư mục/năm là hợp lý.
  • Tránh over-partitioning. Partition lồng quá sâu (year/month/day/hour/region/product) chia nhỏ dữ liệu đến mức mỗi file tí hon; metadata phình, planning chậm. Quy tắc thực dụng: mỗi partition nên chứa ít nhất vài trăm MB dữ liệu.
  • Phân biệt partition column vs bucket/sort column. Partition prune ở mức thư mục; cột cardinality cao nên xử lý bằng bucketing hoặc sort trong file (mục dưới), không phải partition.

Partition tốt = filter thường dùng nhất, cardinality đủ thấp để mỗi thư mục đủ lớn. Sai khoá partition khó sửa hồi tố vì phải viết lại toàn bộ layout.

Định dạng cột vs định dạng dòng

Định dạng file quyết định đòn bẩy #2 và #3. Chia làm hai họ:

Row-orientedColumn-oriented
Ví dụCSV, JSON, AvroParquet, ORC
Cách lưuLần lượt từng dòng đủ mọi cộtGom giá trị cùng một cột lại với nhau
Đọc vài cột / nhiều dòng (OLAP)Phải đọc cả dòng rồi bỏ cột thừaChỉ đọc đúng cột cần
Ghi/đọc cả một bản ghi (OLTP, streaming ingest)TốtKém hơn
NénTrung bìnhRất tốt (cùng cột → cùng kiểu dữ liệu)
SchemaCSV không có; Avro/JSON cóCó, kèm thống kê

Nguyên tắc chọn: dữ liệu analytics ở tầng silver/gold, đọc theo kiểu quét-vài-cột-nhiều-dòng → Parquet/ORC. Dữ liệu landing/bronze thô, hoặc trao đổi bản ghi từng cái (message queue, CDC), hoặc cần schema evolution mạnh khi ghi → Avro (row, có schema). CSV/JSON tiện cho trao đổi liên thông và debug nhưng không nên là định dạng lưu trữ phân tích: không nén tốt, không thống kê, không pushdown.

Bên trong một file Parquet

Hiểu cấu trúc Parquet giải thích vì sao nó nhanh. Một file Parquet chia thành các row group (khối vài chục–vài trăm MB dòng). Trong mỗi row group, dữ liệu lưu theo column chunk — mỗi cột một khối liên tục, chia tiếp thành page. Ở cuối file là footer chứa schema và thống kê (min/max, số null, count) cho từng column chunk của từng row group.

Sơ đồ trên gói cả hai đòn bẩy còn lại:

  • Column pruning: query chỉ đụng amounttxn_type → engine chỉ đọc hai column chunk đó, không chạm acct_id, channel. Bảng 40 cột mà query 3 cột thì tiết kiệm ~90% I/O — điều bất khả với CSV/JSON.
  • Predicate pushdown: filter amount > 1e9 được đẩy xuống tầng đọc file. Engine so với thống kê min/max: Row group 1 có max(amount)=50k < 1e9toàn bộ row group bị bỏ qua, không giải nén một byte nào. Chỉ Row group 2 (max=3e9) mới có khả năng khớp nên mới mở ra đọc.

Để pushdown hiệu quả, dữ liệu nên có tính cục bộ (clustered) theo cột filter — nếu giá trị amount rải ngẫu nhiên khắp các row group thì min/max của group nào cũng rộng, chẳng prune được gì. Đây là lý do của sort/bucketing (mục dưới).

Nén

Định dạng cột nén rất tốt vì trong một column chunk các giá trị cùng kiểu, thường lặp lại. Parquet dùng encoding (dictionary, run-length, bit-packing) rồi mới đến codec nén khối:

  • Snappy — mặc định phổ biến; nén nhẹ, giải nén rất nhanh, cân bằng tốt cho query nóng.
  • Zstd — tỉ lệ nén cao hơn Snappy mà tốc độ vẫn tốt; ngày càng là lựa chọn ưa thích cho dữ liệu lạnh/lưu lâu.
  • Gzip — nén cao, giải nén chậm hơn; ít dùng cho hot path.

Nén không chỉ tiết kiệm dung lượng lưu trữ (tiền object store) mà còn giảm số byte phải đọc qua mạng — thường là nút cổ chai thật sự trên lake tách rời compute–storage.

Small files problem & compaction

Đây là căn bệnh kinh điển của mọi data lake. Job streaming micro-batch, hay job batch song song cao, hay ghi đè partition liên tục sẽ đẻ ra rất nhiều file nhỏ (hàng KB tới vài MB). Mỗi lần đọc, engine tốn chi phí cố định cho mỗi file: mở kết nối, đọc footer, lập lịch một task. Với object store như S3/GCS, độ trễ mỗi request (list, open) là đáng kể; hàng trăm nghìn file nhỏ khiến thời gian liệt kê + mở file lấn át thời gian đọc dữ liệu thật. Metadata/planning cũng phình.

Compaction (còn gọi bin-packing / file consolidation) là quá trình đọc nhiều file nhỏ trong một partition rồi ghi lại thành ít file lớn hơn đạt kích thước mục tiêu. Có thể chạy như một bước bảo trì định kỳ, hoặc tự động trong table format hiện đại:

Điểm mấu chốt: compaction phải chạy an toàn với reader đồng thời. Table format ACID commit file mới rồi mới ẩn file cũ trong một transaction nguyên tử, nên query đang chạy không thấy trạng thái dở dang — đây chính là giá trị của table format so với thư mục Parquet trần.

Kích thước file tối ưu

Không có con số vàng tuyệt đối, nhưng nguyên tắc định tính rõ ràng:

  • Quá nhỏ (KB–vài MB) → chi phí mở file/lập lịch lấn át; small files problem.
  • Quá lớn (nhiều GB) → một file = một đơn vị đọc khó chia, giảm song song; sửa/ghi lại tốn kém.
  • Vùng ngọt thực dụng thường trong khoảng ~128 MB đến ~1 GB mỗi file (nhiều đội chọn quanh 256 MB–512 MB), căn cho row group nằm gọn để pushdown và đọc song song hiệu quả. Hãy coi đây là điểm khởi đầu để đo, không phải luật.

Bucketing & sort: xử lý cột cardinality cao trong file

Partition không hợp với cột cardinality cao (như account_id). Hai kỹ thuật bù:

  • Sort trong file (sort dữ liệu theo cột filter trước khi ghi): làm giá trị cục bộ hoá, khiến min/max mỗi row group hẹp → predicate pushdown prune được nhiều row group. Đây là cách rẻ và luôn nên cân nhắc cho cột hay lọc theo range.
  • Bucketing (hash cột khoá vào N bucket cố định, mỗi bucket một file): hai bảng cùng bucket theo cùng khoá có thể join mà không cần shuffle (bucket-join), vì các dòng cùng khoá đã nằm cùng vị trí — bổ trợ cho phần shuffle & partition trong Spark. Đổi lại, bucketing cứng nhắc (số bucket cố định) và phải khớp giữa các bảng.

Iceberg tổng quát hoá ý này bằng hidden partitioning + partition transform (ví dụ bucket(N, id), days(ts)): người viết query không cần biết cột partition vật lý, engine tự áp transform — tránh bẫy quên filter đúng cột partition.

Table format cho ACID batch: Delta / Iceberg

Thư mục Parquet + partition trần đã đủ nhanh cho đọc, nhưng thiếu ba thứ mà batch nghiêm túc cần: ghi nguyên tử, cô lập reader/writer, và quản lý metadata/thống kê ở quy mô lớn. Table format (Delta Lake, Apache Iceberg, Hudi) thêm một lớp metadata/transaction log lên trên các file dữ liệu (vẫn là Parquet/ORC bên dưới) để có:

  • ACID commit: INSERT OVERWRITE một partition, compaction, hay MERGE đều là commit nguyên tử — không bao giờ lộ trạng thái dở dang cho reader (nền tảng cho idempotency ở batch-04).
  • Metadata prune ở quy mô lớn: thay vì list hàng trăm nghìn thư mục trên object store (chậm), engine đọc file metadata liệt kê sẵn file nào chứa khoảng giá trị nào → prune nhanh hơn partition thư mục thuần.
  • Schema/partition evolution: đổi khoá partition mà không viết lại lịch sử (Iceberg), thêm cột an toàn.
  • Time travel & snapshot: đọc lại phiên bản cũ để audit/backfill — xem ACID & time travel của Iceberg.

Nói cách khác, table format không thay thế các nguyên lý ở trên — nó tự động hoá chúng (compaction, thống kê, pruning) và bọc trong giao dịch. Các nguyên lý partition/định dạng cột/kích thước file vẫn đúng nguyên; bạn chỉ điều khiển chúng qua API bảng thay vì tự quản thư mục.

Code: viết layout tối ưu bằng PySpark

Ví dụ ghi bảng giao dịch tầng silver: partition theo ngày, sort trong partition để tăng pushdown, nén Zstd, và căn kích thước file.

from pyspark.sql import functions as F

# df: giao dịch đã làm sạch ở tầng silver
# 1) Căn số file mỗi partition ngày ~ dung lượng mong muốn.
#    repartition theo đúng cột partition để mỗi thư mục ngày gọn số file,
#    tránh mỗi task ghi rải rác nhiều file nhỏ (small files).
out = (
    df
    .repartition("event_date")          # gom theo khoá partition
    .sortWithinPartitions("event_date", "account_id")  # cục bộ hoá cho pushdown
)

# 2) Ghi Parquet, Hive-style partition theo ngày, nén zstd.
(
    out.write
       .format("parquet")
       .partitionBy("event_date")       # -> event_date=YYYY-MM-DD/
       .option("compression", "zstd")
       .mode("overwrite")               # ghi đè idempotent phạm vi partition
       .save("s3://bank-dwh/silver/transactions")
)

# 3) Đọc: engine tự partition-prune + column-prune + predicate pushdown.
q = (
    spark.read.parquet("s3://bank-dwh/silver/transactions")
         .where((F.col("event_date") == "2026-07-21")   # partition pruning
                & (F.col("amount") > 1_000_000_000))     # predicate pushdown
         .select("account_id", "amount", "txn_type")     # column pruning
)

Cùng bố cục đó dưới dạng bảng giao dịch (Delta) để có ACID + compaction tự động:

# Ghi Delta thay cho Parquet trần
out.write.format("delta").partitionBy("event_date").mode("overwrite") \
   .save("s3://bank-dwh/silver/transactions_delta")

# Compaction định kỳ + co-locate theo account_id (chạy như bước bảo trì)
spark.sql("""
  OPTIMIZE delta.`s3://bank-dwh/silver/transactions_delta`
  WHERE event_date >= '2026-07-01'
  ZORDER BY (account_id)
""")

Chi tiết cú pháp cụ thể theo phiên bản engine có thể khác nhau; xem I/O & nguồn dữ liệu trong PySpark và tài liệu Delta/Iceberg cho phiên bản bạn dùng.

Use case thực tế

Bối cảnh NCB (số liệu minh hoạ): bảng silver.transactions lưu giao dịch mọi kênh, ~120 triệu dòng/ngày, giữ 3 năm ≈ ~130 tỉ dòng, ~40 cột. Job EOD phân loại nợ và báo cáo giám sát thường chỉ cần 1 ngày, ~5 cột (account_id, amount, txn_type, channel, event_date).

  • Trước tối ưu: lưu một cây thư mục JSON gzip không partition. Query 1 ngày phải liệt kê và quét toàn bộ bảng ~ nhiều TB; job EOD chạy hàng giờ, thường timeout khi backfill.
  • Sau tối ưu: chuyển sang Parquet + partitionBy(event_date) + zstd + sortWithinPartitions theo account_id.
    • Partition pruning: query 1 ngày chỉ mở 1/1095 thư mục.
    • Column pruning: đọc 5/40 cột.
    • Predicate pushdown: filter amount prune tiếp phần lớn row group nhờ dữ liệu đã sort.
    • Kết quả minh hoạ: byte thực đọc giảm từ cả bảng xuống dưới 1%; EOD từ hàng giờ còn vài phút.
  • Small files: job micro-batch bổ sung trong ngày đẻ ~800 file nhỏ/partition. Thêm bước compaction ban đêm (OPTIMIZE trên Delta) gộp về ~256 MB/file; thời gian đọc partition nóng giảm rõ vì hết chi phí mở nghìn file.
  • Backfill: khi sửa logic phân loại nợ và cần chạy lại 30 ngày, chỉ INSERT OVERWRITE 30 partition liên quan — nguyên tử, không đụng lịch sử khác (xem backfill).

Ghi nhớ

  • Job batch chậm phần lớn vì I/O, và I/O được quyết định bởi bố cục lưu trữ đặt ra trước khi engine chạy — không phải bởi số executor.
  • Ba đòn bẩy cộng dồn: partition pruning (đọc ít file), column pruning + định dạng cột (đọc ít cột), predicate pushdown + thống kê min/max (đọc ít dòng).
  • Partition theo cột hay lọc, cardinality vừa phải (thường là ngày); tránh over-partitioning và partition theo cột cardinality cao.
  • Dùng Parquet/ORC cho tầng analytics; Avro cho row-level/ingest; CSV/JSON chỉ để trao đổi, không để lưu phân tích. Nén Snappy/Zstd.
  • Small files problem giết hiệu năng đọc trên object store; giải bằng compaction về vùng file ~128 MB–1 GB, chạy an toàn với reader nhờ commit ACID.
  • Sort trong file làm hẹp min/max → tăng pushdown; bucketing cho phép join không shuffle với cột khoá cardinality cao.
  • Table format (Delta/Iceberg) không thay thế các nguyên lý trên mà tự động hoá + bọc giao dịch (ACID overwrite partition, compaction, pruning bằng metadata, time travel).

Nguồn tham khảo

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.

13 thg 7, 2026 10

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.

13 thg 7, 2026 9

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.

13 thg 7, 2026 8

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.

13 thg 7, 2026 7

Cảm nhận của bạn

Bình luận

Bạn cần để viết bình luận.

Chưa có bình luận. Hãy là người đầu tiên chia sẻ!