Airflow-kernconcepten

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Maak kennis met je instructeur

  Profielfoto van Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ jaar data engineer in gaming
  • Werkt met Airflow sinds versie 1.x
  • Spreker, mentor en nieuwsbrieflead bij Data Engineer Things
Data-pijplijnen bouwen met Airflow

Wat je gaat bouwen

 

  • Schrijf Dags met de TaskFlow API
  • Bouw dynamische workflows met task mapping en asset-based scheduling
  • Handel fouten af met retries en callbacks
  • Draai SQL-workloads via Airflow

Visualisatie van de inhoud per hoofdstuk

Data-pijplijnen bouwen met Airflow

Voor we beginnen

 

$$

  • Thuis in Dags, taken en operators
  • Bekend met basis van plannen

Introductie tot Airflow - cursuspagina

Data-pijplijnen bouwen met Airflow

Korte opfrisser

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()
  • Een Dag is een verzameling taken met afhankelijkheden
  • Taken zijn losse werkstappen
  • Operators/decorators bepalen wat elke taak doet
  • Afhankelijkheden bepalen de uitvoervolgorde

$$

Eenvoudige Dag

Data-pijplijnen bouwen met Airflow

Airflow-architectuur

$$

Airflow 3-architectuur

 

$$

  • Orchestratie: Scheduler, Dag Processor
  • Uitvoering: Worker, Triggerer
  • Interface & opslag: API Server, Metadata DB
Data-pijplijnen bouwen met Airflow

Planningsaanpakken

# Geen automatische runs: handmatig triggeren (standaard)
@dag(schedule=None)
def my_pipeline(): ...

# Tijdgebaseerd: draait elke dag om 6:00 @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Data-aware: draait wanneer een Asset wordt geüpdatet @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Data-pijplijnen bouwen met Airflow

Twee manieren om Dags te schrijven

Klassieke operators

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Kies dit als er geen decorators beschikbaar zijn

TaskFlow API

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

data = extract()
transform(data)
  • Eenvoudige Python-decorators
  • Minder boilerplate
  • Combineerbaar met klassieke operators
Data-pijplijnen bouwen met Airflow

Oefeningen in deze cursus

$$

IDE-oefening screenshot

 

  • IDE-oefeningen: bewerk echte .py-bestanden
  • Klik "Run this file" of gebruik python3 filename.py
  • dag.test() draait de hele Dag in één proces
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...