# Thực hành: pipeline nạp tăng dần chạy lại bao nhiêu lần cũng không trùng dữ liệu

> Khách sẽ đưa bạn ba năm dữ liệu cũ và một job chạy hằng đêm, và cả hai chỉ an toàn khi lần chạy thứ hai cho đúng kết quả như lần đầu.

Bản gốc: https://fdetimes.net/vi/bach-khoa/thuc-hanh-pipeline-idempotent-incremental-backfill/

Hãy hình dung tuần đầu ở một khách hàng bán lẻ. Job nạp đơn hàng đêm qua chết giữa chừng, người vận hành bấm chạy lại, và sáng nay dashboard báo doanh thu gấp đôi. Chẳng ai viết sai câu SQL nào cả. Lỗi nằm ở chỗ job này chưa bao giờ được thiết kế để chạy hai lần.

Nếu làm FDE, bạn nên tính trước rằng mình sẽ gặp tình huống này.

Rồi sẽ đến lúc khách hỏi: "Nạp giúp dữ liệu từ đầu năm ngoái được không?" Nếu pipeline chạy lại mà không sinh bản trùng, bạn có thể nhận lời mà không phải thức trắng đêm.

Bài này dựng một pipeline đơn hàng theo ngày, có ba tính chất: nạp tăng dần, chạy lại an toàn, và backfill bằng chính đoạn code chạy hằng đêm. Các đoạn SQL dưới đây đã được giản lược, còn cú pháp cụ thể sẽ khác nhau tùy kho dữ liệu của khách.

## Bạn cần gì trước khi bắt đầu?

IBM mô tả một data pipeline gồm ba chặng: nạp dữ liệu thô từ nhiều nguồn, biến đổi, rồi lưu vào kho. Pipeline ở đây chạy theo kiểu batch, tức là nạp từng lô dữ liệu theo những khoảng thời gian định sẵn. Việc kích hoạt từng lô do job scheduler đảm nhận, tức phần mềm tự chạy các job nền theo lịch mà không cần ai ngồi canh.

Bạn cần một database hỗ trợ transaction và MERGE, cùng hai bảng: `staging_orders` chứa dữ liệu thô vừa kéo về, `orders` là bảng đích. Mỗi đơn hàng có `order_id` làm business key và `order_date` làm ngày sự kiện. Thêm một script Python nhỏ để gọi job theo từng ngày là đủ.

Trước khi viết dòng nào, hãy nắm định nghĩa này: một thao tác là idempotent nếu áp dụng nhiều lần mà kết quả không đổi. Nhờ vậy nó có thể retry mà không gây tác dụng phụ hay làm hỏng dữ liệu.

Có một điểm hay bị bỏ qua: idempotence là thuộc tính của từng thao tác, không tự động đúng cho cả hệ thống. Vì thế bước nào trong pipeline cũng phải được thiết kế riêng cho nó.

## Bước 1: lấy partition theo ngày làm đơn vị công việc

Mỗi lần chạy chỉ xử lý đúng một ngày logic, gọi là `batch_date`. Scheduler truyền ngày này vào làm tham số, job không tự đi tìm "hôm nay là ngày mấy".

```sql
-- Mọi câu lệnh trong job đều lọc theo tham số này
SELECT order_id, customer_id, amount, updated_at, order_date
FROM staging_orders
WHERE order_date = :batch_date;
```

**Kiểm tra:** chạy câu trên với `2026-03-01` hai lần, số dòng trả về phải giống nhau. Nếu khác, nguồn dữ liệu đang thay đổi trong lúc bạn đọc, và bạn cần biết điều đó trước khi đi tiếp.

## Bước 2: dùng MERGE để chạy lại không ra kết quả mới

INSERT thường là nguyên nhân của dashboard doanh thu gấp đôi. Giả sử ngày 2026-03-01 có 1.000 đơn. Chạy INSERT hai lần thì bảng đích có 2.000 dòng. MERGE (upsert) thì khác: key mới được chèn vào, key đã có được cập nhật.

```sql
MERGE INTO orders AS t
USING (
SELECT order_id, customer_id, amount, updated_at, order_date
FROM staging_orders
WHERE order_date = :batch_date
) AS s
ON t.order_id = s.order_id
WHEN MATCHED THEN UPDATE SET
customer_id = s.customer_id,
amount      = s.amount,
updated_at  = s.updated_at
WHEN NOT MATCHED THEN INSERT
(order_id, customer_id, amount, updated_at, order_date)
VALUES (s.order_id, s.customer_id, s.amount, s.updated_at, s.order_date);
```

