Kiểm tra chất lượng dữ liệu

Xây dựng Data Pipeline với Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Vì sao cổng chất lượng quan trọng

 

  • Pipeline không có kiểm tra là pipeline không thể tin cậy
  • Một thay đổi upstream có thể tạo ra null hoặc giá trị âm
  • Mọi dashboard downstream đều hiển thị sai số
  • Cổng chất lượng bắt lỗi ngay tại nguồn

Dữ liệu xấu chảy xuống hạ nguồn

Xây dựng Data Pipeline với Airflow

SQLColumnCheckOperator

from airflow.providers.common.sql.operators.sql import
    (SQLExecuteQueryOperator, SQLColumnCheckOperator)

check_columns = SQLColumnCheckOperator( task_id="check_columns", conn_id="duckdb_analytics", table="daily_summary", column_mapping={ "total_revenue": {"min": {"greater_than": 0}}, "total_orders": {"null_check": {"equal_to": 0}}, }, )
Xây dựng Data Pipeline với Airflow

Các loại kiểm tra cột

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: xác thực giá trị nhỏ nhất trong một cột
  • null_check: đếm số giá trị null (khẳng định equal_to: 0 để không có null)
  • distinct_check: đếm số giá trị khác nhau
  • Kiểm tra hỗ trợ: null_check, unique_check, distinct_check, min, max
  • Toán tử so sánh: equal_to, greater_than, geq_to, less_than, leq_to
Xây dựng Data Pipeline với Airflow

SQLTableCheckOperator

from airflow.providers.common.sql.operators.sql import SQLTableCheckOperator

check_table = SQLTableCheckOperator( task_id="check_table", conn_id="duckdb_analytics", table="daily_summary", checks={ "row_count_check": { "check_statement": "COUNT(*) > 0", }, }, )
  • Xác thực các thuộc tính của toàn bộ bảng
  • check_statement: bất kỳ biểu thức SQL nào trả về đúng hoặc sai
Xây dựng Data Pipeline với Airflow

Xâu chuỗi các kiểm tra chất lượng

@dag(schedule="@daily", template_searchpath="/path/to/sql")
def sales_pipeline():

    load = SQLExecuteQueryOperator(
        task_id="load_daily_sales",
        conn_id="duckdb_analytics",
        sql="aggregate_sales.sql",
    )

    check_columns = SQLColumnCheckOperator(...)
    check_table = SQLTableCheckOperator(...)

    load >> check_columns >> check_table
  • Nạp trước, rồi xác thực
  • Nếu kiểm tra nào thất bại, pipeline sẽ dừng
  • Người dùng downstream không bao giờ thấy dữ liệu xấu
Xây dựng Data Pipeline với Airflow

Kiểm tra cột vs. bảng

SQLColumnCheckOperator

check_columns = SQLColumnCheckOperator(
    task_id="check_columns",
    conn_id="duckdb_analytics",
    table="daily_summary",
    column_mapping={
        "total_revenue": {"min": {"greater_than": 0}},
        "total_orders": {"null_check": {"equal_to": 0}},
    },
)
  • null không?
  • Giá trị có nằm trong phạm vi không?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • Bảng có dòng không?
  • Tổng có nằm trong ngưỡng kỳ vọng không?
Xây dựng Data Pipeline với Airflow

Giám sát chất lượng ở quy mô lớn với Astro

$$

  • Astronomer Astro Observe giúp theo dõi tình trạng pipeline
  • Theo dõi SLA với cảnh báo trước hạn
  • Dòng dõi pipeline qua DAG và bảng
  • Giám sát chất lượng dữ liệu ở quy mô lớn
  • Phân tích nguyên nhân gốc dựa trên AI

Dòng dõi trong Astro Observe

Xây dựng Data Pipeline với Airflow

Cùng luyện tập nào!

Xây dựng Data Pipeline với Airflow

Preparing Video For Download...