Batch 3 — Xử lý gia tăng (Incremental) & CDC
Mô hình tinh thần: đừng chép lại thứ đã chép
Ở bài tổng quan series chúng ta đã thống nhất rằng batch là "gom một khối dữ liệu rồi xử lý một lần". Câu hỏi tiếp theo, và cũng là câu hỏi định hình toàn bộ chi phí vận hành, là: khối đó lớn cỡ nào?
Hình dung bảng giao dịch của một ngân hàng: mỗi ngày thêm vài chục triệu dòng, tích lũy nhiều năm thành hàng tỷ dòng. Job ETL đêm nào cũng đọc lại toàn bộ bảng đó, biến đổi, rồi ghi đè bảng đích — đó là full load (nạp toàn bộ). Nó đơn giản đến mức hấp dẫn: không cần nhớ lần chạy trước làm tới đâu, kết quả luôn phản ánh trạng thái nguồn hiện tại, chạy lại bao nhiêu lần cũng ra y hệt. Nhưng cái giá là bạn quét lại tỷ dòng lịch sử không hề đổi mỗi đêm chỉ để lấy vài chục triệu dòng mới. Trên warehouse tính tiền theo lượng dữ liệu quét hoặc thời gian compute, đó là lãng phí biến thẳng thành hóa đơn — và tệ hơn, thời gian chạy dài dần đến lúc không kịp cửa sổ đêm (EOD).
Incremental load (nạp gia tăng) đảo ngược mặc định: dữ liệu cũ đã xử lý rồi thì đừng đụng vào nữa. Mỗi lần chạy chỉ lấy phần mới hoặc vừa thay đổi kể từ lần trước. "Xây thêm tầng" thay vì "xây lại cả tòa nhà mỗi sáng". Cái giá phải trả là độ phức tạp: bạn phải tự trả lời hai câu hỏi mà full load lo hộ — "dòng nào là mới?" và "khi một dòng cũ thay đổi/bị xóa thì xử lý sao?". Trả sai là trùng hoặc thiếu dữ liệu, những lỗi âm thầm khó phát hiện.
| Tiêu chí | Full load | Incremental load |
|---|---|---|
| Lượng dữ liệu đọc mỗi lần | Toàn bộ nguồn | Chỉ phần mới/đổi |
| Thời gian & chi phí | Tăng theo kích thước tích lũy | Tăng theo lượng thay đổi/ngày |
| Trạng thái cần nhớ | Không | Vị trí lần chạy trước (watermark) |
| Xử lý delete | Tự nhiên (ghi đè) | Phải xử lý riêng |
| Rủi ro trùng/thiếu | Thấp | Cao nếu logic sai |
| Khi nào chọn | Bảng nhỏ, dim thay đổi toàn phần | Bảng sự kiện lớn, append nhiều |
Quy tắc thực dụng: bắt đầu bằng full load, chỉ chuyển sang incremental khi full load thật sự đắt hoặc không kịp cửa sổ. Đừng tối ưu sớm một bảng vài triệu dòng.
High-water-mark: kỹ thuật incremental nền tảng
Cách phổ biến nhất để trả lời "dòng nào là mới" là high-water-mark (mốc nước cao — HWM). Ý tưởng: chọn một cột chỉ tăng (monotonic) trên bảng nguồn — thường là updated_at (timestamp) hoặc một id/sequence tự tăng — rồi nhớ giá trị lớn nhất đã xử lý ở lần chạy trước. Lần sau chỉ lấy các dòng có giá trị lớn hơn mốc đó.
-- Lần chạy: đọc mốc đã lưu, lấy dữ liệu mới hơn mốc
SELECT *
FROM core.transactions
WHERE updated_at > :last_high_water_mark -- mốc lần trước, ví dụ '2026-07-20 23:00:00'
ORDER BY updated_at;
-- Sau khi ghi thành công, cập nhật mốc = max(updated_at) của batch vừa lấy
Điểm mấu chốt là thứ tự thao tác: chỉ cập nhật mốc sau khi đã ghi thành công vào đích. Nếu job chết giữa chừng mà mốc chưa dời, lần sau chạy lại từ mốc cũ — có thể lấy trùng vài dòng, nhưng nếu bước ghi là idempotent (xem bài batch-04) thì trùng cũng vô hại. Đây là lý do incremental và idempotency luôn đi cặp.
Những cái bẫy của HWM
Cột id tự tăng đơn giản nhưng chỉ bắt được insert, không bắt được update (id cũ không đổi khi dòng bị sửa). Muốn bắt cả update phải dùng updated_at — và điều này buộc nguồn phải kỷ luật cập nhật cột đó ở mọi thao tác ghi, kể cả trigger. Nếu ứng dụng quên cập nhật updated_at, dòng đó "tàng hình" với pipeline.
Cái bẫy tinh vi hơn là thời điểm commit vs thời điểm timestamp. Nhiều DB gán updated_at = giờ bắt đầu transaction chứ không phải giờ commit. Một transaction mở lúc 22:59, commit lúc 23:05, mang timestamp 22:59. Nếu job chạy lúc 23:00 lấy > 22:59 rồi dời mốc lên 23:00, dòng commit-muộn kia có timestamp 22:59 < 23:00 sẽ không bao giờ được lấy ở lần sau. Đây chính là bài toán dẫn tới khái niệm watermark dưới đây.
Watermark: mốc có vùng đệm an toàn
Watermark trong ngữ cảnh batch là một high-water-mark được cố ý lùi lại một khoảng an toàn để dung nạp độ trễ và sai lệch đồng hồ. Thay vì "lấy tất cả > mốc rồi dời mốc lên max", ta chừa một vùng chồng lấn (overlap window):
-- Lấy dữ liệu trong cửa sổ có chồng lấn lùi 2 giờ,
-- chấp nhận đọc lại một ít để KHÔNG bỏ sót dòng commit muộn
SELECT *
FROM core.transactions
WHERE updated_at > :last_hwm - INTERVAL '2' HOUR; -- lùi mốc 2h làm đệm
-- Vì đọc lại chồng lấn, bước ghi BẮT BUỘC phải upsert (idempotent),
-- nếu không sẽ nhân bản dòng trong vùng overlap.
Đánh đổi ở đây rất rõ: vùng đệm càng rộng thì càng an toàn khỏi mất dữ liệu nhưng càng đọc lại nhiều (tốn compute) và bắt buộc phải upsert. Đây là biểu hiện batch của cùng ý niệm "watermark" trong streaming (ranh giới "đã thấy đủ dữ liệu tới thời điểm T"), chỉ khác là ở batch ta hiện thực nó bằng một khoảng lùi tĩnh thay vì ước lượng động.
Lưu ý thuật ngữ: "high-water-mark" và "watermark" thường được dùng lẫn lộn. Trong bài này, HWM = mốc lớn nhất đã xử lý; watermark = HWM đã lùi một khoảng đệm. Điều quan trọng không phải là tên gọi mà là có hay không vùng đệm cho dữ liệu đến trễ.
Change Data Capture kiểu batch
HWM giải quyết insert và update nếu nguồn có cột thời gian tin cậy. Nhưng nó không bắt được delete (dòng biến mất, chẳng để lại dấu vết updated_at), và nó phụ thuộc vào kỷ luật của ứng dụng nguồn. Change Data Capture (CDC) là họ kỹ thuật bắt mọi thay đổi (insert/update/delete) một cách đầy đủ hơn. Trong ngữ cảnh batch có hai trường phái chính.
Cách 1 — So sánh snapshot (snapshot-diff / query-based CDC)
Chụp toàn bộ trạng thái bảng hôm nay, so với snapshot đã lưu hôm qua, rồi suy ra thay đổi bằng một phép join theo khóa:
-- So sánh 2 snapshot để suy ra Insert / Update / Delete
SELECT
COALESCE(t.account_id, y.account_id) AS account_id,
CASE
WHEN y.account_id IS NULL THEN 'INSERT' -- có hôm nay, không có hôm qua
WHEN t.account_id IS NULL THEN 'DELETE' -- có hôm qua, biến mất hôm nay
WHEN t.row_hash <> y.row_hash THEN 'UPDATE'
ELSE 'UNCHANGED'
END AS change_type
FROM snapshot_today t
FULL OUTER JOIN snapshot_yesterday y USING (account_id)
WHERE t.row_hash IS DISTINCT FROM y.row_hash
OR t.account_id IS NULL
OR y.account_id IS NULL;
Mẹo hiệu năng: thay vì so từng cột, tính một row_hash (ví dụ md5 của các cột nghiệp vụ) rồi chỉ so hash — nhanh và gọn. Ưu điểm của snapshot-diff: không cần quyền đặc biệt trên DB nguồn, bắt được cả delete, dễ hiểu. Nhược điểm: phải chụp toàn bảng mỗi kỳ (đắt như full load ở phía đọc), và chỉ thấy trạng thái đầu–cuối, mất các thay đổi trung gian trong ngày (một dòng đổi 3 lần thì chỉ thấy 1).
Cách 2 — Đọc log / CDC feed (log-based CDC)
Mọi DB giao dịch đều ghi một nhật ký thay đổi (redo log của Oracle, WAL của PostgreSQL, binlog của MySQL) để phục hồi. Log-based CDC đọc chính nhật ký đó để lấy dòng sự kiện thay đổi — insert/update/delete đều tường minh, kèm ảnh before/after, mà không chạy query nào lên bảng nguồn (gần như không thêm tải). Đây là cách chính xác và ít xâm lấn nhất, nền tảng của các công cụ như Debezium hay Oracle GoldenGate.
Trong series này ta bàn CDC ở góc độ nguyên lý batch; phần hiện thực cụ thể trên Oracle (GoldenGate, Debezium, redo/archive log, cấu hình supplemental logging) được đào sâu ở series Oracle CDC. Khi CDC feed chảy liên tục qua Kafka thì nó đã bước sang địa hạt streaming — nhưng rất thường ta tiêu thụ feed đó theo lô (gom sự kiện của một khoảng thời gian rồi merge một lần), nên ranh giới batch/streaming ở đây khá mờ.
| Snapshot-diff (query) | Log-based | |
|---|---|---|
| Bắt delete | Có | Có (tường minh) |
| Bắt thay đổi trung gian | Không (chỉ đầu–cuối) | Có (từng sự kiện) |
| Tải lên DB nguồn | Cao (quét toàn bảng) | Rất thấp (đọc log) |
| Quyền/cấu hình nguồn | Chỉ cần SELECT | Cần bật logging, quyền đặc biệt |
| Độ trễ | Theo chu kỳ snapshot | Gần như tức thời |
| Độ phức tạp vận hành | Thấp | Cao hơn |
Xử lý UPDATE và DELETE ở đích
Bắt được thay đổi mới xong một nửa; nửa còn lại là áp chúng vào bảng đích cho đúng.
Update ở đích không thể là "insert thêm" — sẽ có hai phiên bản của cùng một khóa. Phải ghi đè dòng cũ (upsert) hoặc, nếu cần giữ lịch sử, mở một bản ghi mới và đóng bản cũ (SCD Type 2 — chủ đề của series data modeling). Ở đây ta tập trung kiểu ghi đè trạng thái hiện tại.
Delete khó hơn vì có hai lựa chọn triết lý:
- Hard delete: xóa hẳn dòng ở đích. Trung thực với nguồn nhưng mất dấu vết — báo cáo lịch sử, kiểm toán, phân tích "tài khoản đã đóng" không còn dữ liệu. Với ngân hàng, hiếm khi chấp nhận được.
- Soft delete: không xóa, chỉ đánh dấu
is_deleted = true(vàdeleted_at). Bảng đích giữ nguyên dòng nhưng lọc ra khỏi khung nhìn nghiệp vụ. Đây là mặc định an toàn cho dữ liệu tài chính: giữ được lịch sử và khả năng kiểm toán. Query nghiệp vụ luôn thêmWHERE NOT is_deleted.
Lưu ý: HWM query-based không thấy delete, nên nếu dùng HWM mà nguồn có xóa vật lý, bạn buộc phải bổ sung snapshot-diff định kỳ hoặc chuyển sang log-based để bắt delete.
MERGE / UPSERT trong warehouse
Công cụ để áp một lô thay đổi (có cả I/U/D) vào bảng đích trong một câu lệnh nguyên tử là MERGE (chuẩn SQL, còn gọi upsert = update + insert). Hầu hết warehouse hiện đại (Snowflake, BigQuery, Redshift, Databricks/Delta, Oracle) đều hỗ trợ. Ý tưởng: join nguồn thay đổi với đích theo khóa, rồi matched thì update, not-matched thì insert, và xử lý delete bằng cờ soft-delete.
MERGE INTO dwh.dim_account AS tgt
USING (
-- Nguồn: batch thay đổi đã khử trùng, mỗi khóa chỉ 1 dòng mới nhất
SELECT account_id, customer_name, status, change_type, updated_at
FROM staging.account_changes
QUALIFY ROW_NUMBER() OVER (
PARTITION BY account_id ORDER BY updated_at DESC) = 1
) AS src
ON tgt.account_id = src.account_id
-- 1) Dòng bị xóa ở nguồn -> soft delete ở đích (giữ lịch sử)
WHEN MATCHED AND src.change_type = 'DELETE' THEN UPDATE SET
tgt.is_deleted = TRUE,
tgt.deleted_at = src.updated_at,
tgt.updated_at = src.updated_at
-- 2) Dòng đổi và MỚI HƠN bản đang có -> update (guard chống ghi đè ngược)
WHEN MATCHED AND src.change_type = 'UPDATE'
AND src.updated_at > tgt.updated_at THEN UPDATE SET
tgt.customer_name = src.customer_name,
tgt.status = src.status,
tgt.is_deleted = FALSE,
tgt.updated_at = src.updated_at
-- 3) Khóa chưa có -> insert
WHEN NOT MATCHED AND src.change_type <> 'DELETE' THEN INSERT
(account_id, customer_name, status, is_deleted, updated_at)
VALUES
(src.account_id, src.customer_name, src.status, FALSE, src.updated_at);
Ba chi tiết khiến câu MERGE này an toàn khi chạy lại (idempotent):
- Khử trùng nguồn trước khi merge (
ROW_NUMBER ... = 1). MERGE hầu hết phương ngữ báo lỗi hoặc cho kết quả bất định nếu một khóa đích khớp nhiều dòng nguồn. Luôn đảm bảo nguồn có một dòng mới nhất trên mỗi khóa. - Guard
src.updated_at > tgt.updated_at: nếu vì đọc chồng lấn (watermark) hay chạy lại mà một thay đổi cũ lọt vào, điều kiện này chặn nó ghi đè bản mới hơn. Đây là chốt chặn chống dữ liệu đến trễ ghi lộn. - Soft delete bằng UPDATE cờ, không phải
DELETE, để giữ dòng phục vụ kiểm toán.
Nếu warehouse của bạn dùng dbt, chính cơ chế này được gói lại thành materialization incremental với unique_key và các chiến lược merge/delete+insert/insert_overwrite — xem dbt: Incremental models & hiệu năng.
Late-arriving data (dữ liệu đến trễ)
Dữ liệu đến trễ là dòng có thời điểm nghiệp vụ thuộc quá khứ nhưng đến pipeline muộn: giao dịch phát sinh 23:58 nhưng chi nhánh đồng bộ lên trung tâm lúc 00:30 hôm sau; POS ngoại tuyến gom cả ngày, tối mới đẩy về. Nếu đã "chốt sổ" ngày hôm trước, dòng này rơi vào khoảng trống.
Có mấy chiến lược, chọn theo yêu cầu chính xác:
- Vùng đệm watermark (đã bàn): lùi mốc một khoảng đủ dung nạp độ trễ thường gặp. Rẻ, đơn giản, nhưng chỉ cứu được trễ trong khoảng đệm.
- Merge theo thời điểm nghiệp vụ, không theo thời điểm nạp: partition và upsert dữ liệu vào đúng ngày nghiệp vụ của nó, dù nó đến hôm nay. Dòng trễ của ngày 20/07 vẫn được gộp vào partition
event_date = 2026-07-20. - Backfill / reprocessing có chủ đích: định kỳ chạy lại N ngày gần nhất để cuốn hết dòng trễ. Đây là lý do incremental thường đi kèm một job backfill — chủ đề trọn vẹn ở bài batch-05.
- Restatement: chấp nhận rằng con số của một ngày có thể được điều chỉnh lại trong vài ngày sau (T+n), và đánh dấu rõ trạng thái "sơ bộ" vs "chốt". Ngành ngân hàng quen mô hình này (số liệu EOD sơ bộ rồi chốt sau đối soát).
Nguyên tắc bao trùm: thời điểm dữ liệu đến (arrival time) và thời điểm nghiệp vụ (event time) là hai trục khác nhau. Incremental theo arrival time để hiệu quả, nhưng phải merge và phân vùng theo event time để đúng.
Use case thực tế
(Số liệu minh họa.) NCB có bảng transactions ~1,8 tỷ dòng, mỗi ngày thêm ~40 triệu dòng. Job đối chiếu EOD ban đầu full load mất ~2 giờ 40 phút và ngày càng sát trần cửa sổ đêm.
Chuyển sang incremental theo HWM trên updated_at, đọc với watermark lùi 3 giờ để dung nạp dữ liệu chi nhánh đồng bộ trễ, rồi MERGE vào bảng đích theo transaction_id:
- Lượng dữ liệu đọc mỗi đêm giảm từ 1,8 tỷ xuống ~45 triệu dòng (40 triệu mới + ~5 triệu trong vùng chồng lấn) → thời gian chạy còn ~9 phút.
- Vì đọc chồng lấn, bước ghi bắt buộc là
MERGEvới guardupdated_at; ~5 triệu dòng đọc lại được upsert vô hại thay vì nhân bản. - Delete (giao dịch bị hủy/điều chỉnh) không bắt được qua HWM, nên bổ sung một job snapshot-diff hằng tuần trên tập tài khoản active để soft-delete các bản ghi biến mất, giữ nguyên lịch sử cho kiểm toán.
- Dữ liệu POS ngoại tuyến đến trễ được merge theo
event_datenên báo cáo ngày cũ được cập nhật đúng thay vì rơi mất; con số EOD được đánh dấu "sơ bộ" đến khi backfill T+2 chốt lại.
Chi tiết vòng đời một job EOD ngân hàng đầu–cuối được kể ở bài batch-09.
Ghi nhớ
- Full load đơn giản và an toàn nhưng chi phí tăng theo kích thước tích lũy; incremental đổi độ phức tạp lấy hiệu quả — chỉ chuyển khi full load thật sự đắt/không kịp.
- High-water-mark: nhớ giá trị lớn nhất đã xử lý của một cột chỉ tăng (
updated_at/id), lần sau lấy phần lớn hơn. Chỉ dời mốc sau khi ghi thành công. - Cột
idchỉ bắt insert; muốn bắt update phải cóupdated_atđược nguồn cập nhật kỷ luật. HWM không bắt được delete. - Watermark = HWM + vùng đệm: lùi mốc một khoảng để không sót dòng commit/đến trễ, đổi lại phải upsert vì đọc chồng lấn.
- CDC batch có hai trường phái: snapshot-diff (không cần quyền đặc biệt, bắt delete, nhưng mất thay đổi trung gian và quét toàn bảng) và log-based (chính xác, ít xâm lấn, thấy từng sự kiện, nhưng cần bật logging).
- Xử lý delete bằng soft delete (
is_deleted) cho dữ liệu tài chính để giữ lịch sử/kiểm toán; áp thay đổi bằng MERGE/upsert với nguồn đã khử trùng và guard chống ghi đè ngược. - Dữ liệu đến trễ: phân biệt arrival time vs event time — incremental theo arrival để nhanh, merge/partition theo event time để đúng; kết hợp backfill và restatement.
- Incremental và idempotency luôn đi cặp; incremental và backfill cũng vậy.
Nguồn tham khảo
- Fundamentals of Data Engineering — Joe Reis & Matt Housley (O'Reilly), chương ingestion & CDC.
- Designing Data-Intensive Applications — Martin Kleppmann (O'Reilly), chương về logs & change capture.
- The Data Warehouse Toolkit — Kimball & Ross, phần change tracking / slowly changing dimensions.
- Debezium Documentation — CDC log-based (https://debezium.io/documentation/).
- dbt Documentation — Incremental models & materializations (https://docs.getdbt.com/docs/build/incremental-models).
- Delta Lake Documentation —
MERGE INTO/ upsert (https://docs.delta.io/latest/delta-update.html). - "The Log: What every software engineer should know about real-time data's unifying abstraction" — Jay Kreps.
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ẻ!