La Task Execution API d'Airflow 3.0
Airflow 3.0 introduit une Task Execution API qui découple l'exécution des tâches du coeur du scheduler. Un nouveau Task SDK Python remplace l'ancien couplage direct, et cette architecture ouvre la voie à des SDKs dans d'autres langages (Go, Java).
Le Python Task SDK
Le Task SDK Python est un package léger (apache-airflow-task-sdk) qui communique avec le scheduler via une API HTTP. Les tâches s'exécutent dans des processus isolés, ce qui améliore la stabilité et la sécurité.
from airflow.sdk import DAG, task
from datetime import datetime
with DAG(
dag_id='sdk_demo',
schedule='@hourly',
start_date=datetime(2025, 5, 1),
) as dag:
@task
def fetch_data() -> dict:
"""Le SDK gère la sérialisation automatiquement."""
import requests
resp = requests.get('https://api.example.com/data')
return resp.json()
@task
def process(data: dict) -> list:
"""Les tâches communiquent via XCom typé."""
return [item for item in data['results'] if item['active']]
@task
def store(records: list):
print(f"Stockage de {len(records)} enregistrements")
# Chaînage fluide avec le TaskFlow API
store(process(fetch_data()))
# Le SDK communique avec le scheduler via HTTP,
# pas d'import du coeur Airflow dans le worker.
Vers des SDKs multi-langages
Grâce à l'API HTTP standardisée, des SDKs dans d'autres langages pourront définir et exécuter des tâches Airflow. Un SDK Go est en cours de développement. Cela permet d'intégrer des équipes non-Python dans les pipelines Airflow.
// Aperçu du futur SDK Go (en développement)
package main
import (
"fmt"
airflow "github.com/apache/airflow-go-sdk"
)
func ProcessData(ctx airflow.TaskContext) error {
data := ctx.GetXCom("fetch_data")
fmt.Printf("Processing %d records\n", len(data))
// Traitement en Go pour les performances
ctx.PushXCom("result", processed)
return nil
}
// Les tâches Go communiquent avec le scheduler
// via la même API HTTP que le SDK Python.
