Podstawowe pojęcia Airflow

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Twój instruktor

  Volker Janz profile photo

$$

Volker Janz

$$

  • Senior Developer Advocate, Astronomer
  • 14+ lat jako inżynier danych w branży gier
  • Pracuje z Airflow od wersji 1.x
  • Prelegent, mentor i autor newslettera Data Engineer Things
Budowanie potoków danych z Airflow

Co zbudujesz

 

  • Tworzenie Dagów z TaskFlow API
  • Budowanie dynamicznych przepływów z mapowaniem zadań i harmonogramowaniem opartym na zasobach
  • Obsługa błędów za pomocą ponowień i callbacków
  • Uruchamianie obciążeń SQL przez Airflow

Visualization of each chapter content

Budowanie potoków danych z Airflow

Zanim zaczniemy

 

$$

  • Swobodna praca z Dagami, zadaniami i operatorami
  • Znajomość podstaw harmonogramowania

Introduction to Airflow - course page

Budowanie potoków danych z Airflow

Szybkie przypomnienie

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 to zbiór zadań z zależnościami
  • Zadania to pojedyncze jednostki pracy
  • Operatory / dekoratory definiują, co robi każde zadanie
  • Zależności określają kolejność wykonywania

$$

Simple Dag

Budowanie potoków danych z Airflow

Architektura Airflow

$$

Airflow 3 architecture

 

$$

  • Orkiestracja: Scheduler, Dag Processor
  • Wykonanie: Worker, Triggerer
  • Interfejs i przechowywanie: API Server, Metadata DB
Budowanie potoków danych z Airflow

Podejścia do harmonogramowania

# Bez automatycznych uruchomień: wyzwalanie ręczne (domyślnie)
@dag(schedule=None)
def my_pipeline(): ...

# Czasowe: uruchamiane codziennie o 6:00 @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Oparte na danych: uruchamiane po aktualizacji zasobu @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Budowanie potoków danych z Airflow

Dwa sposoby pisania Dagów

Klasyczne operatory

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • Używaj, gdy dekoratory nie są dostępne

TaskFlow API

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

data = extract()
transform(data)
  • Proste dekoratory Pythona
  • Mniej powtarzalnego kodu
  • Można łączyć z klasycznymi operatorami
Budowanie potoków danych z Airflow

Ćwiczenia w tym kursie

$$

IDE exercise screenshot

 

  • Ćwiczenia w IDE: edytujesz prawdziwe pliki .py
  • Kliknij „Run this file" lub użyj python3 filename.py
  • dag.test() uruchamia cały Dag w jednym procesie
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...