การทดสอบโค้ด Airflow

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

Volker Janz

Senior Developer Advocate at Astronomer

การทดสอบ 3 ระดับ

พีระมิดการทดสอบ DAG

  • Integrity: DAG โหลดได้โดยไม่มีข้อผิดพลาดหรือไม่?
  • Unit: Business logic ให้ผลลัพธ์ที่ถูกต้องหรือไม่?
  • Integration: DAG ทำงานได้ครบทุกขั้นตอนหรือไม่?
การสร้าง Data Pipeline ด้วย Airflow

Integrity test ด้วย DagBag

from airflow.models import DagBag

dag_bag = DagBag(include_examples=False)
def test_no_import_errors(): assert len(dag_bag.import_errors) == 0
def test_dag_loaded(): assert "daily_etl" in dag_bag.dags
การสร้าง Data Pipeline ด้วย Airflow

สาเหตุที่ DAG พังตอน import

ModuleNotFoundError

  • ModuleNotFoundError: แพ็กเกจ provider หายหรือ import path ผิด
  • NameError: ตัวแปร ฟังก์ชัน หรือชื่อสะกดผิด
  • ImportError: มี circular import ระหว่างไฟล์ DAG
  • Integrity test ช่วยป้องกันปัญหาเหล่านี้
การสร้าง Data Pipeline ด้วย Airflow

Unit test สำหรับฟังก์ชัน task

ในไฟล์ DAG:

def clean_record(record):
    return {
        "name": record["name"].strip(),
        "email": record["email"].lower(),
    }

@task
def transform(records):
    return [clean_record(r)
            for r in records]
  • แยก business logic ออกมา

ในไฟล์ทดสอบ:

from dags.data_cleaning import (
    clean_record,
)

def test_strips_whitespace():
    result = clean_record(
      {"name": "  Alice  ",
       "email": "[email protected]"}
    )
    assert result["name"] == "Alice"
  • เน้น unit test ที่ business logic
การสร้าง Data Pipeline ด้วย Airflow

Integration test ด้วย dag.test()

import pytest
from airflow.models import DagBag
from pendulum import datetime

dag_bag = DagBag(include_examples=False)

def test_etl_pipeline(): dag = dag_bag.get_dag("etl_output") assert dag is not None dag.test(logical_date=datetime(2026, 1, 15)) output = Path("/tmp/etl_results.json") assert output.exists() results = json.loads(output.read_text()) assert len(results) == 2
  • รัน DAG จริงด้วย input ที่กำหนด แล้วตรวจสอบผลลัพธ์
การสร้าง Data Pipeline ด้วย Airflow

การทดสอบใน CI

CI pipeline

  • Integrity + unit: ทุก commit (เร็ว)
  • Integration: pull request หรือรายคืน (ช้ากว่า)
  • ไม่มี DAG ใดขึ้น production โดยไม่ผ่านทั้ง 3 ระดับ
การสร้าง Data Pipeline ด้วย Airflow

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

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

Preparing Video For Download...