Iceberg 3 — ACID, Time Travel & MERGE

13 thg 7, 2026 3 lượt xem
#data-engineering
#merge
#acid
#time-travel
#iceberg

Vì sao ACID trên data lake lại khó

Data lake cổ điển ghi thẳng file Parquet lên object storage (S3, GCS, HDFS) rồi để engine quét cả thư mục. Mô hình này không có giao dịch: nếu một job Spark đang ghi 200 file mà chết giữa chừng, người đọc thấy một tập file dở dang — nửa cũ nửa mới. Không có khái niệm "trạng thái nhất quán tại một thời điểm", không có cách rollback, và hai job ghi song song có thể giẫm lên nhau. Đây chính là lý do các nghiệp vụ nhạy cảm như đối soát cuối ngày hay báo cáo tuân thủ ngại đặt lên lake thô.

Apache Iceberg giải bài toán này bằng một ý tưởng gọn: không sửa dữ liệu tại chỗ, mà tạo phiên bản mới của metadata trỏ tới tập file hợp lệ. Mỗi phiên bản đó là một snapshot. Việc "công bố" thay đổi thu về đúng một thao tác nguyên tử: đổi con trỏ ở catalog từ metadata cũ sang metadata mới. Bài này mổ xẻ cơ chế ACID đó, rồi khai thác ba năng lực mà nó mở ra và cực kỳ giá trị cho ngân hàng: time travel, rollback, và row-level MERGE/UPDATE/DELETE.

ACID = snapshot + đổi con trỏ nguyên tử

Nhắc lại kiến trúc metadata đã trình bày ở bài kiến trúc: một bảng Iceberg gồm nhiều tầng — file metadata.json ở đỉnh, bên dưới là manifest list của từng snapshot, rồi các manifest file liệt kê data file. Mỗi lần ghi, Iceberg không sửa các file cũ. Nó viết thêm data file mới, viết manifest mới, viết một metadata.json mới chứa snapshot mới, rồi yêu cầu catalog trỏ bảng sang metadata mới đó.

Điểm mấu chốt của ACID nằm ở bước cuối: atomic commit. Catalog (Hive Metastore, JDBC, Glue, Nessie, REST catalog…) thực hiện một phép hoán đổi con trỏ có tính nguyên tử — kiểu compare-and-swap "chỉ đổi từ metadata v5 sang v6 nếu con trỏ hiện tại đang là v5". Trước khi con trỏ đổi, mọi người đọc vẫn thấy snapshot cũ trọn vẹn; sau khi đổi, mọi người đọc thấy snapshot mới trọn vẹn. Không có trạng thái ở giữa. Đó là chữ A (atomicity)C/I (consistency, isolation) của ACID.

  • Atomicity: một commit hoặc thấy toàn bộ, hoặc không thấy gì. Job chết giữa chừng chỉ để lại data file "mồ côi" chưa được snapshot nào tham chiếu — người đọc không bao giờ thấy chúng (dọn sau bằng maintenance, xem bài 6).
  • Consistency: người đọc luôn thấy một snapshot hoàn chỉnh và tự nhất quán.
  • Isolation: reader gắn vào snapshot tại thời điểm mở query, writer commit snapshot mới song song không ảnh hưởng reader đang chạy.
  • Durability: dữ liệu và metadata nằm trên object storage bền vững; catalog lưu con trỏ.

Mức cô lập: snapshot vs serializable

Iceberg hỗ trợ hai mức isolation, cấu hình qua thuộc tính bảng:

  • Snapshot isolation (mặc định cho nhiều thao tác): writer đọc snapshot base, tính toán, rồi commit. Chỉ xung đột khi có bên khác xóa/ghi đè cùng data file mà mình dựa vào.
  • Serializable (chặt hơn): còn kiểm cả trường hợp bên khác thêm dữ liệu vào phạm vi mình đang thao tác (ví dụ có row mới lọt vào điều kiện WHERE của một DELETE), tránh phenomenon "phantom". Đổi lại tỉ lệ xung đột — và retry — cao hơn.

