Batch 4 — Idempotency, Retry & Exactly-once
Mô hình tinh thần: job sẽ chạy lại — hãy thiết kế cho điều đó
Trong thế giới batch, câu hỏi không phải "job có bao giờ chạy lại không" mà là "khi nó chạy lại thì kết quả có sai không". Một job sẽ chạy lại: scheduler retry vì timeout, kỹ sư trực đêm bấm re-run vì nghi ngờ dữ liệu, backfill quét lại 30 ngày quá khứ, hoặc container bị OOM-kill giữa chừng rồi Kubernetes khởi động lại pod. Nếu mỗi lần chạy lại cộng thêm một bản ghi trùng, tổng dư nợ EOD sẽ phình lên, báo cáo NHNN sai, và không ai dám bấm nút re-run nữa.
Idempotency là tính chất: chạy một thao tác một lần hay nhiều lần đều cho cùng một trạng thái cuối. Về mặt toán học, f(f(x)) = f(x). Một pipeline idempotent biến việc chạy lại từ rủi ro thành thao tác an toàn, buồn tẻ — và đó chính là mục tiêu.
Nguyên tắc nền: trong hệ phân tán, lỗi là chuyện thường ngày, nên retry là bắt buộc; mà retry chỉ an toàn khi thao tác idempotent. Idempotency không phải tính năng xa xỉ — nó là điều kiện để retry tồn tại.
Phân biệt với hai khái niệm hay bị lẫn:
| Khái niệm | Ý nghĩa |
|---|---|
| Idempotent | Chạy lại N lần = chạy 1 lần (về trạng thái cuối) |
| Deterministic (tất định) | Cùng input → cùng output, không phụ thuộc thời điểm/thứ tự/ngẫu nhiên |
| Commutative (giao hoán) | Thứ tự áp dụng không đổi kết quả |
Idempotency thường cần tính tất định làm nền: nếu transform dùng now() hay rand(), hai lần chạy cho hai kết quả khác nhau, không thể idempotent.
Ba mức đảm bảo phân phối: at-least-once, at-most-once, exactly-once
Khi một job đọc nguồn, xử lý, rồi ghi đích, giữa các bước luôn có khả năng crash. Tùy cách bố trí "xử lý" và "đánh dấu đã xong" (acknowledge/commit offset), ta rơi vào một trong ba mức:
- At-most-once: đánh dấu xong trước khi xử lý. Crash giữa chừng → bản ghi mất luôn, không bao giờ trùng. Chấp nhận được với metric best-effort, không bao giờ chấp nhận với giao dịch tiền.
- At-least-once: xử lý xong mới đánh dấu. Crash sau khi ghi, trước khi ack → lần retry ghi lại → trùng. Đây là mức phổ biến và rẻ nhất để đạt được. Hầu hết hệ thống thực tế cho at-least-once ở tầng vận chuyển.
- Exactly-once: mỗi bản ghi tác động đúng một lần lên kết quả cuối. Đây là mức khó nhất.
Điểm cốt lõi thường bị hiểu sai: "exactly-once delivery" (phân phối đúng một lần qua mạng) gần như bất khả thi trong hệ phân tán — bài toán Two Generals. Cái ta thực sự đạt được và cần là "exactly-once processing / effect": nguồn có thể gửi một bản ghi nhiều lần (at-least-once), nhưng hiệu ứng lên trạng thái cuối chỉ xảy ra một lần.
Công thức thực dụng: at-least-once delivery + xử lý idempotent = exactly-once effect. Bạn không chống được việc dữ liệu tới hai lần; bạn làm cho lần thứ hai không thay đổi gì. Đây là cách các sink batch đạt exactly-once trên thực tế.
Vì sao "chỉ INSERT" là cái bẫy — minh hoạ nhân đôi
Xét job EOD nạp giao dịch ngày 2026-07-20 vào bảng fct_transactions. Cách viết ngây thơ: INSERT INTO fct_transactions SELECT ... WHERE txn_date = '2026-07-20'. Chạy một lần thì đúng. Nhưng scheduler timeout ở phút thứ 9 (job cần 10 phút), retry lần 2 chạy lại nguyên vẹn:
Với INSERT thuần, retry cộng dồn. Với overwrite theo phạm vi bạn kiểm soát (partition của ngày đó), retry thay thế thay vì cộng thêm — nên chạy 1 lần hay 5 lần, partition vẫn chứa đúng ảnh chụp mới nhất. Đây là kỹ thuật idempotent nền tảng nhất của batch.
Bộ công cụ làm job idempotent
Không có một "viên đạn bạc" — có một hộp công cụ, chọn theo hình dạng dữ liệu.
1. Overwrite theo partition (insert-overwrite)
Chia bảng theo cột thời gian (thường date) và mỗi lần chạy ghi đè trọn partition của phạm vi đó thay vì append. Đây là cách phổ biến nhất cho ETL theo ngày. Chi tiết cơ chế partition và I/O xem Batch 6 — Partitioning & I/O.
-- Spark SQL / Hive: thay TOÀN BỘ partition ngày đó, các ngày khác không đụng
INSERT OVERWRITE TABLE fct_transactions PARTITION (txn_date = '2026-07-20')
SELECT txn_id, account_id, amount, channel
FROM staging_transactions
WHERE txn_date = '2026-07-20';
Lưu ý cấu hình quan trọng trên Spark: chế độ ghi đè partition động (partitionOverwriteMode = dynamic) chỉ thay các partition có mặt trong dữ liệu ra, không xoá nhầm ngày khác. Với mode="overwrite") mặc định (static) trên Spark cũ, một cú ghi có thể xoá cả bảng — nhớ kiểm tra kỹ. (Về mặt nguyên lý: phạm vi ghi đè phải khớp đúng phạm vi bạn tính lại.)
# PySpark: ghi đè động, chỉ đụng các partition có trong df_out
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
(df_out
.write
.mode("overwrite") # thay, không append
.partitionBy("txn_date")
.format("parquet")
.saveAsTable("fct_transactions"))
2. MERGE (upsert) thay INSERT
Khi không thể chia sạch theo partition — ví dụ cập nhật một tập key rải rác (trạng thái khoản vay đổi ở nhiều ngày mở khác nhau) — dùng MERGE: khớp theo khóa nghiệp vụ, có thì UPDATE, chưa có thì INSERT. Chạy lại lần hai không tạo bản trùng vì key đã tồn tại → rơi vào nhánh UPDATE với cùng giá trị. MERGE cũng là xương sống của incremental/CDC, xem Batch 3 — Incremental & CDC.
MERGE INTO fct_loan_status AS t
USING staging_loan_status AS s
ON t.loan_id = s.loan_id
WHEN MATCHED THEN UPDATE SET
t.status = s.status, t.dpd = s.dpd, t.updated_at = s.updated_at
WHEN NOT MATCHED THEN INSERT
(loan_id, status, dpd, updated_at)
VALUES (s.loan_id, s.status, s.dpd, s.updated_at);
MERGE idempotent nếu nguồn s đã dedup theo loan_id (một key một dòng) — nếu không, chuẩn SQL báo lỗi hoặc cho kết quả không tất định. Đó là lý do bước dedup dưới đây thường đứng trước MERGE.
3. Dedup theo key (khử trùng chủ động)
Khi at-least-once cho phép nguồn lặp bản ghi, khử trùng bằng khóa duy nhất trước khi ghi. Giữ một đại diện theo tiêu chí tất định (bản mới nhất theo timestamp, tie-break bằng một cột ổn định):
WITH ranked AS (
SELECT *,
ROW_NUMBER() OVER (
PARTITION BY txn_id -- khóa nghiệp vụ
ORDER BY ingested_at DESC, source_file ASC -- tie-break tất định
) AS rn
FROM staging_transactions
)
SELECT * FROM ranked WHERE rn = 1;
Tie-break phải tất định: nếu chỉ ORDER BY ingested_at mà hai bản cùng timestamp, mỗi lần chạy có thể giữ bản khác nhau → mất tính idempotent một cách tinh vi.
4. Transform tất định (deterministic)
Idempotency sụp đổ nếu logic phụ thuộc yếu tố thay đổi giữa các lần chạy. Tránh:
now(),current_timestampgiữa thân transform → dùng tham số ngày logic (logical_date/dscủa job) truyền vào, không dùng đồng hồ thực.rand(), UUID ngẫu nhiên làm khóa → dùng surrogate key tất định (hash của khóa tự nhiên).- Phụ thuộc thứ tự dòng không có
ORDER BYrõ ràng. - Đọc "trạng thái hiện tại của bảng đích rồi cộng thêm" (
balance = balance + x) — thao tác này không idempotent; hãy tính lại từ nguồn bất biến.
5. Atomic / transactional write (all-or-nothing)
Vấn đề độc hại nhất là ghi dở: job chết ở giữa, để lại nửa partition. Retry sẽ append lên phần dở → hỗn loạn. Giải pháp: khiến kết quả hoặc hiện đủ, hoặc không hiện gì (nguyên tử).
- Ghi tạm rồi đổi tên/hoán đổi: ghi ra thư mục/bảng staging, thành công thì
swap/renamenguyên tử. Con đọc không bao giờ thấy trạng thái dở. - Bảng transaction (ACID table) như Delta Lake, Apache Iceberg, Apache Hudi: mỗi lần ghi là một commit nguyên tử; job chết giữa chừng thì snapshot cũ vẫn nguyên, không có bản ghi rác. Đây là lý do lakehouse format lên ngôi cho batch tin cậy.
Retry và poison pill
Idempotency mở khóa cho retry, nhưng retry cần kỷ luật:
- Backoff (giãn cách tăng dần) + jitter: đừng retry tức thì dồn dập — làm ngập nguồn đang hồi phục (thundering herd). Tăng khoảng chờ theo cấp số nhân, thêm nhiễu ngẫu nhiên.
- Giới hạn số lần: retry vô hạn che giấu lỗi thật.
- Poison pill (viên thuốc độc): một bản ghi/đơn vị công việc luôn làm job chết (schema hỏng, chia cho 0, ký tự lạ). Retry nó vô nghĩa — nó sẽ giết mọi lần chạy. Cần: sau N lần, tách nó sang Dead Letter Queue / bảng quarantine để phần còn lại chạy tiếp, rồi điều tra riêng. Không cô lập poison pill → một dòng xấu chặn cả EOD.
- Phân biệt lỗi tạm thời và lỗi vĩnh viễn: timeout mạng thì retry; lỗi logic/validation thì retry vô ích, nên fail nhanh và báo động.
Cấu hình retry/backoff và phân loại lỗi thuộc tầng điều phối — xem Batch 7 — Orchestration về retry ở mức DAG, sensor và alerting.
Idempotency key — chốt chặn cuối
Khi bản thân thao tác khó làm idempotent (gọi API bên ngoài, ghi vào hệ không hỗ trợ MERGE), dùng idempotency key: mỗi đơn vị công việc mang một khóa duy nhất, ổn định; đích ghi nhận key đã xử lý và bỏ qua nếu gặp lại.
- Key phải tất định theo nội dung công việc, không phải ngẫu nhiên mỗi lần gọi — nếu không, retry sinh key mới và mất tác dụng chống trùng. Ví dụ:
hash(txn_id + logical_date)chứ không phảiuuid4(). - Đích lưu bảng
processed_keys(hoặc dùng ràng buộc UNIQUE): trước khi áp dụng, kiểm tra key; đã có thì skip. Ràng buộcUNIQUE(txn_id)ở tầng DB là một idempotency key thụ động — lần chèn trùng thứ hai bị từ chối. - Đây chính là cơ chế các cổng thanh toán dùng để bạn bấm "Trả" hai lần mà không bị trừ tiền hai lần.
-- Ràng buộc UNIQUE làm idempotency key thụ động
ALTER TABLE fct_transactions
ADD CONSTRAINT uq_txn UNIQUE (txn_id);
-- chèn lại txn_id đã có sẽ bị chặn / ON CONFLICT DO NOTHING (Postgres)
INSERT INTO fct_transactions (txn_id, account_id, amount)
VALUES ('T-9001', 'A-55', 1500000)
ON CONFLICT (txn_id) DO NOTHING;
Use case thực tế
Bối cảnh (số liệu minh hoạ): Job EOD của NCB nạp ~1,2 triệu giao dịch/ngày vào fct_transactions rồi tính bảng phân loại nợ. Một đêm, mạng tới core banking chập chờn khiến job vượt SLA 10 phút và bị Airflow retry.
- Trước: job dùng
INSERT ... SELECT. Retry chạy lại trọn ngày → partition2026-07-20có ~2,4 triệu dòng, tổng dư nợ EOD gấp đôi (~sai lệch hàng nghìn tỷ trên báo cáo nội bộ). Đội vận hành phảiDELETEthủ công lúc 2h sáng, dò tay bản trùng — rủi ro xoá nhầm. - Sau: chuyển sang
INSERT OVERWRITE PARTITION (txn_date=...)+ dedupROW_NUMBERtheotxn_id, ghi trên bảng Iceberg (commit nguyên tử). Cùng sự cố mạng, Airflow retry 2 lần: mỗi lần thay trọn partition bằng ảnh chụp mới nhất. Kết quả cuối: đúng ~1,2 triệu dòng, dù chạy 1 hay 3 lần. Đội trực không cần can thiệp; nút "Clear & re-run" trở thành thao tác an toàn ai cũng dám bấm. - Poison pill: một file nguồn có 3 dòng
amountrỗng làm bước ép kiểu chết. Thay vì fail cả EOD, bước validate tách 3 dòng sang bảngquarantine_transactions, EOD hoàn tất đúng giờ, sáng hôm sau phân tích 3 dòng lỗi riêng.
Ghi nhớ
- Giả định job sẽ chạy lại (retry, backfill, re-run tay, pod restart). Thiết kế để lần chạy thứ N cho cùng trạng thái cuối như lần đầu.
- At-least-once delivery + xử lý idempotent = exactly-once effect. Đừng đuổi theo "exactly-once delivery" — hãy làm lần lặp thứ hai không thay đổi gì.
- Công cụ idempotent, chọn theo hình dạng dữ liệu: insert-overwrite partition (theo ngày), MERGE/upsert (key rải rác), dedup ROW_NUMBER (khử trùng nguồn), transform tất định, atomic/transactional write.
INSERTthuần là bẫy nhân đôi. Overwrite thay thế phạm vi bạn kiểm soát; MERGE khớp key rồi update; cả hai chống trùng khi retry.- Transform phải tất định: tránh
now()/rand()/UUID ngẫu nhiên/x = x + delta; dùng ngày logic và surrogate key hash. - Atomic write (rename staging, hay ACID table Delta/Iceberg/Hudi) diệt vấn đề ghi-dở — con đọc không thấy trạng thái nửa vời.
- Retry cần backoff + jitter + giới hạn; cô lập poison pill sang DLQ/quarantine để một dòng xấu không chặn cả pipeline.
- Idempotency key tất định (hash nội dung, hoặc ràng buộc UNIQUE) là chốt chặn cuối khi thao tác không tự idempotent.
Nguồn tham khảo
- Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly, 2022): chương về idempotency, reproducibility và độ tin cậy pipeline.
- Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly, 2017): các mức đảm bảo, exactly-once, dedup và bài toán hai vị tướng.
- Streaming Systems — Tyler Akidau, Slava Chernyak, Reuven Lax (O'Reilly, 2018): exactly-once, idempotent sink, two-phase commit (áp dụng chung cho batch).
- Apache Spark SQL docs —
INSERT OVERWRITE,partitionOverwriteMode(dynamic/static): https://spark.apache.org/docs/latest/sql-ref-syntax-dml-insert-overwrite-table.html - Delta Lake documentation — ACID transactions,
MERGE INTO, atomic commit: https://docs.delta.io/latest/delta-update.html - Apache Iceberg documentation — snapshot, atomic commit, row-level operations: https://iceberg.apache.org/docs/latest/
- "Questioning the Lambda Architecture" — Jay Kreps (O'Reilly Radar, 2014): reprocessing và tính chạy-lại-được của pipeline.
- Stripe API — Idempotent Requests: https://docs.stripe.com/api/idempotent_requests
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.
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.
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.
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.
Cảm nhận của bạn
Bình luận
Chưa có bình luận. Hãy là người đầu tiên chia sẻ!