# Thực hành: dựng bốn lớp cảnh báo cho dữ liệu khách đến trễ, thiếu dòng, đổi schema hoặc ngừng chạy

> File đơn hàng của khách có thể đến trễ, bị cắt mất một nửa hay bị đổi tên cột mà pipeline vẫn báo chạy xanh. Hướng dẫn này giúp bạn bắt được lỗi trước khi khách tự phát hiện.

Bản gốc: https://fdetimes.net/vi/bach-khoa/thuc-hanh-giam-sat-pipeline-du-lieu/

Pipeline báo chạy xanh không có nghĩa là dữ liệu đúng. Job vẫn có thể chạy thành công trên file của hôm qua, trên một file chỉ có một nửa số dòng, hoặc trên một file mà khách vừa đổi tên cột `amount` thành `total_amount`. Lúc đó dashboard vẫn hiện số liệu, chỉ là số liệu sai.

Với FDE, đây là loại sự cố gây mất niềm tin nhanh nhất. Khách không cần biết hệ thống của bạn có bao nhiêu service. Họ chỉ nhớ lần đầu tiên chính họ phát hiện ra số liệu sai trước bạn. Vì thế, nếu được giao một pipeline nhận dữ liệu từ khách, đây là thứ nên dựng sớm.

Hướng dẫn này đi qua một tình huống giả định. Một chuỗi bán lẻ gửi bảng `orders` vào warehouse mỗi đêm và gửi danh mục sản phẩm mỗi tuần. Bạn sẽ dựng ba lớp kiểm tra cho ba kiểu hỏng: đến trễ, thiếu dòng và đổi schema. Sau đó bạn thêm lớp thứ tư làm lưới an toàn cho trường hợp job không chạy chút nào.

## Bạn cần chuẩn bị gì?

Bạn cần một dự án dbt đã khai báo source, một môi trường Python có cài Great Expectations, và một công cụ giám sát. Công cụ đó có thể là Prometheus, Datadog hoặc Sentry, tùy khách đang dùng gì. Không cần đủ cả ba. Ở khách hàng thật, bạn nên dùng công cụ mà đội vận hành của họ đã quen theo dõi.

Điều kiện quan trọng nhất nằm ở dữ liệu: mỗi bảng phải có một cột ghi thời điểm dòng đó được nạp vào, ví dụ `_etl_loaded_at`. Nếu chưa có cột này, việc đầu tiên là thêm nó vào bước ingest.

## Bước 1: freshness trong dbt bắt dữ liệu đến trễ

Tài liệu của dbt mô tả cấu hình `freshness` là nơi bạn khai báo dữ liệu nguồn cần mới đến mức nào. Hai ngưỡng chính là `warn_after`, tức tuổi tối đa của bản ghi mới nhất trước khi dbt cảnh báo, và `error_after`, tức mốc mà dbt báo lỗi. Mỗi ngưỡng cần có cả `count` lẫn `period`.

```yaml
sources:
- name: retail_client
loaded_at_field: _etl_loaded_at
freshness:
warn_after: {count: 12, period: hour}
error_after: {count: 24, period: hour}
tables:
- name: orders
freshness:  # chặt hơn cho orders
warn_after: {count: 6, period: hour}
error_after: {count: 12, period: hour}
filter: _etl_loaded_at >= current_date - 3
- name: product_catalog
freshness: null
```

Đoạn trên đã được đơn giản hóa. Vị trí khóa `freshness` có thể khác nhau giữa các phiên bản dbt, còn cú pháp ngày tháng trong `filter` phụ thuộc vào warehouse. Cấu hình này có ba chi tiết bạn nên để ý.

**Ngưỡng riêng cho từng bảng.** dbt cho phép đặt ngưỡng theo bảng và cho phép tắt kiểm tra. Vì thế `orders` bị kiểm chặt hơn, còn `product_catalog`, vốn chỉ cập nhật hằng tuần, được bỏ qua.

**`filter` để kiểm tra rẻ.** Khóa này thêm một mệnh đề `WHERE` vào truy vấn freshness, nên dbt không phải quét toàn bộ một bảng lớn.

**`loaded_at_field` trỏ vào thời điểm nạp.** Đây là chỗ hay bị bỏ sót nhất. dbt cần cột này để biết dữ liệu đến lần cuối khi nào, trong trường hợp adapter không cung cấp metadata đó, và bạn nên trỏ nó vào thời điểm nạp chứ không phải thời điểm phát sinh đơn hàng.