Có thể đặt riêng cho đường ghi, ví dụ write.delete.isolation-level / write.update.isolation-level / write.merge.isolation-level = serializable hoặc snapshot.

Optimistic concurrency: nhiều writer cùng ghi

Iceberg dùng optimistic concurrency control. Writer không khóa bảng trước. Nó làm việc trên snapshot base rồi cố commit; catalog chỉ chấp nhận nếu con trỏ vẫn ở đúng snapshot base đó. Nếu trong lúc mình tính toán có writer khác đã commit trước, phép compare-and-swap thất bại → Iceberg kiểm xung đột (hai commit có đụng cùng file không?) và retry: nạp lại snapshot mới nhất, thẩm định lại, ghi metadata mới rồi thử commit lại.

Cơ chế này cho throughput tốt khi các writer chạm vùng dữ liệu khác nhau (mỗi job ghi partition ngày riêng — gần như không xung đột). Nhưng nếu nhiều job cùng ghi đè một partition nóng, retry storm có thể xảy ra (chỉnh qua commit.retry.num-retries, commit.retry.min-wait-ms…). Bài học vận hành: thiết kế phân vùng và lịch job để giảm chồng lấn ghi, đừng dựa vào retry để cứu.

Time travel: đọc bảng như nó từng là

Vì mọi snapshot cũ vẫn còn (đến khi bị expire), ta đọc được bảng tại một thời điểm quá khứ — gọi là time travel. Chọn snapshot theo snapshot-id (chính xác tuyệt đối) hoặc theo timestamp (Iceberg tìm snapshot còn hiệu lực tại mốc thời gian đó).

Cú pháp minh họa (Spark SQL):

-- Đọc theo phiên bản (snapshot-id) — MINH HOẠ, không chạy trong sandbox
SELECT * FROM db.loan_balance VERSION AS OF 3821550127947089009;

-- Đọc theo thời điểm
SELECT * FROM db.loan_balance TIMESTAMP AS OF '2026-06-30 23:59:59';

Trino dùng cú pháp FOR ...:

-- Trino — MINH HOẠ
SELECT * FROM iceberg.db.loan_balance FOR TIMESTAMP AS OF TIMESTAMP '2026-06-30 23:59:59 Asia/Ho_Chi_Minh';
SELECT * FROM iceberg.db.loan_balance FOR VERSION AS OF 3821550127947089009;

Lưu ý: các câu SQL trong bài là cú pháp Spark/Trino trên Iceberg, không phải PostgreSQL — nên chúng chỉ để minh họa, không đánh dấu chạy được trong sandbox.

Muốn biết có những snapshot nào để chọn, đọc bảng metadata historysnapshots:

-- MINH HOẠ (Spark)
SELECT committed_at, snapshot_id, operation, summary
FROM db.loan_balance.snapshots ORDER BY committed_at;

Cột operation cho biết mỗi snapshot sinh ra bởi thao tác gì (append, overwrite, delete, replace), còn summary chứa thống kê (số record thêm/xóa, số file). Đây là "sổ cái" thay đổi của bảng — cực hữu ích cho audit.

Ứng dụng điển hình của time travel: audit ("cột dư nợ khách này chiều qua là bao nhiêu?"), tái tạo báo cáo ("in lại đúng con số đã nộp ngày chốt kỳ"), và so sánh theo thời gian (chạy cùng query trên hai snapshot để tìm dữ liệu đã đổi).

Rollback & cherry-pick: sửa khi ghi sai

Nếu một job ghi hỏng (nạp nhầm file, transform sai logic) và đã commit, bạn không cần khôi phục từ backup — chỉ cần trỏ bảng về snapshot tốt trước đó. Đây là rollback, một thao tác metadata tức thời:

