Pipeline có hai thứ cần kiểm tra: code biến đổi có đúng logic không và dữ liệu thật đi qua có đúng kỳ vọng không. Unit test chạy trên dữ liệu giả nhỏ trước khi merge; data test chạy trên dữ liệu thật mỗi lần pipeline chạy. Hai loại bổ sung cho nhau, không thay nhau được.
1. Unit test cho logic biến đổi. Tách phần biến đổi thành hàm nhận DataFrame, trả DataFrame, rồi test bằng pytest với vài dòng tự dựng, gồm cả ca biên (NULL, trùng, giá trị âm).
def test_dedup_keeps_latest(spark):
df = spark.createDataFrame(
[("o1", "2026-09-01 10:00", 100), ("o1", "2026-09-01 11:00", 120)],
["order_id", "updated_at", "amount"],
)
out = dedup_latest(df, key="order_id", order_col="updated_at")
assert [r["amount"] for r in out.collect()] == [120]spark là fixture tạo một SparkSession local dùng chung cho cả lượt chạy test. dbt từ v1.8 cũng có unit test: khai dữ liệu đầu vào ở given và kết quả mong đợi ở expect trong YAML.
2. Data test trên dữ liệu thật, chạy sau mỗi bước: khoá không trùng, không NULL, giá trị thuộc tập cho phép, khoá ngoại tồn tại. dbt có sẵn unique, not_null, accepted_values, relationships; test tự viết là một câu SQL trả về các dòng vi phạm, có dòng nào là fail.
3. Test end-to-end: chạy cả DAG trên staging với một lát dữ liệu thật (một ngày), so số dòng và tổng tiền với nguồn, so output với bản production hiện tại trước khi thay.
4. Contract/schema test ở ranh giới với hệ nguồn: nguồn đổi tên hay xoá cột thì pipeline dừng sớm với lỗi rõ ràng thay vì ghi ra NULL.
Lưu ý: test chỉ bắt được những lỗi bạn đã nghĩ tới. Cần thêm monitoring trên production (số dòng mỗi ngày, độ tươi, tỉ lệ NULL so với trung bình 7 ngày) để phát hiện những lỗi chưa có test.