TaskFlow API

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Det klassiska tillvägagångssättet

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
Bygg datapipelines med Airflow

TaskFlow-metoden

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Bygg datapipelines med Airflow

Varför TaskFlow?

 

  • Mindre standardkod – dekoratorer ersätter operatorinstanser
  • Implicita beroenden – returvärden kopplar ihop uppgifter automatiskt
  • Läsbar – DAG:en läses som ett vanligt Python-skript
  • Klassiska operatorer finns kvar för providerintegration

TaskFlow API - visual

Bygg datapipelines med Airflow

Airflow UI: Grid-vy

Airflow Grid-vy

 

  • Varje kolumn är en DAG-körning
  • Varje rad är en uppgift
  • Färger visar status: grön = lyckad, röd = misslyckad
  • Klicka på en ruta för att se loggar och detaljer
Bygg datapipelines med Airflow

Airflow UI: Grafvy

Airflow Grid-vy

 

  • Fokuserar på en DAG-körning
  • Visar beroenden och parallell körning
  • Justera detaljer i inställningarna
Bygg datapipelines med Airflow

Airflow UI: DAG-versionshantering

Versionsindikator för DAG i Airflow UI

 

  • Spårar strukturella ändringar automatiskt
  • Varje körning kopplas till den version som var aktiv vid tillfället
  • Hitta versioner i panelen DAG-detaljer
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...