Khi chạy lại cùng batch, mỗi dòng chỉ được cập nhật thành đúng giá trị nó đang giữ, nên lần chạy thứ hai không thay đổi gì. Trong ví dụ trên, bảng đích vẫn giữ 1.000 dòng sau bao nhiêu lần chạy đi nữa.

**Kiểm tra:** chạy MERGE hai lần rồi tìm bản trùng. Kết quả phải rỗng.

```sql
SELECT order_id, COUNT(*)
FROM orders
GROUP BY order_id
HAVING COUNT(*) > 1;
```

## Bước 3: khi không có key, ghi đè cả partition

Có những bảng không có business key ổn định, ví dụ bảng tổng hợp doanh thu theo ngày. Cách làm lúc này là coi cả partition là đơn vị công việc và thay thế nó toàn bộ. Chạy lại một ngày thì ngày đó được thay mới chứ không bị nhân đôi.

```sql
BEGIN;
DELETE FROM daily_revenue WHERE order_date = :batch_date;
INSERT INTO daily_revenue (order_date, total_amount, order_count)
SELECT order_date, SUM(amount), COUNT(*)
FROM orders
WHERE order_date = :batch_date
GROUP BY order_date;
COMMIT;
```

`BEGIN` và `COMMIT` là phần quan trọng nhất của đoạn này. DELETE rồi INSERT chỉ an toàn khi nằm trong một transaction. Nếu job sập giữa hai câu lệnh, partition sẽ bị bỏ trống, và sáng hôm sau khách thấy doanh thu ngày đó bằng 0.

**Kiểm tra:** trong môi trường thử, cho job dừng ngay sau DELETE (ví dụ chèn một lỗi cố ý), rồi xác nhận dữ liệu cũ của ngày đó vẫn còn nguyên.

## Bước 4: bỏ NOW() ra khỏi câu ghi

Rất nhiều job ghi `loaded_at = NOW()` vào từng dòng. Pipeline làm vậy sẽ cho ra dữ liệu khác nhau ở mỗi lần chạy lại, nghĩa là không còn idempotent nữa. Hãy dùng thời gian của sự kiện hoặc ngày logic của batch. Đây vẫn là bước ghi của bước 3, chỉ đổi câu INSERT bên trong khối `BEGIN ... COMMIT`:

```sql
-- Trước: ... SELECT order_date, SUM(amount), COUNT(*), NOW()
-- Sau: ghi ngày logic mà scheduler truyền vào
INSERT INTO daily_revenue (order_date, total_amount, order_count, batch_date)
SELECT order_date, SUM(amount), COUNT(*), :batch_date
FROM orders
WHERE order_date = :batch_date
GROUP BY order_date;
```

**Kiểm tra:** chạy một ngày hai lần, xuất kết quả ra file rồi so sánh hai file. Nếu khác nhau dù chỉ một cột, vẫn còn một nguồn không tất định nào đó trong job.

## Vì sao watermark không cứu được bạn?

Cách nạp tăng dần quen thuộc là lưu một watermark, chẳng hạn `updated_at` lớn nhất đã nạp, rồi lần sau chỉ lấy những dòng mới hơn. Một bài phân tích của Algoscale nhận xét logic watermark trông đơn giản nhưng mong manh đến bất ngờ.

Nó dễ vỡ ở timestamp nằm đúng ranh giới, ở những dòng đến trễ, ở chênh lệch múi giờ, và đã có một trường hợp trên Fabric sinh ra dữ liệu trùng.

Vì thế, đừng bắt watermark gánh phần đúng đắn của dữ liệu. Khi bước ghi đã idempotent, mỗi business key chỉ có một dòng, thì một dòng có bị chọn hai lần cũng không tạo bản trùng. Theo Algoscale, lúc đó việc chọn trùng không còn làm sai dữ liệu nữa, cùng lắm chỉ tốn thêm chút công đọc.

**Điểm mấu chốt:** Đừng cố chọn dữ liệu cho hoàn hảo; hãy làm cho việc ghi lặp lại trở nên vô hại.

