Concepte de bază în Airflow

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Cunoaște-ți instructorul

  Fotografia de profil a lui Volker Janz

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ ani ca inginer de date în gaming
  • Lucrează cu Airflow din versiunea 1.x
  • Speaker, mentor și responsabil newsletter la Data Engineer Things
Construirea pipeline-urilor de date cu Airflow

Ce vei construi

 

  • Creează Dag-uri cu TaskFlow API
  • Construiește fluxuri dinamice cu task mapping și planificare bazată pe resurse
  • Gestionează eșecurile cu reîncercări și callback-uri
  • Rulează sarcini SQL prin Airflow

Vizualizarea conținutului fiecărui capitol

Construirea pipeline-urilor de date cu Airflow

Înainte să începem

 

$$

  • Familiarizat cu Dag-uri, task-uri și operatori
  • Cunoști noțiunile de bază ale planificării

Pagina cursului Introduction to Airflow

Construirea pipeline-urilor de date cu Airflow

Recapitulare rapidă

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()
  • Un Dag este o colecție de task-uri cu dependențe
  • Task-urile sunt unități individuale de lucru
  • Operatorii / decoratorii definesc ce face fiecare task
  • Dependențele stabilesc ordinea de execuție

$$

Dag simplu

Construirea pipeline-urilor de date cu Airflow

Arhitectura Airflow

$$

Arhitectura Airflow 3

 

$$

  • Orchestrare: Scheduler, Dag Processor
  • Execuție: Worker, Triggerer
  • Interfață și stocare: API Server, Metadata DB
Construirea pipeline-urilor de date cu Airflow

Metode de planificare

# 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(): ...
Construirea pipeline-urilor de date cu Airflow

Două moduri de a scrie Dag-uri

Operatori clasici

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Folosește când decoratorii nu sunt disponibili

TaskFlow API

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

data = extract()
transform(data)
  • Decoratori Python simpli
  • Mai puțin cod boilerplate
  • Compatibil cu operatorii clasici
Construirea pipeline-urilor de date cu Airflow

Exercițiile din acest curs

$$

Captură de ecran exercițiu IDE

 

  • Exerciții IDE: editezi fișiere .py reale
  • Apasă "Run this file" sau folosește python3 filename.py
  • dag.test() rulează întregul Dag într-un singur proces
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...