Creación de canalizaciones de datos con Airflow
Volker Janz
Senior Developer Advocate at Astronomer

$$
$$

$$

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()
$$

$$

$$
# 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(): ...
Operadores clásicos
extract = PythonOperator(
task_id="extract",
python_callable=extract_fn)
extract >> transform
TaskFlow API
@task
def extract():
return {"users": 150}
data = extract()
transform(data)
$$

.py realespython3 filename.pydag.test() ejecuta el Dag completo en un solo procesoCreación de canalizaciones de datos con Airflow