Airflows kärnkoncept

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Möt din instruktör

  Volker Janz profilbild

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ års erfarenhet som datatekniker inom spelindustrin
  • Arbetat med Airflow sedan version 1.x
  • Talare, mentor och nyhetsbrevansvarig hos Data Engineer Things
Bygg datapipelines med Airflow

Det här bygger du

 

  • Skriv Dags med TaskFlow API
  • Bygg dynamiska arbetsflöden med task-mappning och tillgångsbaserad schemaläggning
  • Hantera fel med återförsök och callbacks
  • Kör SQL-arbetsbelastningar via Airflow

Visualisering av kursens kapitelinnehåll

Bygg datapipelines med Airflow

Innan vi börjar

 

$$

  • Bekväm med Dags, tasks och operatorer
  • Bekant med grundläggande schemaläggning

Introduktion till Airflow – kurssida

Bygg datapipelines med Airflow

Snabb repetition

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()
  • En Dag är en samling tasks med beroenden
  • Tasks är enskilda arbetsenheter
  • Operatorer / dekoratorer definierar vad varje task gör
  • Beroenden bestämmer körningsordningen

$$

Enkel Dag

Bygg datapipelines med Airflow

Airflow-arkitektur

$$

Airflow 3-arkitektur

 

$$

  • Orkestrering: Scheduler, Dag Processor
  • Körning: Worker, Triggerer
  • Gränssnitt & lagring: API Server, Metadata DB
Bygg datapipelines med Airflow

Schemaläggningsmetoder

# 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(): ...
Bygg datapipelines med Airflow

Två sätt att skriva Dags

Klassiska operatorer

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Välj när dekoratorer saknas

TaskFlow API

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

data = extract()
transform(data)
  • Enkla Python-dekoratorer
  • Mindre boilerplate-kod
  • Kan kombineras med klassiska operatorer
Bygg datapipelines med Airflow

Övningar i den här kursen

$$

Skärmdump av IDE-övning

 

  • IDE-övningar: redigera riktiga .py-filer
  • Klicka på "Run this file" eller använd python3 filename.py
  • dag.test() kör hela Dag:en i en och samma process
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...