Trên thực tế, khi bước ghi đã là MERGE, bạn có thể chủ động lùi watermark lại một khoảng để vớt các dòng đến trễ. Cái giá là đọc thêm một ít dữ liệu, đổi lại không phải lo mất dòng.

## Bước 5: backfill chỉ là một vòng lặp

Đến đây việc backfill dữ liệu lịch sử của khách trở nên đơn giản: chạy lại từng partition. Data Vidhya viết rằng các đội có pipeline idempotent làm backfill vào một buổi chiều thứ Ba bình thường. Đoạn script dưới đây đã được giản lược, `run_partition` là hàm gọi các bước 1 đến 4 với một ngày cụ thể.

```python
from datetime import date, timedelta

def run_partition(batch_date):
# Gọi MERGE vào orders, rồi ghi đè daily_revenue cho batch_date
...

d = date(2025, 1, 1)
end = date(2025, 12, 31)
while d <= end:
run_partition(d)
d += timedelta(days=1)
```

**Kiểm tra:** chạy vòng lặp hai lần cho một khoảng ngắn, chẳng hạn ba ngày. Số dòng và tổng tiền phải giống hệt nhau giữa hai lần. Nếu vòng lặp chết ở ngày thứ hai trăm, bạn chỉ cần chạy lại từ ngày đó mà không phải dọn dẹp gì.

## Những chỗ hay vấp

Chỗ vấp phổ biến nhất là chọn sai business key. Nếu hệ thống nguồn của khách tái sử dụng `order_id` giữa các chi nhánh, MERGE sẽ ghi đè nhầm đơn hàng, và key đúng phải là cặp chi nhánh với mã đơn. Hãy hỏi khách điều này ngay từ buổi customer discovery, đừng tự đoán.

Cũng hay gặp cảnh DELETE và INSERT nằm ở hai job riêng, hoặc chạy ở chế độ autocommit, tức là mất luôn tấm lưới transaction của bước 3. Và nhiều đội chỉ thử đường đi thuận lợi. Một pipeline chưa từng bị chạy hai lần cho cùng một ngày thì xem như chưa được kiểm thử.

## Khi bạn đã ở chỗ khách hàng

Việc đầu tiên ở một khách mới không phải là viết pipeline mới mà là hỏi: "Nếu chạy lại job đêm qua thì chuyện gì xảy ra?" Câu trả lời cho bạn biết hệ thống của họ có chịu được retry, có chịu được backfill hay không.

Trước khi nhận lời nạp ba năm dữ liệu cũ, hãy kiểm tra từng bảng đích đang ghi bằng INSERT, MERGE hay ghi đè partition.

Nếu bạn đang ứng tuyển vị trí FDE, hãy để ý những JD nhắc tới "data ingestion", "backfill" hay "pipeline reliability". Trong CV, câu "xây pipeline ETL" gần như không nói lên điều gì. Thay vào đó, hãy viết bạn đã chuyển job từ INSERT sang MERGE theo business key, và nhờ vậy backfill được một năm dữ liệu mà không có bản trùng.

Pipeline nào rồi cũng có lúc hỏng. Điều đáng để bạn khoe là pipeline của mình hỏng xong chỉ cần chạy lại, không ai phải ngồi dọn dữ liệu.

**Thử ngay tuần này:**

- Chọn một job nạp dữ liệu bạn đang có, chạy nó hai lần liên tiếp cho cùng một ngày rồi chạy câu GROUP BY ... HAVING COUNT(*) > 1 để đếm bản trùng
- Tìm mọi chỗ dùng NOW() hoặc thời điểm hiện tại trong job và thay bằng tham số batch_date
- Viết script backfill ba ngày bằng vòng lặp gọi run_partition, chạy hai lần và so sánh số dòng

## Nguồn

- [Why idempotence was important to DevOps (DEV Community)](https://dev.to/startpher/why-idempotence-was-important-to-devops-2jn3)

- [What is a Data Pipeline? (IBM)](https://www.ibm.com/topics/data-pipeline)

- [Job scheduler (Wikipedia)](https://en.wikipedia.org/wiki/Job_scheduler)

- [Watermark Bugs in Fabric Incremental Loads (Algoscale Blog)](https://algoscale.com/go/blog/fabric-watermark-incremental-load-duplicates/)

- [Idempotent Loads · Data Warehousing (Data Vidhya)](https://datavidhya.com/learn/data-warehouse/loading-the-warehouse/idempotent-loads/)
