Batch 8 — Data Quality & Kiểm thử pipeline
Vì sao data quality là một phần của pipeline, không phải việc dọn sau
Ở Batch 2 — ETL vs ELT ta biết dữ liệu chảy từ nguồn qua transform tới đích; ở Batch 7 — Orchestration ta biết DAG điều phối các bước theo thứ tự và điều kiện. Bài này trả lời câu hỏi mà mọi báo cáo cuối cùng đều dựa vào: dữ liệu chảy ra khỏi pipeline có tin được không?
Mô hình tinh thần cần nắm trước: chất lượng dữ liệu không phải một bước "làm sạch" chạy một lần, mà là một tập kiểm soát (control) gắn vào từng chặng của pipeline. Giống như dây chuyền sản xuất có trạm QC ở nhiều công đoạn chứ không chỉ kiểm ở cuối, pipeline batch cần kiểm ở lúc ingest (dữ liệu vào có đúng hình dạng không), sau transform (logic có giữ được ràng buộc không), và trước khi publish (bảng vàng có đạt chuẩn để báo cáo/nộp cơ quan không). Nếu chỉ kiểm ở cuối, lỗi đã lan khắp các bảng phụ thuộc và việc truy nguyên tốn gấp nhiều lần.
Trong ngân hàng, một dữ liệu sai không chỉ khó chịu — nó thành hậu quả pháp lý. Một số dư lệch, một báo cáo phân loại nợ sai nhóm, một bảng sao kê thiếu giao dịch: mỗi lỗi có thể thành finding kiểm toán hoặc khoản phạt. Vì thế bài này coi data quality là cơ chế phòng thủ tự động đặt ngay trong DAG, không phải công việc thủ công cuối tháng.
Sáu chiều chất lượng dữ liệu
Cộng đồng data quality (và các khung như DAMA-DMBOK) mô tả chất lượng theo nhiều chiều (dimension) đo lường được. Sáu chiều dưới đây là tập phổ biến nhất; mỗi chiều trả lời một câu hỏi khác nhau và cần một loại test khác nhau.
| Chiều | Câu hỏi | Ví dụ ngân hàng | Cách đo |
|---|---|---|---|
| Completeness (đầy đủ) | Có thiếu dữ liệu không? | Mọi giao dịch EOD đều có account_id; số bản ghi khớp số nguồn | Tỷ lệ NULL; đếm bản ghi so với kỳ vọng |
| Accuracy (chính xác) | Dữ liệu có đúng thực tế không? | Số dư trong DW khớp số dư core banking | Đối soát với nguồn tin cậy (system of record) |
| Consistency (nhất quán) | Các bản sao/bảng có khớp nhau không? | Tổng dư nợ ở bảng khoản vay = tổng ở sổ cái | So sánh chéo giữa các bảng/hệ thống |
| Timeliness (kịp thời) | Dữ liệu có mới đủ độ không? | Feed EOD về trước 02:00 để kịp phân loại nợ | Độ trễ (freshness) so với SLA |
| Uniqueness (duy nhất) | Có bị trùng lặp không? | Một transaction_id chỉ xuất hiện một lần | Đếm khóa trùng; test unique |
| Validity (hợp lệ) | Dữ liệu có đúng định dạng/miền giá trị không? | currency ∈ {VND, USD, EUR}; ngày ≤ hôm nay | Regex, ràng buộc miền (accepted values), range check |
Hai chiều hay bị nhầm: accuracy là "khớp sự thật bên ngoài" (cần một nguồn chân lý để so), còn consistency là "khớp nội bộ giữa các bản sao" (không cần nguồn ngoài, chỉ cần các bảng khớp nhau). Một cột có thể consistent mà vẫn inaccurate nếu cả hai bảng cùng sai giống nhau.
Sơ đồ dưới đặt sáu chiều vào vị trí kiểm tương ứng trong pipeline batch:
Ý tưởng cốt lõi: đẩy kiểm tra càng sớm (shift-left) thì lỗi bị chặn càng gần nguồn, chi phí sửa càng thấp. Nhưng vài chiều — nhất là accuracy qua đối soát — chỉ đo được ở cuối khi đã có bảng tổng hợp để so với nguồn.
Hai lớp kiểm thử: unit test transform và data test
Người mới hay gộp "test pipeline" thành một thứ. Thực ra có hai lớp khác bản chất:
1. Unit test cho logic transform — kiểm code biến đổi, chạy trên dữ liệu giả cố định, không đụng dữ liệu production. Ví dụ: hàm phân loại nhóm nợ theo số ngày quá hạn (DPD) — cho input DPD=95 phải ra "Nhóm 3". Đây là test tất định, chạy trong CI mỗi lần đổi code, độc lập với dữ liệu thật.
# Unit test một hàm transform thuần (pytest) — dữ liệu giả, tất định
def classify_npl_group(days_past_due: int) -> str:
if days_past_due <= 10: return "Nhóm 1" # đủ tiêu chuẩn
if days_past_due <= 90: return "Nhóm 2" # cần chú ý
if days_past_due <= 180: return "Nhóm 3" # dưới tiêu chuẩn
if days_past_due <= 360: return "Nhóm 4" # nghi ngờ
return "Nhóm 5" # có khả năng mất vốn
def test_npl_boundaries():
assert classify_npl_group(10) == "Nhóm 1"
assert classify_npl_group(11) == "Nhóm 2"
assert classify_npl_group(95) == "Nhóm 3" # ranh giới hay sai
assert classify_npl_group(400) == "Nhóm 5"
2. Data test (kiểm dữ liệu) — kiểm dữ liệu thực sau khi chảy qua pipeline, khẳng định các bất biến (invariant). Bốn nhóm data test kinh điển:
| Loại | Khẳng định | Ví dụ |
|---|---|---|
| not null | Cột bắt buộc không rỗng | account_id, txn_date không NULL |
| unique | Khóa không trùng | transaction_id duy nhất |
| accepted values | Giá trị nằm trong tập cho phép | status ∈ {active, closed, frozen} |
| referential (relationships) | FK tồn tại ở bảng cha | mọi account_id trong txn có ở dim_account |
Khác biệt then chốt: unit test bắt lỗi code; data test bắt lỗi dữ liệu (nguồn bẩn, edge case ngoài dự đoán, drift theo thời gian) — loại lỗi mà code đúng vẫn để lọt. Cả hai đều cần; xem thêm DataOps — Testing & Quality cho tổ chức hai lớp này trong quy trình.
Reconciliation — đối soát tổng với nguồn
Data test kiểm hình dạng dữ liệu; reconciliation (đối soát) kiểm độ chính xác bằng cách so tổng/số đếm giữa đích và nguồn tin cậy. Đây là kiểm quan trọng nhất trong ngân hàng vì nó là bằng chứng "không mất, không thêm, không lệch" khi dữ liệu đi qua nhiều chặng.
Ba mức đối soát thường dùng:
- Row count — số bản ghi đích khớp nguồn (bắt mất/nhân đôi bản ghi).
- Control total — tổng một cột số (tổng số tiền, tổng dư nợ) khớp giữa nguồn và đích; cho phép sai số làm tròn có ngưỡng.
- Hash/checksum theo khóa — với dữ liệu nhạy cảm, so hash từng nhóm để bắt lệch giá trị mà tổng vẫn tình cờ bằng nhau.
-- Đối soát control total: tổng số tiền giao dịch DW vs staging nguồn theo ngày
WITH src AS (
SELECT txn_date, COUNT(*) AS n_src, SUM(amount) AS sum_src
FROM staging.core_transactions
WHERE txn_date = DATE '2026-07-20'
GROUP BY txn_date
),
dw AS (
SELECT txn_date, COUNT(*) AS n_dw, SUM(amount) AS sum_dw
FROM warehouse.fact_transaction
WHERE txn_date = DATE '2026-07-20'
GROUP BY txn_date
)
SELECT s.txn_date,
s.n_src, d.n_dw, (d.n_dw - s.n_src) AS row_diff,
s.sum_src, d.sum_dw, (d.sum_dw - s.sum_src) AS amount_diff
FROM src s JOIN dw d USING (txn_date)
WHERE d.n_dw <> s.n_src
OR ABS(d.sum_dw - s.sum_src) > 0.01; -- ngưỡng sai số làm tròn
-- Trả về 0 dòng = đối soát khớp; có dòng = pipeline FAIL, phải chặn.
Đối soát nên chạy như một test có ngưỡng trong DAG, không phải báo cáo người xem sau. Xem thêm gov-03 — Data Quality mô tả đối soát số dư nguồn–DW chi tiết hơn.
Quality gate — chặn pipeline khi fail
Quality gate là điểm quyết định trong DAG: chỉ khi các test pass thì dữ liệu mới được đẩy tiếp xuống hạ nguồn; nếu fail, pipeline dừng và không có dữ liệu bẩn nào được publish. Ở Batch 7 ta đã đặt quality gate như một task trong DAG với luồng rẽ nhánh; bài này nói quyết định gì khi gate fail.
Một pattern rất hữu ích là Write–Audit–Publish (WAP): ghi dữ liệu vào vùng tạm/nhánh ẩn (write), chạy test trên vùng đó (audit), chỉ khi pass mới hoán đổi/công bố cho người dùng (publish). Người dùng không bao giờ nhìn thấy dữ liệu chưa qua audit.
Điểm tinh tế: không phải mọi fail đều nên chặn toàn bộ pipeline. Cần phân mức nghiêm trọng:
- Error/blocking (ví dụ reconciliation lệch, khóa NULL): chặn cứng — thà không có báo cáo còn hơn báo cáo sai. Giữ nguyên bản gold cũ (bản tốt gần nhất), báo alert.
- Warn/non-blocking (ví dụ tỷ lệ NULL một cột phụ vượt ngưỡng nhẹ): cho chảy tiếp nhưng ghi cảnh báo để điều tra.
Trong dbt, đây chính là severity: error vs severity: warn; trong Great Expectations là kết quả validation gắn với action chặn hay chỉ log.
Circuit breaker & quarantine dữ liệu xấu
Khi một feed nguồn hỏng nặng (ví dụ file EOD của một hệ thống thẻ về rỗng hoặc lệch định dạng), có hai cơ chế phòng thủ bổ sung cho quality gate:
Circuit breaker — nếu số lỗi vượt ngưỡng, "ngắt mạch" toàn bộ nhánh pipeline liên quan thay vì cứ cố xử lý từng bản ghi và làm hỏng bảng hạ nguồn. Giống cầu dao điện: một cú ngắt dứt khoát an toàn hơn để dòng lỗi âm ỉ lan ra. Ví dụ: nếu >5% giao dịch thiếu account_id, dừng luôn, không nạp vào fact — vì lỗi mức đó thường là sự cố nguồn, xử lý từng dòng sẽ ra bảng sai lệch.
Quarantine (cách ly) — khi chỉ một phần dữ liệu xấu, tách các bản ghi vi phạm ra một bảng cách ly kèm lý do, cho phần tốt chảy tiếp. Lợi ích: (1) pipeline không chết vì vài dòng lỗi; (2) dữ liệu xấu được lưu lại để điều tra và reprocess/backfill sau khi sửa nguồn, thay vì mất luôn.
-- Tách dòng vi phạm sang bảng quarantine kèm lý do, giữ dòng tốt cho hạ nguồn
INSERT INTO quarantine.transactions_bad
SELECT t.*, CURRENT_TIMESTAMP AS quarantined_at,
CASE
WHEN t.account_id IS NULL THEN 'missing_account_id'
WHEN t.amount < 0 THEN 'negative_amount'
WHEN t.currency NOT IN ('VND','USD','EUR') THEN 'invalid_currency'
END AS reason
FROM staging.core_transactions t
WHERE t.account_id IS NULL
OR t.amount < 0
OR t.currency NOT IN ('VND','USD','EUR');
-- Chỉ nạp dòng hợp lệ vào bảng sạch
INSERT INTO warehouse.fact_transaction
SELECT t.* FROM staging.core_transactions t
WHERE t.account_id IS NOT NULL
AND t.amount >= 0
AND t.currency IN ('VND','USD','EUR');
Nguyên tắc: quarantine phải luôn kèm đếm và cảnh báo. Loại bỏ âm thầm dòng xấu là phản pattern nguy hiểm — báo cáo trông "sạch" nhưng thiếu dữ liệu mà không ai biết. Số dòng bị quarantine phải xuất hiện trên dashboard và kích alert khi vượt ngưỡng.
Công cụ: dbt tests, Great Expectations, Soda
Ba công cụ phổ biến, khác triết lý; đều trung lập với engine (chạy được trên nhiều warehouse), chọn theo bối cảnh:
| Công cụ | Triết lý | Hợp khi |
|---|---|---|
| dbt tests | Test khai báo trong YAML cạnh model, chạy cùng transform ELT | Bạn đã dùng dbt cho transform; muốn test đi liền model (xem dbt-04) |
| Great Expectations (GX) | "Expectation Suite" phong phú + Data Docs; kiểm ở ranh giới ingest/pipeline | Cần thư viện expectation đa dạng, tài liệu chất lượng tự sinh, kiểm cả ngoài dbt |
| Soda | Ngôn ngữ khai báo SodaCL gọn, thiên monitoring & alert | Muốn khai báo check ngắn gọn, tích hợp cảnh báo/observability |
dbt có sẵn 4 test generic (not_null, unique, accepted_values, relationships) — đúng bốn nhóm data test ở trên — và cho viết test tuỳ biến bằng SQL/macro:
# schema.yml — data test khai báo ngay cạnh model dbt
models:
- name: fact_transaction
columns:
- name: transaction_id
tests:
- not_null
- unique
- name: account_id
tests:
- not_null
- relationships: # referential integrity
to: ref('dim_account')
field: account_id
- name: currency
tests:
- accepted_values:
values: ['VND', 'USD', 'EUR']
config:
severity: error # fail = chặn (vs warn)
Great Expectations diễn đạt cùng ý bằng "expectation" — mỗi expectation là một khẳng định về dữ liệu, gom thành suite:
# Great Expectations — một expectation suite tối giản cho bảng giao dịch
validator.expect_column_values_to_not_be_null("account_id")
validator.expect_column_values_to_be_unique("transaction_id")
validator.expect_column_values_to_be_in_set("currency", ["VND", "USD", "EUR"])
validator.expect_column_values_to_be_between("amount", min_value=0)
result = validator.validate()
if not result["success"]:
raise ValueError("Quality gate FAIL — chặn publish") # để orchestrator bắt
Dù công cụ nào, ba nguyên tắc chung: (1) test là code, versioned cùng pipeline; (2) kết quả test gắn với hành động (chặn/quarantine/alert), không chỉ để xem; (3) test chạy tự động trong DAG, không phụ thuộc con người nhớ chạy. Xem thêm DataOps adv — Quality & Contracts về data contract ràng buộc kỳ vọng giữa producer và consumer.
Giám sát & cảnh báo
Test chặn được lỗi biết trước; giám sát (monitoring) phát hiện điều bất thường mà không ai nghĩ tới viết test. Vài tín hiệu nên theo dõi liên tục:
- Freshness — bảng cập nhật lần cuối lúc nào; feed EOD có trễ SLA không (chiều timeliness).
- Volume — số bản ghi mỗi lần chạy; sụt/tăng đột biến so với đường cơ sở là dấu hiệu nguồn hỏng.
- Distribution drift — phân phối cột số (min/max/mean/tỷ lệ NULL) lệch bất thường so với lịch sử.
- Tỷ lệ test fail & số dòng quarantine theo thời gian — xu hướng tăng báo hiệu chất lượng nguồn xấu đi.
Nguyên tắc cảnh báo: alert phải hướng hành động và có người chịu trách nhiệm. Alert phân mức (blocking → gọi on-call ngay; warn → gom báo cáo sáng hôm sau), tránh "alert fatigue" khiến cảnh báo thật bị bỏ qua. Kết hợp với observability/incident của pipeline ở DataOps — Observability & Incident.
Use case thực tế
Bối cảnh (minh hoạ, số liệu giả định): Pipeline EOD của NCB gom ~5,2 triệu giao dịch/ngày từ core banking, hệ thống thẻ và khoản vay để dựng fact_transaction và bảng phân loại nợ nộp báo cáo. Yêu cầu: dữ liệu phải chính xác tuyệt đối trước 06:00 để kịp báo cáo.
Áp dụng khung bài này:
- Ingest checks: mỗi feed kiểm schema + đếm bản ghi so file control; feed thẻ về 4,98 triệu trong khi control ghi 5,01 triệu → lệch 0,6%.
- Circuit breaker: ngưỡng lệch cho phép ở mức nạp file là 0,1%; 0,6% vượt ngưỡng → ngắt nhánh feed thẻ, không nạp vào fact, giữ bản gold hôm trước.
- Alert: on-call nhận cảnh báo blocking lúc 02:14, phát hiện nguồn thẻ xuất thiếu một batch giờ chót.
- Reprocess: nguồn xuất lại đủ file; chạy backfill riêng partition ngày đó; reconciliation control total khớp (chênh 0,00 sau làm tròn) → quality gate pass → publish.
- Quarantine song song: 1.240 giao dịch thiếu
account_id(lỗi ánh xạ) bị tách sang bảng cách ly, phần còn lại vẫn publish; đội nguồn sửa ánh xạ hôm sau rồi nạp bù.
Kết quả (minh hoạ): báo cáo phân loại nợ ra đúng hạn với con số đã đối soát, không có bản ghi sai lọt xuống báo cáo NHNN; 1.240 dòng lỗi được truy vết và bù thay vì mất âm thầm. Chi tiết pipeline EOD ngân hàng ở Batch 9 — Banking EOD.
Ghi nhớ
- Data quality là control gắn trong pipeline, kiểm ở nhiều chặng (ingest/transform/publish), không phải bước dọn một lần ở cuối.
- Sáu chiều: completeness, accuracy, consistency, timeliness, uniqueness, validity — mỗi chiều một câu hỏi và một loại test; đừng nhầm accuracy (khớp sự thật ngoài) với consistency (khớp nội bộ).
- Hai lớp kiểm thử: unit test bắt lỗi code transform trên dữ liệu giả; data test (not null/unique/accepted values/referential) bắt lỗi dữ liệu thật. Cần cả hai.
- Reconciliation (row count / control total / checksum) là bằng chứng accuracy — chạy như test có ngưỡng, không phải báo cáo xem sau.
- Quality gate theo pattern Write–Audit–Publish chỉ công bố khi pass; phân mức fail: blocking chặn cứng, warn cho qua có ghi log.
- Circuit breaker ngắt nhánh khi lỗi vượt ngưỡng; quarantine tách dòng xấu (kèm lý do + đếm + alert) để reprocess sau — không loại âm thầm.
- Công cụ trung lập engine: dbt tests (test cạnh model), Great Expectations (suite + Data Docs), Soda (SodaCL, thiên monitoring); test là code versioned, gắn với hành động, chạy tự động trong DAG.
- Giám sát freshness/volume/drift bắt lỗi ngoài dự đoán; alert phân mức, hướng hành động, có người chịu trách nhiệm.
Nguồn tham khảo
- Fundamentals of Data Engineering — Joe Reis & Matt Housley, O'Reilly (chương data quality & vòng đời kỹ thuật dữ liệu).
- Designing Data-Intensive Applications — Martin Kleppmann, O'Reilly (tính đúng đắn, ràng buộc dữ liệu).
- DAMA-DMBOK, Data Management Body of Knowledge — chương Data Quality (các chiều chất lượng).
- dbt — tài liệu Tests & data tests: https://docs.getdbt.com/docs/build/data-tests
- Great Expectations — tài liệu chính thức: https://docs.greatexpectations.io/
- Soda — tài liệu SodaCL & checks: https://docs.soda.io/
- Netflix Tech Blog — "Write-Audit-Publish" pattern cho data pipeline chất lượng cao.
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ẻ!