-- Spark — MINH HOẠ
CALL catalog.system.rollback_to_snapshot('db.loan_balance', 3821550127947089009);
CALL catalog.system.set_current_snapshot('db.loan_balance', 3821550127947089009);
CALL catalog.system.rollback_to_timestamp('db.loan_balance', TIMESTAMP '2026-06-30 20:00:00');
  • rollback_to_snapshot quay về một snapshot là tổ tiên của snapshot hiện tại.
  • set_current_snapshot trỏ tới bất kỳ snapshot-id nào (kể cả nhánh khác), linh hoạt hơn.
  • Cherry-pick (cherrypick_snapshot) áp thay đổi của một snapshot đã tồn tại (thường sinh ở branch) lên bảng chính mà không tái thực thi công việc — nền tảng cho quy trình write-audit-publish bên dưới.

Rollback bản thân nó cũng tạo một snapshot mới (con trỏ tiến lên, nội dung trỏ về file cũ), nên lịch sử vẫn liền mạch và có thể audit được — bạn thấy rõ "đã rollback lúc nào, về đâu".

Row-level: MERGE, UPDATE, DELETE

Lake thô gần như không sửa được một dòng — phải ghi đè cả partition. Iceberg hỗ trợ row-level operations thật sự: MERGE INTO (upsert), UPDATE, DELETE. Đây là điều kiện tiên quyết để làm CDC (change data capture), xử lý GDPR "right to be forgotten", hay hiệu chỉnh bút toán.

Ví dụ MERGE để hợp nhất dữ liệu thay đổi từ nguồn vào bảng đích:

-- Spark — MINH HOẠ (upsert dư nợ khoản vay)
MERGE INTO db.loan_balance t
USING staging.loan_cdc s
ON t.loan_id = s.loan_id
WHEN MATCHED AND s.op = 'D' THEN DELETE
WHEN MATCHED AND s.op = 'U' THEN UPDATE SET t.balance = s.balance, t.updated_at = s.ts
WHEN NOT MATCHED THEN INSERT (loan_id, balance, updated_at) VALUES (s.loan_id, s.balance, s.ts);

Copy-on-write vs merge-on-read

Câu hỏi cốt lõi: khi UPDATE/DELETE một số dòng nằm trong một data file lớn, Iceberg xử lý thế nào? Có hai chiến lược, đánh đổi trực tiếp giữa tốc độ ghitốc độ đọc:

Copy-on-write (CoW) — mặc định cho nhiều engine. Bất kỳ data file nào chứa dòng bị ảnh hưởng sẽ được đọc lại, sửa, rồi ghi lại toàn bộ thành file mới (không có dòng bị xóa, có dòng đã cập nhật). Snapshot mới trỏ tới file mới.

  • Ghi đắt: sửa 1 dòng có thể phải viết lại cả file trăm MB → write amplification lớn.
  • Đọc rẻ: file kết quả đã sạch, engine chỉ việc quét, không phải hợp nhất gì.
  • Hợp cho bảng ghi thưa, đọc nhiều (báo cáo, bảng chốt số).

Merge-on-read (MoR) — ghi thêm delete file đánh dấu dòng nào bị xóa/thay, KHÔNG viết lại data file. Lúc đọc, engine gộp data file với delete file để loại các dòng đã đánh dấu.

  • Ghi nhanh: chỉ viết delete file nhỏ → độ trễ thấp, hợp streaming/CDC tần suất cao.
  • Đọc đắt hơn: mỗi lần đọc phải áp delete file; delete file tích tụ làm chậm dần → cần compaction định kỳ (bài 6).

Cấu hình qua thuộc tính bảng, tách riêng theo từng loại thao tác: write.update.mode, write.delete.mode, write.merge.mode nhận copy-on-write hoặc merge-on-read.

Equality vs positional deletes (chỉ MoR)

