Standardmäßig ist der Prompt ausgewählt, der zuerst die Quelle prüft. Sie können zu einem direkten Befehl wechseln oder eine lokale Kopie herunterladen.
Quelldateien prüfen
Lesen Sie SKILL.md und alle von SkillsMP angezeigten Begleitdateien, bevor Sie sich für eine Installation entscheiden.
Mit Codex oder Claude installieren Kopieren Sie diesen Prompt, fügen Sie ihn in Codex, Claude oder einen anderen Assistant ein und lassen Sie die Skill-Seite prüfen und installieren.
Ein direkter Befehl überspringt den Prüf-Prompt. Prüfen Sie die Quelle, bevor Sie ihn ausführen.
Compétence Apache Airflow : Orchestration de Pipelines de Données
Vue d'ensemble
Apache Airflow est la plateforme d'orchestration la plus utilisée en data engineering. Elle permet de définir, planifier et monitorer des pipelines de données complexes sous forme de DAGs (Directed Acyclic Graphs).
Architecture clé :
Scheduler : planifie les exécutions des DAGs.
Worker (Executor) : exécute les tâches (LocalExecutor, CeleryExecutor, KubernetesExecutor).
Web Server : interface utilisateur pour le monitoring.
Metastore : base de données (PostgreSQL) pour l'état des DAGs/tâches.
Queue (Celery) : file de messages pour la distribution des tâches.
Cette compétence couvre : TaskFlow API (2.0+), operators custom, sensors, Dynamic Task Mapping, execution environments (Docker / K8s Pod operator), gestion des secrets (Airflow Connections / Variables), alerting et SLA.
Quand l'utiliser
Activez cette compétence lorsque l'utilisateur :
Veut orchestrer des pipelines ETL/ELT complexes multi-systèmes.
A besoin de planification, retry, alerting pour des jobs batch.
Demande d'interagir avec Spark, Kafka, dbt, Snowflake via Airflow.
Veut gérer des dépendances inter-DAG ou intra-DAG.
Pose des questions sur le choix de l'executor (Celery vs Kubernetes), les pools, les SLAs.
export AIRFLOW_HOME=~/airflow
airflow db init
airflow users create --username eva --password eva --role Admin \
--email eva@thehive.local --firstname EVA --lastname Admin
airflow standalone # scheduler + webserver en un
# docker-compose.yml pour Airflow Celeryversion:'3.8'services:airflow-scheduler:image:apache/airflow:2.10command:schedulerenvironment:AIRFLOW__CORE__EXECUTOR:CeleryExecutorAIRFLOW__CELERY__BROKER_URL:redis://redis:6379/0AIRFLOW__CELERY__RESULT_BACKEND:db+postgresql://airflow:airflow@postgres:5432/airflowAIRFLOW__CELERY__WORKER_CONCURRENCY:8airflow-worker:image:apache/airflow:2.10command:celeryworkerdeploy:replicas:3environment:AIRFLOW__CORE__EXECUTOR:CeleryExecutorredis:image:redis:7-alpinepostgres:image:postgres:15environment:POSTGRES_DB:airflow
5. Gestion des Secrets et Variables
5.1 Airflow Connections
from airflow.models import Connection
from airflow.settings import Session
# Créer une connexion programmatique
conn = Connection(
conn_id='postgres_dw',
conn_type='postgres',
host='postgres-1.data.internal',
login='airflow',
password='${AIRFLOW_SECRET_POSTGRES_PW}',
port=5432,
extra={
"sslmode": "require",
"keepalives_idle": 30,
}
)
# Utiliser en Variable d'environnement# export AIRFLOW_CONN_POSTGRES_DW='postgres://airflow:pass@host:5432/dw'
5.2 Variables et Paramètres de DAG
from airflow.models import Variable
# Définir
Variable.set("slack_webhook", "https://hooks.slack.com/services/...")
Variable.set("data_config", {
"batch_size": 10000,
"retention_days": 90,
"output_format": "parquet"
})
# Utiliser dans un DAGclassConfigurableDAG:
@taskdefuse_config():
config = Variable.get("data_config", deserialize_json=True)
print(f"Batch size: {config['batch_size']}")
5.3 Paramètres (Airflow 2.4+)
from airflow.models.param import Param
@dag(
params={
'date_range': Param(
default='2026-07-01,2026-07-22',
type='string',
description='Plage de dates (YYYY-MM-DD,YYYY-MM-DD)'),
'full_refresh': Param(
default=False,
type='boolean',
),
}
)defconfigurable_etl():
@taskdefextract(params):
dates = params['date_range'].split(',')
print(f"Extraction du {dates[0]} au {dates[1]}")
# Limiter le parallélisme : 2 tâches Spark simultanées max
run_spark_1 = SparkSubmitOperator(
task_id='spark_load_1',
pool='spark_pool',
...
)
run_spark_2 = SparkSubmitOperator(
task_id='spark_load_2',
pool='spark_pool',
...
)
# Config du pool : airflow pools set spark_pool 2 "Pool pour les jobs Spark"
Pièges Courants (Pitfalls)
DAGs non idempotents.
Erreur : Un backfill ou un retry produit des doublons dans la base cible.
Correction : Tous les DAGs doivent être idempotents. Utiliser INSERT OVERWRITE, MERGE, ou des clés déterministes.
Dépendances extérieures sans gestion de timeout.
Erreur : Un Sensor attend indéfiniment un fichier qui n'arrive jamais → occupe un worker slot.
Correction : Toujours configurer timeout et mode='reschedule' pour les sensors longue durée.
Surcharge de la base de données metastore.
Erreur : Trop de DAGs actifs (500+) ou des tâches avec des logs énormes → PostgreSQL ralentit.
Correction : Purger régulièrement : airflow db clean --before-date 2026-01-01. Optimiser les connexions avec PGBouncer.
xcoms trop volumineux.
Erreur : Passer un DataFrame (plusieurs Mo) via xcom → sature la base et ralentit le scheduler.
Correction : Stocker les données volumineuses dans S3/S3 et ne passer que le chemin via xcom.
CeleryExecutor — perte de tâches sur restart worker.
Erreur : Un worker Celery redémarre pendant une tâche longue → la tâche est marquée en échec mais l'effet secondaire persiste.
Correction : Configurer task_acks_late=True (ack après complétion) avec task_reject_on_worker_lost=True.
Liste de Vérification (Checklist)
DAGs idempotents (retry et backfill sans doublons).
catchup=False sauf pour backfill explicite.
Sensors avec timeout et mode='reschedule'.
SLAs configurés pour les pipelines critiques.
Pools définis pour limiter la concurrence.
xcom limités aux métadonnées, pas aux données.
Notifications sur échec (Slack/Teams/Email).
retries configuré (au moins 2 pour la production).