Giả sử khách gửi bù đơn của ba ngày trước: nếu dùng thời điểm phát sinh, dữ liệu trông như đã cũ ba ngày; nếu dùng thời điểm nạp, bạn mới thấy đúng lúc file thực sự đến.

**Kiểm tra:** chạy `dbt source freshness`, sau đó tạm hạ `warn_after` xuống 1 giờ để chắc chắn bạn thấy được trạng thái cảnh báo.

## Bước 2: số dòng bất thường thường báo hiệu file bị cắt

Freshness chỉ cho biết dữ liệu có đến hay không. Nó không cho biết dữ liệu có đủ không. Nếu file xuất bị cắt giữa chừng, bản ghi mới nhất vẫn có timestamp mới và kiểm tra freshness vẫn qua.

Great Expectations có expectation kiểm tra số dòng của bảng có nằm trong một khoảng hay không, và khoảng này tính cả hai giá trị biên. `min_value` có thể đặt một mình, nên bạn chỉ cần cận dưới để phát hiện thiếu dòng.

```python
import great_expectations as gx

row_check = gx.expectations.ExpectTableRowCountToBeBetween(
min_value=30000,
max_value=None,
)
```

Đoạn này đã lược bỏ phần kết nối batch và chạy validation. Hãy thử tính con số với dữ liệu giả định. Nếu 14 đêm gần nhất `orders` dao động từ 38.000 đến 45.000 dòng, mốc 30.000 vẫn để lại khoảng dư cho một đêm vắng khách. Nhưng nếu file đêm nay chỉ có 19.000 dòng vì bị cắt đôi, kiểm tra sẽ báo lỗi.

Ngưỡng cố định sẽ gặp khó khi lượng dữ liệu thay đổi theo ngày trong tuần, chẳng hạn cuối tuần gấp đôi ngày thường. Datadog có anomaly monitor học từ dữ liệu lịch sử để phát hiện hành vi bất thường, phù hợp hơn với kiểu dữ liệu này.

Datadog cũng có monitor cho freshness, row count và các metric ở cấp cột trên nhiều warehouse, nghĩa là cả ba kiểu hỏng có thể được theo dõi trên cùng một công cụ.

**Kiểm tra:** chạy expectation trên một bản sao bảng đã bị xóa một nửa số dòng. Kiểm tra phải trả về thất bại.

## Bước 3: đổi schema cần một danh sách cột cam kết

Hai kiểm tra ở trên chỉ lo độ mới và số dòng. Với schema, cách đơn giản nhất là tự viết một kiểm tra nhỏ. Đoạn dưới đây là ví dụ minh họa, chưa phải giải pháp hoàn chỉnh: lưu danh sách cột đã thống nhất với khách, rồi so sánh với cột thực tế của mỗi lần nạp.

```python
EXPECTED = {"order_id", "store_id", "amount", "created_at", "_etl_loaded_at"}

def schema_diff(actual_columns):
actual = set(actual_columns)
return {"missing": EXPECTED - actual, "new": actual - EXPECTED}
```

Một cột bị mất nên được xử lý như lỗi. Một cột mới chỉ nên tạo cảnh báo, vì nhiều khả năng khách vừa bổ sung thông tin chứ không làm hỏng gì. Nếu khách dùng Datadog, bạn có thể chuyển kiểm tra này thành một custom SQL monitor, loại monitor mà tài liệu của Datadog có liệt kê.

**Kiểm tra:** đổi tên một cột trên bảng thử nghiệm. `missing` phải chứa tên cũ và `new` phải chứa tên mới.

## Bước 4: lưới an toàn khi job không chạy

Cả ba bước trên đều cần job được chạy thì mới có kết quả. Nếu cron chết, sẽ không có kiểm tra nào báo lỗi, đơn giản vì không có kiểm tra nào được thực thi.

Prometheus thu thập metric theo cơ chế pull và dùng mô hình dữ liệu đa chiều, trong đó mỗi chuỗi thời gian được xác định bằng tên metric và các cặp key/value. Pipeline có thể xuất ra một metric ghi thời điểm nạp thành công gần nhất.

