Batch 5 — Backfill & Tái xử lý lịch sử

22 thg 7, 2026 2 lượt xem
#data-engineering
#airflow
#partitioning
#batch
#backfill
#reprocessing

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ốngVí 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ồnBả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 — IdempotencyBatch 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ô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_transaction là 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ôngpartitionOverwriteMode=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ặc pipeline_run_id, processed_at và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):

  1. 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).
  2. 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ó.
  3. 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).
  4. Khi đạt, hoán đổi con trỏ: đổi tên bảng, hoặc trỏ view/alias fact_transaction sang green. Thao tác này gần như tức thời và nguyên tử.
  5. 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 channel cho 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:

  1. Xây bảng shadow loan_classification__v2 với logic mới; giữ bảng blue đang phục vụ báo cáo.
  2. Backfill song song qua Airflow: airflow dags backfill loan_classification_daily --start-date 2026-01-01 --end-date 2026-03-31, đặt max_active_runs=5 để không đụng job EOD; mỗi run overwrite đúng partition ngày (idempotent), ghi kèm logic_version='v2'.
  3. 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.
  4. Hoán đổi view loan_classification sang 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ới dags backfill (bù có chủ đích); dùng max_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.

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