Concepts de base d'Airflow

Créer des pipelines de données avec Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Faites connaissance avec votre instructeur

  Photo de profil de Volker Janz

$$

Volker Janz

$$

  • Responsable principal de la promotion auprès des développeurs, Astronomer
  • 14+ ans comme ingénieur de données dans le jeu vidéo
  • Travaille avec Airflow depuis la version 1.x
  • Conférencier, mentor et responsable de l'infolettre Data Engineer Things
Créer des pipelines de données avec Airflow

Ce que vous allez construire

 

  • Écrire des Dags avec la TaskFlow API
  • Créer des flux de travaux dynamiques avec mappage de tâches et horaire basé sur les actifs
  • Gérer les échecs avec reprises et rappels
  • Exécuter des charges SQL via Airflow

Visualisation du contenu de chaque chapitre

Créer des pipelines de données avec Airflow

Avant de commencer

 

$$

  • À l'aise avec les Dags, tâches et opérateurs
  • Connaître les bases de l'ordonnancement

Introduction à Airflow - page du cours

Créer des pipelines de données avec Airflow

Petit rappel

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 est un ensemble de tâches avec des dépendances
  • Les tâches sont des unités de travail individuelles
  • Les opérateurs / décorateurs définissent l'action de chaque tâche
  • Les dépendances définissent l'ordre d'exécution

$$

Dag simple

Créer des pipelines de données avec Airflow

Architecture d'Airflow

$$

Architecture d'Airflow 3

 

$$

  • Orchestration : ordonnanceur, processeur de Dag
  • Exécution : travailleur, déclencheur
  • Interface et stockage : serveur d'API, BD des métadonnées
Créer des pipelines de données avec Airflow

Approches d'ordonnancement

# Aucune exécution automatique : déclenchement manuel (par défaut)
@dag(schedule=None)
def my_pipeline(): ...

# À base de temps : s'exécute tous les jours à 6 h @dag(schedule="0 6 * * *") def daily_pipeline(): ...
# Sensible aux données : s'exécute quand un actif est mis à jour @dag(schedule=[Asset("my_asset")]) def downstream_pipeline(): ...
Créer des pipelines de données avec Airflow

Deux façons d'écrire des Dags

Opérateurs classiques

extract = PythonOperator(
    task_id="extract",
    python_callable=extract_fn)
extract >> transform
  • À choisir quand aucun décorateur n'est offert

TaskFlow API

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

data = extract()
transform(data)
  • Décorateurs Python simples
  • Moins de code passe-partout
  • Peut se combiner aux opérateurs classiques
Créer des pipelines de données avec Airflow

Exercices dans ce cours

$$

Capture d'écran d'un exercice IDE

 

  • Exercices IDE : modifier de vrais fichiers .py
  • Cliquer sur "Run this file" ou utiliser python3 filename.py
  • dag.test() exécute le Dag complet dans un seul processus
Créer des pipelines de données avec Airflow

Passons à la pratique !

Créer des pipelines de données avec Airflow

Preparing Video For Download...