Probar código de Airflow

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Tres niveles de testing

Pirámide de tests de Dag

  • Integridad: ¿carga el Dag sin errores?
  • Unitarias: ¿la lógica de negocio da resultados correctos?
  • Integración: ¿el Dag completo corre de extremo a extremo?
Creación de canalizaciones de datos con Airflow

Tests de integridad con 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
Creación de canalizaciones de datos con Airflow

Por qué fallan los Dags al importar

ModuleNotFoundError

  • ModuleNotFoundError: falta un paquete de provider o ruta de import incorrecta
  • NameError: variable o función renombrada, o un typo
  • ImportError: imports circulares entre archivos de Dag
  • Los tests de integridad ayudan a evitarlo
Creación de canalizaciones de datos con Airflow

Tests unitarios de funciones de tareas

En el archivo del 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]
  • Extrae la lógica de negocio

En el archivo de tests:

from dags.data_cleaning import (
    clean_record,
)

def test_strips_whitespace():
    result = clean_record(
      {"name": "  Alice  ",
       "email": "[email protected]"}
    )
    assert result["name"] == "Alice"
  • Centra las pruebas unitarias en la lógica de negocio
Creación de canalizaciones de datos con Airflow

Tests de integración con 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
  • Ejecuta el Dag real con entrada controlada y valida la salida
Creación de canalizaciones de datos con Airflow

Testing en CI

Canalización de CI

  • Integridad + unitarias: cada commit (rápidos)
  • Integración: pull requests o nightly (más lento)
  • Ningún Dag llega a producción sin pasar las tres
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...