การตรวจสอบคุณภาพข้อมูล

การสร้าง Data Pipeline ด้วย Airflow

Volker Janz

Senior Developer Advocate at Astronomer

ทำไม Quality Gate จึงสำคัญ

 

  • Pipeline ที่ไม่มีการตรวจสอบคือ Pipeline ที่ ไม่น่าเชื่อถือ
  • การเปลี่ยนแปลงต้นทางอาจทำให้เกิด null หรือ ค่าติดลบ
  • Dashboard ปลายทางทุกอันจะแสดง ตัวเลขผิดพลาด
  • Quality gate จะดักปัญหา ตั้งแต่ต้นทาง

ข้อมูลเสียที่ไหลไปยังปลายทาง

การสร้าง Data Pipeline ด้วย 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}}, }, )
การสร้าง Data Pipeline ด้วย Airflow

ประเภทการตรวจสอบคอลัมน์

column_mapping={
    "total_revenue": {"min": {"greater_than": 0}},
    "total_orders": {"null_check": {"equal_to": 0}},
    "order_date": {"distinct_check": {"greater_than": 0}},
}
  • min: ตรวจสอบค่าต่ำสุดในคอลัมน์
  • null_check: นับค่า null (ใช้ equal_to: 0 เพื่อยืนยันว่าไม่มี null)
  • distinct_check: นับค่าที่ไม่ซ้ำกัน
  • การตรวจสอบที่ใช้ได้: null_check, unique_check, distinct_check, min, max
  • ตัวเปรียบเทียบที่ใช้ได้: equal_to, greater_than, geq_to, less_than, leq_to
การสร้าง Data Pipeline ด้วย 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", }, }, )
  • ตรวจสอบคุณสมบัติของ ตารางทั้งหมด
  • check_statement: นิพจน์ SQL ใดก็ได้ที่ให้ผลลัพธ์เป็น true หรือ false
การสร้าง Data Pipeline ด้วย Airflow

การเชื่อมต่อ Quality Check

@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
  • โหลดข้อมูลก่อน แล้วจึงตรวจสอบ
  • หากการตรวจสอบใดล้มเหลว Pipeline จะหยุดทำงาน
  • ผู้ใช้ปลายทางไม่มีทางเห็นข้อมูลที่ผิดพลาด
การสร้าง Data Pipeline ด้วย Airflow

การตรวจสอบระดับคอลัมน์ vs ระดับตาราง

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 อยู่หรือไม่?
  • ค่าอยู่ ในช่วงที่กำหนด หรือไม่?

SQLTableCheckOperator

check_table = SQLTableCheckOperator(
    task_id="check_table",
    conn_id="duckdb_analytics",
    table="daily_summary",
    checks={
        "row_count_check": {
            "check_statement": "COUNT(*) > 0",
        },
    },
)
  • ตารางมี แถวข้อมูล หรือไม่?
  • ค่ารวมอยู่ใน ช่วงที่คาดไว้ หรือไม่?
การสร้าง Data Pipeline ด้วย Airflow

การติดตามคุณภาพในระดับ Scale ด้วย Astro

$$

  • Astro Observe ของ Astronomer ช่วยติดตามสุขภาพของ Pipeline
  • การติดตาม SLA พร้อมแจ้งเตือนก่อนถึงกำหนด
  • Pipeline lineage ข้าม DAG และตาราง
  • การตรวจสอบคุณภาพข้อมูล ในระดับ Scale
  • การวิเคราะห์สาเหตุต้นตอ ด้วย AI

Astro Observe lineage

การสร้าง Data Pipeline ด้วย Airflow

มาฝึกกันเถอะ!

การสร้าง Data Pipeline ด้วย Airflow

Preparing Video For Download...