데이터 품질 검사

Airflow로 데이터 파이프라인 구축하기

Volker Janz

Senior Developer Advocate at Astronomer

품질 게이트가 중요한 이유

 

  • 검사 없는 파이프라인은 신뢰할 수 없습니다
  • 업스트림 변경 하나가 null 또는 음수 값을 유발합니다
  • 하위 대시보드 전체에 잘못된 수치가 표시됩니다
  • 품질 게이트는 문제를 원점에서 차단합니다

잘못된 데이터가 하위로 흘러가는 모습

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}}, }, )
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 값의 개수를 셉니다 (null 없음은 equal_to: 0으로 검증)
  • distinct_check: 고유 값의 개수를 셉니다
  • 사용 가능한 검사: null_check, unique_check, distinct_check, min, max
  • 사용 가능한 비교 연산자: equal_to, greater_than, geq_to, less_than, leq_to
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 표현식
Airflow로 데이터 파이프라인 구축하기

품질 검사 연결하기

@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
  • 로드검증합니다
  • 검사가 실패하면 파이프라인이 중단됩니다
  • 하위 사용자는 잘못된 데이터를 볼 수 없습니다
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",
        },
    },
)
  • 테이블에 행이 있나요?
  • 합계가 예상 범위 내에 있나요?
Airflow로 데이터 파이프라인 구축하기

Astro로 대규모 품질 모니터링하기

$$

  • Astronomer의 Astro Observe로 파이프라인 상태를 모니터링합니다
  • 마감 전 알림을 제공하는 SLA 추적
  • DAG와 테이블에 걸친 파이프라인 계보
  • 대규모 데이터 품질 모니터링
  • AI 기반 근본 원인 분석

Astro Observe 계보

Airflow로 데이터 파이프라인 구축하기

연습해 봅시다!

Airflow로 데이터 파이프라인 구축하기

Preparing Video For Download...