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.
- 1Nhận batch_dateScheduler truyền ngày logic vào, mỗi lần chạy chỉ xử lý một partition
- 2Chọn dữ liệuLọc theo ngày, có thể lùi watermark để vớt dòng đến trễ
- 3Ghi idempotentMERGE theo key hoặc ghi đè partition trong transaction, dùng batch_date thay NOW()
- 4Chạy lại để kiểm traChạy cùng một ngày hai lần, đếm bản trùng và so sánh kết quả
- 5BackfillLặp qua từng ngày lịch sử, gọi lại đúng job chạy hằng đêm
Khi bước ghi đã idempotent, chạy lại hay backfill đều cho cùng một kết quả.
Đồ hoạ: FDE Times
Tóm tắt nhanh
- Watermark chỉ giúp chọn dữ liệu, còn chống trùng phải nằm ở bước ghi.
- MERGE theo business key, hoặc DELETE rồi INSERT cả partition trong một transaction, giúp chạy lại không đổi kết quả.
- Khi từng partition idempotent, backfill chỉ là chạy lại từng ngày theo thứ tự.
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”.
-- 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.
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.
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.
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:
-- 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.
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ể.
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.