Batch 5 — Backfill & Tái xử lý lịch sử
Batch 5 — Backfill & Tái xử lý lịch sử
Mô hình tinh thần
Một pipeline batch không chỉ chạy cho hôm nay. Sớm hay muộn bạn sẽ phải trả lời câu hỏi: "Dữ liệu 6 tháng vừa rồi cần tính lại thì làm thế nào?". Đó là backfill — chạy lại (một phần hoặc toàn bộ) pipeline cho một khoảng thời gian trong quá khứ.
Ba tình huống điển hình buộc phải backfill:
| Tình huống | Ví dụ ngân hàng |
|---|---|
| Đổi logic nghiệp vụ | Đổi công thức phân loại nợ (nhóm 1–5) theo thông tư mới → phải tính lại lịch sử phân loại |
| Sửa lỗi (bug fix) | Job tính phí sao kê nhầm múi giờ → mọi bảng sao kê 3 tháng qua sai giờ cắt EOD |
| Thêm cột / thêm nguồn | Bảng fact_transaction bổ sung cột channel (ATM/POS/IB) → cột mới rỗng cho toàn bộ dữ liệu cũ |
Mô hình tinh thần cốt lõi: coi thời gian là một trục có thể chia nhỏ (partition), và mỗi partition là một đơn vị xử lý độc lập, có thể chạy lại mà không làm hỏng phần còn lại. Nếu pipeline của bạn đạt được tính chất đó, backfill chỉ là "lặp qua các ngày cần tính lại". Nếu không, backfill trở thành một dự án đầy rủi ro.
Backfill là bài kiểm tra khắc nghiệt nhất cho hai tính chất đã bàn ở Batch 4 — Idempotency và Batch 6 — Partitioning & I/O. Pipeline không idempotent + không partition tốt gần như không thể backfill an toàn.
Vì sao backfill khó
Backfill khác một lần chạy bình thường ở quy mô và bối cảnh:
- Khối lượng lớn, dồn cục: một job hằng ngày xử lý 1 ngày dữ liệu; backfill 12 tháng = 365 lần khối lượng đó, chạy trong cửa sổ ngắn. Dễ làm nghẽn cluster, cạn quota kho lưu trữ, hết slot cảnh báo.
- Phụ thuộc downstream: tính lại
fact_transactionlà chưa đủ — mọi bảng aggregate, mart, báo cáo, mô hình ML đọc từ nó cũng phải tính lại theo. Nếu không, downstream giữ số cũ, tạo mâu thuẫn. - Tránh ảnh hưởng job production đang chạy: backfill nặng không được tranh tài nguyên với job EOD hằng ngày, và tuyệt đối không được ghi đè nhầm lên partition hôm nay đang được job thường tạo ra.
- Nhất quán tại ranh giới: chỗ nối giữa "vùng đã backfill (logic mới)" và "vùng chưa đụng (logic cũ)" phải rõ ràng, tránh tình trạng nửa cũ nửa mới trong cùng một bảng mà người dùng không biết.
- Không idempotent = nhân đôi: nếu chạy lại một partition mà pipeline append thay vì overwrite, backfill sẽ nhân bản dữ liệu.
Bốn đòn bẩy để khống chế bốn khó khăn trên là: partition theo thời gian, idempotent overwrite, versioning logic, và shadow/blue-green table. Ta đi lần lượt.
Kỹ thuật 1 — Partition theo ngày & backfill song song
Nền tảng của backfill an toàn là phân vùng dữ liệu theo cột thời gian (thường là ngày nghiệp vụ, ví dụ event_date hoặc partition_date). Khi đó mỗi ngày là một partition độc lập trên storage:
s3://lake/silver/fact_transaction/partition_date=2026-01-01/
s3://lake/silver/fact_transaction/partition_date=2026-01-02/
...
Backfill trở thành: với mỗi ngày trong khoảng cần tính lại, đọc input của đúng ngày đó → tính → ghi đè đúng partition đó. Vì các partition độc lập, ta có thể chạy song song nhiều ngày cùng lúc (giới hạn theo tài nguyên).
Vài nguyên tắc thực chiến:
- Giới hạn mức song song (concurrency): đặt một pool/limit để backfill không nuốt trọn cluster. Thà chậm mà không đụng job hằng ngày.
- Chia lô (batching): với khoảng rất dài, chia thành lô theo tuần/tháng, chạy tuần tự từng lô để dễ kiểm soát và dừng giữa chừng.
- Ưu tiên & tách tài nguyên: chạy backfill trên queue/cluster riêng, hoặc lịch giờ thấp điểm (đêm/cuối tuần), tách khỏi luồng EOD.
- Theo dõi tiến độ theo partition: lưu trạng thái từng ngày (pending/done/failed) để có thể chạy lại chỉ những ngày lỗi.
Kỹ thuật 2 — Idempotent overwrite
Song song chỉ an toàn khi mỗi task idempotent: chạy 1 lần hay 5 lần trên cùng partition đều cho kết quả y hệt, không nhân đôi. Chìa khoá là overwrite theo partition thay vì append mù.
Với Spark, mẫu phổ biến là ghi đè động chỉ những partition được đụng tới:
# Chỉ ghi đè các partition có trong dataframe, giữ nguyên partition khác
spark.conf.set("spark.sql.sources.partitionOverwriteMode", "dynamic")
(df_recomputed
.write
.mode("overwrite") # overwrite thay vì append
.partitionBy("partition_date") # đúng khoá partition
.format("parquet")
.save("s3://lake/silver/fact_transaction"))
Với bảng SQL/warehouse, mẫu tương đương là DELETE + INSERT trong một transaction cho đúng khoảng ngày (idempotent vì lần chạy sau xoá sạch rồi ghi lại):
BEGIN;
DELETE FROM silver.fact_transaction
WHERE partition_date BETWEEN DATE '2026-01-01' AND DATE '2026-01-31';
INSERT INTO silver.fact_transaction
SELECT ... FROM staging.transaction_raw
WHERE partition_date BETWEEN DATE '2026-01-01' AND DATE '2026-01-31';
COMMIT;
Với các định dạng bảng lakehouse (Delta Lake, Apache Iceberg), câu lệnh MERGE/REPLACE WHERE/INSERT OVERWRITE cho ghi đè partition nguyên tử (atomic) — người đọc không bao giờ thấy trạng thái nửa vời. Chi tiết cơ chế idempotent xem Batch 4.
Bẫy hay gặp:
mode("overwrite")không cópartitionOverwriteMode=dynamic(hoặc chưa lọc đúng ngày) sẽ xoá cả bảng. Luôn kiểm thử trên môi trường staging trước khi backfill production.
Kỹ thuật 3 — Versioning logic
Khi backfill là do đổi logic, câu hỏi khó là: "làm sao biết partition nào đã dùng logic mới, partition nào còn logic cũ?". Câu trả lời: gắn phiên bản logic vào chính dữ liệu.
- Thêm cột metadata như
logic_version(ví dụv2) hoặcpipeline_run_id,processed_atvào mỗi bản ghi/partition. Sau backfill, có thể truy vấn xem còn partition nào ởv1. - Quản lý code theo Git tag/commit; ghi lại commit nào tạo ra partition nào để tái lập (reproducibility). Đây là nguyên tắc của DataOps — version hoá.
- Với thay đổi lớn, cân nhắc giữ logic cũ song song một thời gian để đối chiếu (reconcile) số cũ vs số mới trước khi "chốt".
Cột phiên bản còn giúp backfill có thể tạm dừng và tiếp tục: chỉ tính lại những partition chưa đạt phiên bản mục tiêu.
Kỹ thuật 4 — Shadow table / Blue-green khi đổi logic
Ghi đè trực tiếp (in-place) lên bảng production trong lúc backfill có rủi ro: người dùng đọc bảng ngay giữa chừng sẽ thấy dữ liệu nửa cũ nửa mới. Với thay đổi logic lớn, kỹ thuật an toàn là tính ra một bảng bóng (shadow) rồi hoán đổi (blue-green):
- Tạo bảng/thư mục mới
fact_transaction__v2(green), giữ nguyên bảng đang phục vụfact_transaction(blue). - Backfill toàn bộ lịch sử bằng logic mới vào bảng green — mất bao lâu cũng được, không ai đọc nó.
- Kiểm định chất lượng trên green: đếm dòng, đối chiếu tổng tiền, so sai khác với blue (xem Data Quality).
- Khi đạt, hoán đổi con trỏ: đổi tên bảng, hoặc trỏ view/alias
fact_transactionsang green. Thao tác này gần như tức thời và nguyên tử. - Giữ blue một thời gian để rollback nếu phát hiện lỗi.
Chi phí là gấp đôi storage tạm thời và một bước hoán đổi — đổi lại bạn có backfill không downtime, có QC trước khi phơi số, và có đường lùi. Đây là lựa chọn mặc định cho các thay đổi logic ảnh hưởng báo cáo tài chính.
Với thay đổi nhỏ (thêm 1 cột, sửa 1 bug cục bộ), overwrite in-place theo partition (Kỹ thuật 2) thường là đủ và rẻ hơn. Blue-green dành cho thay đổi sâu, rộng, nhạy cảm.
Backfill trong Airflow
Airflow mô hình hoá thời gian qua data interval của mỗi DAG run, nên rất hợp để backfill theo partition ngày. Có hai cơ chế cần phân biệt rõ.
Catchup (bù lịch tự động)
Khi bật catchup=True, nếu start_date ở quá khứ, scheduler sẽ tự tạo và chạy mọi DAG run cho các interval đã trôi qua giữa start_date và hiện tại. Đây thực chất là backfill tự động khi bạn bật một DAG mới hoặc DAG bị tạm dừng lâu.
from airflow import DAG
import pendulum
with DAG(
dag_id="fact_transaction_daily",
start_date=pendulum.datetime(2026, 1, 1, tz="Asia/Ho_Chi_Minh"),
schedule="@daily",
catchup=False, # thường đặt False để tránh "bão" run khi deploy
max_active_runs=3, # chặn số run đồng thời -> bảo vệ cluster
) as dag:
...
Thực tế đa số team đặt catchup=False để tránh scheduler bất ngờ khởi động hàng trăm run lịch sử khi triển khai DAG. Khi cần bù lịch, họ dùng backfill có chủ đích (dưới đây). Chi tiết cơ chế lập lịch & data interval xem Batch 7 — Orchestration & Catchup.
Backfill có chủ đích (command)
Khi muốn tính lại một khoảng cụ thể, chạy backfill qua CLI với mốc đầu–cuối:
# Chạy lại DAG cho từng interval trong khoảng ngày đã chọn
airflow dags backfill fact_transaction_daily \
--start-date 2026-01-01 \
--end-date 2026-01-31
Mỗi ngày trong khoảng thành một DAG run với data_interval tương ứng; task đọc/ghi đúng partition ngày đó. Kết hợp với max_active_runs (hoặc pool) để giới hạn song song, ta có đúng mô hình Kỹ thuật 1 nhưng do scheduler điều phối.
Điều kiện để backfill Airflow an toàn:
- Task idempotent theo interval: mỗi task chỉ đọc/ghi partition của interval của nó (dùng biến logic như
data_interval_start), không dựa vào "ngày hệ thống". - Không phụ thuộc trạng thái ngoài luồng: tránh đọc "bảng hôm nay" cứng, tránh side-effect gửi email/thông báo khi chạy lại lịch sử (dùng cờ để tắt notify trong chế độ backfill).
- Downstream nối bằng dependency/Asset: để khi partition upstream được tính lại, các job aggregate/mart phụ thuộc cũng được kích hoạt tính lại theo, tránh số lệch.
Reprocessing khi schema hoặc logic đổi
Backfill vì thêm cột/đổi schema có vài điểm riêng:
- Thêm cột nullable: rẻ nhất — cột mới để
NULL/mặc định cho dữ liệu cũ, chỉ backfill nếu nghiệp vụ bắt buộc có giá trị lịch sử. Nhiều khi "từ nay có, quá khứ để trống" là chấp nhận được. - Backfill giá trị cột mới: nếu cần điền cột
channelcho lịch sử, phải còn nguồn thô (raw/bronze) để suy ra. Đây là lý do kiến trúc medallion giữ lớp bronze bất biến: bronze là "bộ nhớ" để reprocess bất cứ lúc nào. - Đổi kiểu / đổi ý nghĩa cột (breaking): nên đi đường shadow table + version thay vì sửa tại chỗ, vì downstream có thể vỡ. Phối hợp schema evolution của định dạng bảng (Iceberg/Delta) để thêm cột an toàn.
- Đổi logic có hiệu lực theo mốc thời gian: nếu thông tư mới chỉ áp dụng từ ngày X, thì backfill chỉ vùng ≥ X, và giữ nguyên vùng < X (nguyên tắc bi-temporal: phân biệt "thời điểm sự kiện" và "thời điểm biết/áp dụng").
Nguyên tắc bao trùm: giữ lớp thô bất biến + pipeline có thể tái chạy từ thô = năng lực reprocess vô hạn. Không có bronze, mọi backfill đổi logic đều là canh bạc.
Use case thực tế
Bối cảnh (minh hoạ, số liệu giả định). NCB ban hành cách tính lại nhóm nợ theo quy định mới, áp dụng cho toàn bộ dư nợ từ 2026-01-01. Bảng silver.loan_classification phân vùng theo partition_date, ~90 ngày cần tính lại, mỗi ngày ~2 triệu bản ghi khoản vay.
Cách làm:
- Xây bảng shadow
loan_classification__v2với logic mới; giữ bảng blue đang phục vụ báo cáo. - Backfill song song qua Airflow:
airflow dags backfill loan_classification_daily --start-date 2026-01-01 --end-date 2026-03-31, đặtmax_active_runs=5để không đụng job EOD; mỗi run overwrite đúng partition ngày (idempotent), ghi kèmlogic_version='v2'. - QC trên shadow: đối chiếu phân bố nhóm 1–5 giữa v1 và v2 theo ngày; số khoản "nhảy nhóm" phải khớp giải trình nghiệp vụ. Giả sử phát hiện lệch ~1,8% khoản chuyển nhóm 2→3 — đúng như dự kiến của quy định mới.
- Hoán đổi view
loan_classificationsang v2 (atomic), kích hoạt tính lại các mart/báo cáo phân loại nợ downstream, giữ v1 thêm 14 ngày để rollback.
Kết quả (minh hoạ): backfill 90 ngày hoàn tất trong ~1 cửa sổ cuối tuần, không gián đoạn báo cáo hằng ngày, có số đối chiếu trước khi công bố, và có đường lùi. Nếu làm in-place không shadow, báo cáo giữa chừng sẽ hiển thị số nửa cũ nửa mới suốt cả cuối tuần.
Ghi nhớ
- Backfill = chạy lại pipeline cho quá khứ khi đổi logic, sửa bug, hoặc thêm cột/nguồn.
- Điều kiện tiên quyết: partition theo thời gian + task idempotent — mỗi partition là đơn vị chạy lại độc lập.
- Overwrite theo partition (dynamic overwrite hoặc DELETE+INSERT trong transaction), tuyệt đối không append mù kẻo nhân đôi.
- Giới hạn song song/concurrency và tách tài nguyên để backfill không cướp job production đang chạy.
- Đổi logic sâu, nhạy cảm → shadow table + blue-green swap: QC trước, hoán đổi nguyên tử, giữ đường rollback.
- Version hoá logic (cột
logic_version, Git commit) để biết partition nào đã tính lại và có thể dừng/tiếp tục. - Trong Airflow: phân biệt
catchup(bù lịch tự động) vớidags backfill(bù có chủ đích); dùngmax_active_runs/pool để kiểm soát. - Giữ lớp bronze bất biến là điều kiện để reprocess vô hạn; đừng để mất nguồn thô.
Nguồn tham khảo
- Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly): chương ingestion/transformation, backfill & reprocessing.
- Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly): idempotency, batch reprocessing, derived data.
- Apache Airflow — Docs, "Backfill and Catchup" / DAG runs & data intervals: https://airflow.apache.org/docs/apache-airflow/stable/
- Delta Lake — Docs,
INSERT OVERWRITE/MERGE/ dynamic partition overwrite: https://docs.delta.io/latest/ - Apache Iceberg — Docs, partitioning, schema evolution &
REPLACE/OVERWRITE: https://iceberg.apache.org/docs/latest/ - Apache Spark — SQL & DataFrame,
spark.sql.sources.partitionOverwriteMode: https://spark.apache.org/docs/latest/ - The Data Warehouse Toolkit — Kimball & Ross: xử lý dữ liệu lịch sử, bi-temporal/effective dating.
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ẻ!