Go to main content
Orchestrazione pipeline: Airflow e Prefect - immagine ufficiale della lezione su GinnyTech, creata da AD

Pipeline orchestration: Airflow and Prefect

Data pipeline orchestration: workflow scheduling, dependencies, and retry management.

AD
Created byAndrii Dyshkantiuk
Lesson 133 / 236Level: AdvancedDuration: 22 minPrerequisites: 1

What you will learn

  • 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

Pipeline orchestration: Airflow and 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

  1. Mappa dipendenze, volumi e necessità di backfill prima di scegliere lo strumento.
  2. Modella ogni flusso come codice versionato con ambienti separati.
  3. Configura retry, sensor e branching con parametri dichiarati.
  4. 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.

SQL example: building a control view

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;

Python example: checking stability and anomalies

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

  1. Quale dipendenza tra task deve diventare esplicita nel tuo flusso?
  2. Quale retry e quale sensor proteggono il flusso dai guasti esterni?
  3. Quale metrica tra durata, failure rate e costo guida la scelta dello strumento?
  4. Quale backfill riprocessa lo storico senza bloccare il resto?
Serve una mano concreta?

Bloccato su questo argomento o vuoi applicarlo al tuo caso? Prenota una call di 15 minuti con un analista esperto.

Book a call