Conceptos clave de Airflow

Creación de canalizaciones de datos con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Conoce a tu instructor

  Foto de perfil de Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ años como data engineer en gaming
  • Trabaja con Airflow desde la versión 1.x
  • Ponente, mentor y responsable del boletín en Data Engineer Things
Creación de canalizaciones de datos con Airflow

Lo que vas a construir

 

  • Crear Dags con la TaskFlow API
  • Construir flujos dinámicos con task mapping y planificación basada en activos
  • Gestionar fallos con reintentos y callbacks
  • Ejecutar cargas de trabajo SQL con Airflow

Visualización del contenido de cada capítulo

Creación de canalizaciones de datos con Airflow

Antes de empezar

 

$$

  • Comodidad con Dags, tareas y operadores
  • Familiaridad con nociones básicas de planificación

Introduction to Airflow - course page

Creación de canalizaciones de datos con Airflow

Repaso rápido

from airflow.sdk import dag, task

@dag
def star_wars_dag():


@task def get_star_wars_person(): import requests return requests.get("https://swapi.dev/api/people/1/").json()
@task.bash def print_name(person): return f"echo '{person['name']}'"
person = get_star_wars_person() print_name(person) star_wars_dag()
  • Un Dag es una colección de tareas con dependencias
  • Las tareas son unidades de trabajo individuales
  • Operadores/decoradores definen qué hace cada tarea
  • Las dependencias definen el orden de ejecución

$$

Dag simple

Creación de canalizaciones de datos con Airflow

Arquitectura de Airflow

$$

Arquitectura de Airflow 3

 

$$

  • Orquestación: Scheduler, Dag Processor
  • Ejecución: Worker, Triggerer
  • Interfaz y almacenamiento: API Server, Metadata DB
Creación de canalizaciones de datos con Airflow

Enfoques de planificación

# Sin ejecuciones automáticas: lánzalo manualmente (por defecto)
@dag(schedule=None)
def my_pipeline(): ...

# Basado en tiempo: se ejecuta cada día a las 6 a. m. @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Consciencia de datos: se ejecuta cuando se actualiza un Asset @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Creación de canalizaciones de datos con Airflow

Dos formas de escribir Dags

Operadores clásicos

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Elige esta opción cuando no haya decoradores disponibles

TaskFlow API

@task
def extract():
    return {"users": 150}

data = extract()
transform(data)
  • Decoradores de Python sencillos
  • Menos código boilerplate
  • Puede combinarse con operadores clásicos
Creación de canalizaciones de datos con Airflow

Ejercicios en este curso

$$

Captura del ejercicio en el IDE

 

  • Ejercicios en IDE: edita archivos .py reales
  • Haz clic en "Run this file" o usa python3 filename.py
  • dag.test() ejecuta el Dag completo en un solo proceso
Creación de canalizaciones de datos con Airflow

¡Vamos a practicar!

Creación de canalizaciones de datos con Airflow

Preparing Video For Download...