Khi dùng merge-on-read, delete file có hai dạng:

  • Positional delete: ghi "trong file X, xóa dòng ở vị trí 42, 87…". Chính xác theo vị trí, hiệu quả khi biết rõ dòng nằm đâu (thường Spark/Flink batch tạo ra).
  • Equality delete: ghi "xóa mọi dòng có loan_id = 12345". Không cần biết vị trí — lý tưởng cho streaming/CDC khi upstream chỉ gửi khóa của bản ghi bị xóa. Đọc đắt hơn vì engine phải so khớp giá trị trên toàn bộ data file liên quan.

Streaming ingest (ví dụ Flink) thường sinh equality delete; nếu để tích tụ, đọc sẽ chậm dần — nên lịch compaction phải theo kịp nhịp ghi.

Branching & tagging: quy trình Write-Audit-Publish

Iceberg cho phép đặt tag (nhãn cố định vào một snapshot) và branch (một nhánh ghi độc lập). Đây là nền của mẫu Write-Audit-Publish (WAP) — cực hợp môi trường ngân hàng nơi dữ liệu sai không được lọt ra người dùng.

  • Tag: neo một snapshot làm mốc bất biến, ví dụ eom-2026-06 (cuối kỳ tháng 6). Về sau ... VERSION AS OF 'eom-2026-06' luôn trả đúng ảnh chốt số, kể cả khi bảng đã đi tiếp nhiều snapshot. Có thể đặt thời hạn giữ (retain) để phục vụ chính sách lưu trữ.
  • Branch: tạo audit branch, ghi dữ liệu mới vào đó, chạy kiểm tra chất lượng (data quality) — reconciliation, kiểm null, so tổng — trước khi hợp nhất vào bảng chính. Chỉ khi kiểm đạt mới fast-forward/cherry-pick branch lên bảng chính; nếu fail thì bỏ branch, người dùng cuối không bao giờ thấy dữ liệu lỗi.
-- Spark — MINH HOẠ
ALTER TABLE db.loan_balance CREATE TAG `eom-2026-06` RETAIN 3650 DAYS;
ALTER TABLE db.loan_balance CREATE BRANCH audit;
-- ... ghi & kiểm trên branch audit ...
CALL catalog.system.fast_forward('db.loan_balance', 'main', 'audit');

Quy trình WAP: Write (ghi vào branch) → Audit (chạy kiểm chất lượng trên branch) → Publish (fast-forward/cherry-pick lên main nếu đạt). Người tiêu thụ trên main luôn thấy dữ liệu đã được kiểm.

Use case thực tế

Bối cảnh (số liệu ước lượng minh họa). NCB đưa bảng loan_balance (dư nợ từng khoản vay, ~4,2 triệu dòng, cập nhật hằng ngày) lên Iceberg trên object storage, truy vấn bằng Spark cho ETL và Trino cho báo cáo. Hai bài toán tuân thủ:

1) Tái tạo báo cáo NHNN đúng thời điểm chốt số. Cuối mỗi kỳ, đội báo cáo chốt số dư nợ và nhóm nợ để nộp Ngân hàng Nhà nước. Ba tuần sau, thanh tra hỏi lại "con số dư nợ nhóm 3 tại 23:59:59 ngày 30/06 lấy ở đâu ra?". Trước đây, bảng đã bị hàng chục job ghi đè, không cách nào dựng lại. Với Iceberg, ngay khi chốt số đội đặt một tag eom-2026-06. Khi thanh tra hỏi, chỉ cần:

-- MINH HOẠ (Trino) — dựng lại đúng ảnh chốt số
SELECT nhom_no, COUNT(*) so_khoan, SUM(balance) tong_du_no
FROM iceberg.risk.loan_balance FOR VERSION AS OF 'eom-2026-06'
GROUP BY nhom_no;

