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.
Mỗi kiểu hỏng cần một kiểm tra riêng, và một pipeline báo chạy xanh không bắt được kiểu nào trong số đó.
Đồ hoạ: FDE Times
Tóm tắt nhanh
- Dữ liệu khách thường hỏng theo ba kiểu: đến trễ, thiếu dòng và đổi schema. Mỗi kiểu cần một kiểm tra riêng, cộng thêm một lớp cho trường hợp job không chạy.
- Đo freshness bằng thời điểm nạp dữ liệu, không dùng thời điểm phát sinh sự kiện, và đặt ngưỡng riêng theo SLA của từng feed.
- Chỉ cảnh báo ở những triệu chứng mà khách sẽ thấy, và chừa khoảng trễ để vài phút chậm không làm ai bị gọi dậy lúc nửa đêm.
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.
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.
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.
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.
- 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 đó.
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.