Základní koncepty Airflow

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Představení instruktora

  Profilová fotografie Volkera Janze

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ let jako datový inženýr v herním průmyslu
  • Pracuje s Airflow od verze 1.x
  • Přednášející, mentor a vedoucí newsletteru Data Engineer Things
Tvorba datových pipeline s Airflow

Co si postavíš

 

  • Psát Dagy pomocí TaskFlow API
  • Vytvářet dynamické workflow s mapováním tasků a plánováním na základě assetů
  • Zvládat selhání pomocí opakování a callbacků
  • Spouštět SQL úlohy přes Airflow

Vizualizace obsahu jednotlivých kapitol

Tvorba datových pipeline s Airflow

Než začneme

 

$$

  • Znalost Dagů, tasků a operátorů
  • Základní orientace v plánování (scheduling)

Stránka kurzu Introduction to Airflow

Tvorba datových pipeline s Airflow

Rychlé zopakování

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()
  • Dag je kolekce tasků se závislostmi
  • Tasky jsou jednotlivé pracovní jednotky
  • Operátory / dekorátory definují, co každý task dělá
  • Závislosti určují pořadí spouštění

$$

Jednoduchý Dag

Tvorba datových pipeline s Airflow

Architektura Airflow

$$

Architektura Airflow 3

 

$$

  • Orchestrace: Scheduler, Dag Processor
  • Spouštění: Worker, Triggerer
  • Rozhraní a úložiště: API Server, Metadata DB
Tvorba datových pipeline s Airflow

Přístupy k plánování

# No automatic runs: trigger manually (default)
@dag(schedule=None)
def my_pipeline(): ...

# Time-based: runs every day at 6 AM @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Data-aware: runs when an Asset updates @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Tvorba datových pipeline s Airflow

Dva způsoby psaní Dagů

Klasické operátory

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Použij, když dekorátory nejsou k dispozici

TaskFlow API

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

data = extract()
transform(data)
  • Jednoduché Python dekorátory
  • Méně zbytečného kódu
  • Lze kombinovat s klasickými operátory
Tvorba datových pipeline s Airflow

Cvičení v tomto kurzu

$$

Screenshot cvičení v IDE

 

  • Cvičení v IDE: upravuješ skutečné soubory .py
  • Klikni na „Run this file" nebo použij python3 filename.py
  • dag.test() spustí celý Dag v jednom procesu
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...