Con số khớp tuyệt đối với báo cáo đã nộp, kèm bằng chứng snapshot-id và committed_at trong bảng history — audit trail không thể chối cãi. So sánh chênh lệch giữa hai kỳ cũng chỉ là chạy cùng query trên hai tag rồi EXCEPT.

2) Write-Audit-Publish trước khi publish số liệu. Job nạp dư nợ hằng đêm từng vài lần đẩy nhầm dữ liệu (nguồn CDC lỗi khiến ~1,8% khoản vay có balance âm), và vì ghi thẳng vào bảng chính nên dashboard rủi ro sáng hôm sau hiển thị sai, phải chữa cháy. Đội chuyển sang WAP: mỗi đêm job ghi vào branch audit, rồi chạy bộ kiểm — không có balance âm, tổng dư nợ lệch không quá 2% so với hôm trước, tỉ lệ null khóa = 0. Chỉ khi tất cả đạt, pipeline fast_forward branch lên main; nếu fail, branch bị bỏ và cảnh báo gửi trực ban, main vẫn giữ số liệu tốt hôm trước. Sau khi áp WAP, số sự cố "dữ liệu rủi ro sai lọt ra dashboard" giảm về gần 0 trong quý.

3) Sửa nhanh khi lỡ ghi sai. Một lần transform nhóm nợ chạy nhầm tham số và đã commit lên main. Thay vì phục hồi từ backup (ước tính vài giờ), trực ca chạy rollback_to_snapshot về snapshot trước đó — bảng trở lại trạng thái đúng trong vài giây, dashboard tự khỏi ở lần refresh kế tiếp.

Về chiến lược ghi: bảng loan_balance phục vụ đọc/báo cáo nặng nên để copy-on-write cho UPDATE/DELETE (đọc nhanh, ghi thưa mỗi đêm chấp nhận được). Ngược lại, bảng txn_stream nhận CDC giao dịch gần real-time đặt merge-on-read với equality delete để ghi độ trễ thấp, và lịch compaction chạy mỗi 2 giờ để delete file không tồn đọng làm chậm truy vấn.

Ghi nhớ

  • ACID của Iceberg = snapshot + atomic commit ở catalog. Ghi không sửa file cũ; công bố thay đổi = đổi con trỏ metadata nguyên tử (compare-and-swap). Người đọc luôn thấy một snapshot hoàn chỉnh.
  • Optimistic concurrency: nhiều writer cùng ghi, không khóa trước; commit kiểm xung đột và retry nếu base đã cũ. Thiết kế partition/lịch job để giảm chồng lấn, đừng dựa vào retry.
  • Isolation: snapshot isolation (mặc định) vs serializable (chặt hơn, chặn phantom, xung đột nhiều hơn) — cấu hình theo từng loại ghi.
  • Time travel đọc bảng quá khứ theo VERSION AS OF (snapshot-id) hoặc TIMESTAMP AS OF; Trino dùng FOR VERSION/TIMESTAMP AS OF. Dùng để audit, tái tạo báo cáo, so sánh theo thời gian. Bảng history/snapshots là sổ cái thay đổi.
  • Rollback (rollback_to_snapshot/set_current_snapshot) trỏ bảng về snapshot tốt — sửa lỗi ghi trong vài giây, không cần restore backup.
  • Row-level: MERGE/UPDATE/DELETE thật sự. Copy-on-write viết lại data file (ghi chậm, đọc nhanh) vs merge-on-read ghi delete file (ghi nhanh, đọc chậm, cần compaction). MoR có positional delete (theo vị trí) và equality delete (theo khóa, hợp streaming/CDC).
  • Branching & tagging cho phép Write-Audit-Publish: ghi vào branch, kiểm chất lượng, chỉ publish lên main khi đạt; tag neo mốc cuối kỳ bất biến để dựng lại đúng số đã nộp.
  • Xem thêm: tổng quan Iceberg, kiến trúc metadata, bảo trì & hiệu năng.

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ẻ!