Mục tiêu là pipeline idempotent: chạy một lần hay mười lần với cùng khoảng thời gian thì kết quả vẫn như nhau. Trùng dữ liệu xảy ra vì job dùng INSERT append, chạy lại là chèn thêm lần nữa.
Các cách làm:
- Gắn mỗi lần chạy với một khoảng thời gian cố định (ví dụ ngày 2026-09-24), không dùng "bây giờ". Airflow truyền logical date và data interval cho mỗi lần chạy; dùng now() trong task thì chạy lại hôm sau sẽ lấy dữ liệu khác.
- Ghi đè đúng partition thay vì append: xoá rồi ghi lại partition của ngày đó trong cùng transaction, hoặc dùng INSERT OVERWRITE trên Spark/Hive.
- Upsert theo khoá (MERGE) khi dữ liệu không chia partition theo ngày được.
- Ghi vào bảng tạm rồi đổi tên/swap để người đọc không bao giờ thấy dữ liệu dở dang.
BEGIN;
DELETE FROM fact_orders WHERE order_date = :run_date;
INSERT INTO fact_orders
SELECT * FROM stg_orders WHERE order_date = :run_date;
COMMIT;Backfill: khi logic đổi hoặc job hỏng nhiều ngày, chỉ cần chạy lại từng khoảng ngày cần sửa (backfill hoặc clear task trên Airflow). Vì mỗi lần chạy độc lập theo ngày, có thể chạy song song vài ngày một lúc, miễn là giới hạn để không làm quá tải nguồn.
Lưu ý: nguồn cũng phải đọc lại được. Nếu extract đọc từ API chỉ trả dữ liệu hiện tại, backfill sẽ ra số khác. Vì vậy nên lưu bản thô của mỗi lần extract theo ngày, rồi backfill từ bản thô đó.