Hàm `absent()` hoạt động ngược với trực giác: khi chuỗi đó vẫn còn dữ liệu, nó trả về rỗng; chỉ khi chuỗi đã biến mất, nó mới trả về kết quả. Vì vậy cảnh báo dựa trên `absent()` chỉ bật lên khi feed ngừng báo cáo hẳn.

```yaml
- alert: OrdersFeedStale
expr: time() - pipeline_last_success_timestamp_seconds{feed="orders"} > 6 * 3600
for: 30m
- alert: OrdersFeedMissing
expr: absent(pipeline_last_success_timestamp_seconds{feed="orders"})
```

Đây là phiên bản rút gọn, chưa có label định tuyến và annotation. Nếu khách không dùng Prometheus, Sentry Cron Monitoring có thể theo dõi bất kỳ job định kỳ nào, kể cả job ingest kéo dữ liệu của khách.

## Lỗi hay gặp: cảnh báo quá nhiều

Hướng dẫn về alerting của Prometheus khuyên giữ càng ít cảnh báo càng tốt và chỉ cảnh báo trên các triệu chứng gắn với điều người dùng cuối thực sự gặp phải.

Tài liệu đó cũng khuyên chừa khoảng trễ trong ngưỡng để những dao động nhỏ không làm phát sinh cảnh báo. `for: 30m` và khoảng cách giữa `warn_after` với `error_after` trong các ví dụ trên được đặt ra vì lý do đó.

**Điểm mấu chốt:** Cảnh báo tốt đo thứ khách hàng sẽ thấy đau, không đo mọi thứ có thể đo.

Hai lỗi khác cũng rất phổ biến. Lỗi đầu tiên là áp cùng một ngưỡng cho mọi bảng, khiến bảng cập nhật hằng tuần báo đỏ mỗi ngày cho đến khi cả đội tắt thông báo.

Lỗi thứ hai là gửi cảnh báo vào một kênh không có người chịu trách nhiệm. Mỗi cảnh báo nên đi kèm tên người xử lý và một dòng hướng dẫn việc cần làm đầu tiên.

## Ở khách hàng thật, việc này bắt đầu từ một câu hỏi

Phần khó nhất không nằm ở YAML. Nó nằm ở việc hỏi khách: "File `orders` trễ bao lâu thì bên anh chị bắt đầu gặp vấn đề?" Câu trả lời của khách chính là `error_after`. Cách đặt ngưỡng này cho thấy bạn đang làm customer discovery chứ không chỉ cấu hình công cụ.

Khi đọc JD của một vị trí FDE hoặc data engineer làm việc trực tiếp với khách, hãy tìm các cụm như "data quality", "observability" hay "SLA với khách hàng". Trong CV, nên viết theo kết quả.

Ví dụ: "Dựng kiểm tra freshness, số dòng và schema cho 6 feed của khách, phát hiện file bị cắt trước khi báo cáo buổi sáng được gửi đi." Những con số trong câu này cần là số thật từ dự án của bạn.

Ngày đầu tiên ở một khách hàng mới, đừng cố dựng cả bốn lớp cùng lúc. Hãy bắt đầu với freshness cho feed quan trọng nhất, vì đó là kiểm tra rẻ nhất và giúp bạn biết trước khi khách phải gọi điện báo.

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

- Chọn một source trong dự án dbt bạn đang làm, thêm loaded_at_field cùng warn_after và error_after, rồi chạy dbt source freshness để xem kết quả.
- Lấy số dòng của một feed trong 14 ngày gần nhất, chọn cận dưới thấp hơn ngày ít dòng nhất, rồi viết thành một ExpectTableRowCountToBeBetween.
- Viết một câu trong CV mô tả cảnh báo dữ liệu bạn đã dựng, ghi rõ loại lỗi nó bắt được và thời gian phát hiện.

## Nguồn

- [freshness | dbt Developer Hub](https://docs.getdbt.com/reference/resource-properties/freshness)

- [expect table row count to be between | Great Expectations](https://greatexpectations.io/legacy/v1/expectations/expect_table_row_count_to_be_between/)

- [Monitor Types | Datadog Docs](https://docs.datadoghq.com/monitors/types/)

- [Overview | Prometheus](https://prometheus.io/docs/introduction/overview/)

- [Query functions | Prometheus](https://prometheus.io/docs/prometheus/latest/querying/functions/)

- [Alerting | Prometheus](https://prometheus.io/docs/practices/alerting/)

- [Crons | Sentry Docs](https://docs.sentry.io/product/crons/)
