Testarea codului Airflow

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Trei niveluri de testare

Piramida de testare a DAG-ului

  • Integritate: DAG-ul se încarcă fără erori?
  • Unitate: logica de business produce rezultate corecte?
  • Integrare: DAG-ul rulează complet end-to-end?
Construirea pipeline-urilor de date cu Airflow

Teste de integritate cu 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
Construirea pipeline-urilor de date cu Airflow

De ce eșuează DAG-urile la import

ModuleNotFoundError

  • ModuleNotFoundError: pachet provider lipsă sau cale de import incorectă
  • NameError: variabilă sau funcție redenumită, sau greșeală de scriere
  • ImportError: importuri circulare între fișierele DAG
  • Testele de integritate ajută la evitarea acestor probleme
Construirea pipeline-urilor de date cu Airflow

Testarea de unitate a funcțiilor de task

În fișierul 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]
  • Extrage logica de business

În fișierul de teste:

from dags.data_cleaning import (
    clean_record,
)

def test_strips_whitespace():
    result = clean_record(
      {"name": "  Alice  ",
       "email": "[email protected]"}
    )
    assert result["name"] == "Alice"
  • Concentrează testele de unitate pe logica de business
Construirea pipeline-urilor de date cu Airflow

Teste de integrare cu 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
  • Rulează DAG-ul efectiv cu date de intrare controlate și validează rezultatul
Construirea pipeline-urilor de date cu Airflow

Testarea în CI

Pipeline CI

  • Integritate + unitate: la fiecare commit (rapid)
  • Integrare: la pull request-uri sau nocturn (mai lent)
  • Niciun DAG nu ajunge în producție fără a trece toate cele trei niveluri
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...