TaskFlow API

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Klasický přístup

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
Tvorba datových pipeline s Airflow

Přístup TaskFlow

from airflow.sdk import dag, task

@dag def etl_pipeline(): @task def extract_data(): return {"users": 150}
data = extract_data() print_summary(data)
Tvorba datových pipeline s Airflow

Proč TaskFlow?

 

  • Méně zbytečného kódu, dekorátory nahrazují instance operátorů
  • Implicitní závislosti, návratové hodnoty propojují tasky automaticky
  • Čitelnost, DAG vypadá jako běžný Python skript
  • Klasické operátory stále dostupné pro integrace s providery

TaskFlow API – vizuál

Tvorba datových pipeline s Airflow

Airflow UI: Grid view

Airflow Grid view

 

  • Každý sloupec je jedno spuštění DAGu
  • Každý řádek je jeden task
  • Barvy ukazují stav: zelená = úspěch, červená = selhání
  • Klikni na čtvereček a zobraz logy a detaily
Tvorba datových pipeline s Airflow

Airflow UI: Graph view

Airflow Grid view

 

  • Zaměřuje se na jedno spuštění DAGu
  • Zobrazuje závislosti a paralelní běh
  • Detaily lze upravit v nastavení
Tvorba datových pipeline s Airflow

Airflow UI: správa verzí DAGu

Indikátor verze DAGu v Airflow UI

 

  • Automaticky sleduje strukturální změny
  • Každé spuštění je provázáno s verzí aktivní v daném čase
  • Verze najdeš v panelu Dag details
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...