
Orchestrazione pipeline: Airflow e Prefect
Orchestrazione di pipeline dati: workflow scheduling, dipendenze e retry management.
Cosa imparerai
- Modellare un DAG Airflow con dipendenze, retry e sensor e confrontarlo con Prefect e Dagster
- Misurare durata, failure rate e costo per flusso per scegliere lo strumento di orchestrazione
Collegamenti
Orchestrazione pipeline: Airflow e Prefect
Questa lezione, sul binario ml-tabellare, ti introduce al mestiere di chi governa i flussi: non basta che uno script faccia il suo lavoro, deve farlo nell’ordine giusto, al momento giusto e per mano della persona giusta.
Che cosa cambia quando arriva l’orchestrazione
L’orchestrazione rende espliciti dipendenze, retry e scheduling dei flussi dati versionati. Tradotto: ciò che prima viveva nella testa del collega che ha scritto lo script, adesso vive nel codice e può essere visto, testato e migliorato.
Il percorso di scelta e setup
- Mappa dipendenze, volumi e necessità di backfill prima di scegliere lo strumento.
- Modella ogni flusso come codice versionato con ambienti separati.
- Configura retry, sensor e branching con parametri dichiarati.
- Misura durata, failure rate e costo per flusso con review settimanale.
Perché Airflow è diventato lo standard
Airflow modella i flussi come grafi diretti aciclici in Python con scheduling, dipendenze e operatori per ogni sistema. Un DAG è un grafo di task con un ordine preciso e senza cicli.
from airflow import DAG
from airflow.operators.bash import BashOperator
from datetime import datetime
with DAG('etl_daily', start_date=datetime(2024,1,1), schedule='@daily') as dag:
extract = BashOperator(task_id='extract', bash_command='python extract.py')
transform = BashOperator(task_id='transform', bash_command='dbt run --select staging')
validate = BashOperator(task_id='validate', bash_command='dbt test')
load = BashOperator(task_id='load', bash_command='dbt run --select marts')
extract >> transform >> validate >> load
I pregi sono maturità, community ampia e operatori per ogni sistema. I limiti sono complessità di gestione e scheduling statico con backfill delicato.
Prefect e Dagster: perché esistono
Prefect e Dagster portano scheduling dinamico con parametri a runtime, retry nativo con backoff esponenziale e caching dei task già eseguiti con gli stessi input.
from prefect import flow, task
@task(retries=3, retry_delay_seconds=60)
def extract():
return fetch_from_api()
@flow
def etl_pipeline():
data = extract()
transform(data)
Le primitive restano funzioni Python con decoratori invece di operatori configurati. La scelta dipende da complessità delle dipendenze, volumi, backfill richiesto e competenze del team.
Verdetto: Airflow per maturità ed ecosistemi ampi e Prefect o Dagster per dinamismo, retry nativo e caching.
I quattro pattern che coprono quasi tutto
Quattro pattern coprono quasi tutti i flussi. Il fan-out con fan-in parallelizza per data e poi riunisce i risultati. Il branching condizionale attiva rami diversi in base alla validità dei dati. Il backfill riprocessa lo storico con catchup o flussi dedicati. I sensor attendono condizioni esterne, come file su storage, prima di procedere.
Esempio SQL: costruire una vista di controllo
Il pattern crea una base analitica con metrica, segmento e finestra temporale per confrontare run, durate e costi tra pipeline diverse.
WITH base_events AS (
SELECT
user_id,
account_id,
event_type,
event_time,
DATE_TRUNC('week', event_time) AS week,
source,
device_type
FROM events
WHERE event_time >= CURRENT_DATE - INTERVAL '180 days'
AND user_id IS NOT NULL
),
weekly_user_metrics AS (
SELECT
week,
user_id,
COALESCE(source, 'unknown') AS source,
COALESCE(device_type, 'unknown') AS device_type,
COUNT(*) AS total_events,
COUNT(DISTINCT DATE(event_time)) AS active_days,
COUNT(DISTINCT event_type) AS event_diversity,
MAX(CASE WHEN event_type IN ('purchase', 'subscribe', 'activation') THEN 1 ELSE 0 END) AS reached_key_outcome
FROM base_events
GROUP BY week, user_id, source, device_type
)
SELECT
week,
source,
device_type,
COUNT(DISTINCT user_id) AS users,
ROUND(AVG(active_days), 2) AS avg_active_days,
ROUND(AVG(event_diversity), 2) AS avg_event_diversity,
ROUND(AVG(reached_key_outcome) * 100, 2) AS key_outcome_rate
FROM weekly_user_metrics
GROUP BY week, source, device_type
ORDER BY week, source, device_type;
Esempio Python: controllare stabilità e anomalie
Il controllo su finestra mobile su failure rate o durata media segnala i flussi che degradano.
# df contiene: week, segment, users, key_outcome_rate
# key_outcome_rate espresso in percentuale, es. 12.4
df = df.sort_values(['segment', 'week']).copy()
df['previous_rate'] = df.groupby('segment')['key_outcome_rate'].shift(1)
df['wow_change_pp'] = df['key_outcome_rate'] - df['previous_rate']
df['rolling_mean'] = df.groupby('segment')['key_outcome_rate'].transform(
lambda s: s.rolling(4, min_periods=2).mean()
)
df['rolling_std'] = df.groupby('segment')['key_outcome_rate'].transform(
lambda s: s.rolling(4, min_periods=2).std()
)
df['z_score'] = (df['key_outcome_rate'] - df['rolling_mean']) / df['rolling_std']
anomalies = df[df['z_score'].abs() >= 2].sort_values('z_score')
print(anomalies[['week', 'segment', 'key_outcome_rate', 'wow_change_pp', 'z_score']])
Dalla sofferenza di Airbnb a un ecosistema
Nel 2014 il team dati di Airbnb crea Airflow per sostituire cron sparsi e script senza dipendenze dichiarate. Il progetto diventa open source nel 2015 e standard per workflow batch con retry, scheduling e backfill in Python. Prefect e Dagster arrivano dal 2020 con scheduling dinamico e caching per superare i limiti dello scheduling statico. La lezione resta una sola: l’orchestratore non rende affidabile una pipeline fragile, ma ne rende visibili le fragilità con owner e tempi misurabili.
Domande per riflettere sul tuo flusso
- Quale dipendenza tra task deve diventare esplicita nel tuo flusso?
- Quale retry e quale sensor proteggono il flusso dai guasti esterni?
- Quale metrica tra durata, failure rate e costo guida la scelta dello strumento?
- Quale backfill riprocessa lo storico senza bloccare il resto?
Bloccato su questo argomento o vuoi applicarlo al tuo caso? Prenota una call di 15 minuti con un analista esperto.
Percorso collegato
Lezioni da leggere insieme
Questi collegamenti portano la lezione dentro il resto del corso: basi da riprendere, passaggi successivi e connessioni tematiche tra moduli.