TaskFlow API

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Abordarea clasică

def extract_data():
    return {"users": 150, "events": 4200}

with DAG("etl_pipeline") as dag: t1 = PythonOperator(task_id="extract", python_callable=extract_data) t2 = PythonOperator(task_id="summary", python_callable=print_summary)
t1 >> t2
Construirea pipeline-urilor de date cu Airflow

Abordarea TaskFlow

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Construirea pipeline-urilor de date cu Airflow

De ce TaskFlow?

 

  • Mai puțin cod repetitiv, decoratorii înlocuiesc instanțele de operator
  • Dependențe implicite, valorile returnate conectează task-urile automat
  • Lizibil, DAG-ul arată ca un script Python obișnuit
  • Operatorii clasici rămân disponibili pentru integrările cu provideri

TaskFlow API - vizual

Construirea pipeline-urilor de date cu Airflow

Interfața Airflow: vizualizarea Grid

Vizualizare Grid în Airflow

 

  • Fiecare coloană reprezintă o rulare a DAG-ului
  • Fiecare rând reprezintă un task
  • Culorile indică starea: verde = succes, roșu = eșec
  • Apasă pe un pătrat pentru a vedea jurnalele și detaliile
Construirea pipeline-urilor de date cu Airflow

Interfața Airflow: vizualizarea Graph

Vizualizare Grid în Airflow

 

  • Se concentrează pe o singură rulare a DAG-ului
  • Arată dependențele și execuția paralelă
  • Ajustează detaliile din setări
Construirea pipeline-urilor de date cu Airflow

Interfața Airflow: versionarea DAG-urilor

Indicator de versiune DAG în interfața Airflow

 

  • Urmărește automat modificările structurale
  • Fiecare rulare este legată de versiunea activă la acel moment
  • Găsești versiunile în panoul Dag details
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...