Concetti chiave di Airflow

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Incontra il tuo istruttore

  Foto profilo di Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ anni come data engineer nel gaming
  • Lavora con Airflow dalla versione 1.x
  • Speaker, mentor e curatore della newsletter Data Engineer Things
Creare data pipeline con Airflow

Cosa costruirai

 

  • Scrivi Dags con la TaskFlow API
  • Crea workflow dinamici con task mapping e scheduling basato su asset
  • Gestisci gli errori con retry e callback
  • Esegui carichi SQL tramite Airflow

Visualizzazione dei contenuti di ogni capitolo

Creare data pipeline con Airflow

Prima di iniziare

 

$$

  • A tuo agio con Dags, task e operatori
  • Conoscenze di base di scheduling

Introduzione ad Airflow - pagina del corso

Creare data pipeline con Airflow

Ripasso rapido

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 è una raccolta di task con dipendenze
  • I task sono unità di lavoro autonome
  • Operatori/decorator definiscono cosa fa ogni task
  • Le dipendenze definiscono l'ordine di esecuzione

$$

Dag semplice

Creare data pipeline con Airflow

Architettura di Airflow

$$

Architettura di Airflow 3

 

$$

  • Orchestrazione: Scheduler, Dag Processor
  • Esecuzione: Worker, Triggerer
  • Interfaccia e storage: API Server, Metadata DB
Creare data pipeline con Airflow

Approcci di scheduling

# Nessuna esecuzione automatica: avvio manuale (predefinito)
@dag(schedule=None)
def my_pipeline(): ...

# Basata sul tempo: ogni giorno alle 6 @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Data-aware: parte quando si aggiorna un Asset @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Creare data pipeline con Airflow

Due modi per scrivere Dags

Operatori classici

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Scegli quando non ci sono decorator disponibili

TaskFlow API

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

data = extract()
transform(data)
  • Semplici decorator Python
  • Meno boilerplate
  • Combinabile con operatori classici
Creare data pipeline con Airflow

Esercizi in questo corso

$$

Screenshot esercizio in IDE

 

  • Esercizi in IDE: modifica veri file .py
  • Clicca "Run this file" o usa python3 filename.py
  • dag.test() esegue l'intero Dag in un unico processo
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...