Corso

Immagine dell’autore.
Apache Airflow è un popolare strumento open source di orchestrazione dei dati progettato per creare, programmare e monitorare pipeline di dati. Offre una dashboard che aiuta a gestire lo stato dei workflow, rendendolo uno strumento perfetto per la maggior parte delle esigenze.
Tuttavia, ad Airflow mancano alcune funzionalità importanti che possono essere vitali per requisiti complessi e moderni di orchestrazione dei dati.
In questo tutorial, esploreremo cinque alternative ad Airflow che offrono capacità avanzate e colmano alcune delle sue lacune. Inoltre, impareremo a creare una semplice pipeline ETL con ciascuno strumento, eseguirla e visualizzarla nella loro dashboard.
Perché scegliere un’alternativa ad Airflow?
Airflow è potente per vari workflow di dati, ma presenta alcune limitazioni che potrebbero spingere le aziende a considerare alternative.
Ecco alcuni motivi per cui potresti scegliere un’alternativa:
- Ripida curva di apprendimento: Airflow può essere impegnativo da imparare, soprattutto per chi è alle prime armi con gli strumenti di gestione dei workflow.
- Manutenzione: richiede una manutenzione significativa, specialmente in deployment su larga scala.
- Documentazione insufficiente: diversi utenti segnalano problemi di documentazione che rendono più difficile risolvere problemi o apprendere nuove funzionalità.
- Intensivo in termini di risorse: Airflow può richiedere molte risorse, con notevole uso di CPU e memoria per funzionare in modo efficiente.
- Flessibilità limitata per chi non usa Python: la filosofia “workflow-as-code” si basa fortemente su Python, cosa che può escludere esperti di dominio non avvezzi alla programmazione.
- Scalabilità: alcuni utenti riportano difficoltà a scalare Airflow per workflow di grandi dimensioni.
- Elaborazione real-time limitata: Airflow è pensato principalmente per l’elaborazione batch, non per flussi di dati in tempo reale.
Prima di passare alla parte di codice degli altri strumenti di orchestrazione, è importante imparare a scrivere una pipeline di dati con Apache Airflow seguendo il tutorial Guida introduttiva ad Apache Airflow, così da confrontare correttamente le alternative.
Se sei completamente nuovo ad Airflow, valuta di seguire il breve corso Introduzione ad Airflow in Python per apprendere le basi della creazione e programmazione di pipeline di dati.
Le 5 migliori alternative ad Airflow per l’orchestrazione dei dati
Descriviamo ora le 5 principali alternative ad Airflow e vediamo come usarle con esempi pratici di codice.
1. Prefect
Prefect è uno strumento open source di orchestrazione di workflow in Python pensato per data e ML engineer moderni. Offre una semplice API per creare rapidamente una pipeline di dati e gestirla tramite una dashboard interattiva.
Perfect offre un modello di esecuzione ibrido: puoi eseguire il workflow nel cloud oppure in locale dal repository.
Rispetto ad Airflow, Prefect include funzionalità avanzate come dipendenze dei task automatizzate, trigger basati su eventi, notifiche integrate, infrastrutture specifiche per workflow e condivisione dei dati tra task. Queste capacità lo rendono una soluzione potente per gestire workflow complessi in modo efficiente ed efficace.
Prefect è semplice e ricco di funzionalità potenti. Mi sono bastati 5 minuti per eseguire il codice di esempio. Mi piace in particolare il design della dashboard, come si impostano le notifiche, il riavvio delle pipeline e la gestione e il monitoraggio di tutto dalla Dashboard.
Abid Ali Awan, Author
Leggi il blog Airflow vs Prefect: quale scegliere per il tuo workflow dati per un confronto dettagliato tra questi due strumenti di orchestrazione dei dati.
Per iniziare con Prefect
Inizieremo il progetto Prefect installando il pacchetto Python. Esegui il seguente comando in un terminale.
$ pip install -U prefect
Dopodiché, crea uno script Python chiamato prefect_etl.py e inserisci il seguente codice.
from prefect import task, flow
import pandas as pd
# Extract data
@task
def extract_data():
# Simulating data extraction
data = {
"name": ["Alice", "Bob", "Charlie"],
"age": [25, 30, 35],
"city": ["New York", "Los Angeles", "Chicago"]
}
df = pd.DataFrame(data)
return df
# Transform data
@task
def transform_data(df: pd.DataFrame):
# Example transformation: adding a new column
df["age_plus_ten"] = df["age"] + 10
return df
# Load data
@task
def load_data(df: pd.DataFrame):
# Simulating data load
print("Loading data to target destination:")
print(df)
# Defining the flow
@flow(log_prints=True)
def etl():
raw_data = extract_data()
transformed_data = transform_data(raw_data)
load_data(transformed_data)
# Running the flow
if __name__ == "__main__":
etl()
Il codice sopra definisce le funzioni task extract_data(), transform_data(), e load_data() ed esegue la serie in un flow chiamato etl(). Queste funzioni sono create con i decorator Python di Prefect.
In breve, creiamo un DataFrame pandas, lo trasformiamo e poi mostriamo il risultato finale con una print. È un modo semplice per simulare una pipeline ETL.
Per eseguire il workflow, basta lanciare lo script Python con il seguente comando.
$ python prefect_etl.py
Come si vede, l’esecuzione del workflow è andata a buon fine.

Log di esecuzione del flow in Prefect.
Deployment del flow
Ora effettueremo il deployment del workflow per eseguirlo a orario o attivarlo in base a un evento. Il deployment consente anche di monitorare e gestire più workflow in modo centralizzato.
Per il deployment useremo la CLI di Prefect. La funzione deploy richiede il nome del file Python, il nome della funzione del flow nel file e il nome del deployment. In questo caso lo chiamiamo “simple_etl”.
$ prefect deploy prefect_etl.py:etl -n 'simple_etl'
Dopo aver eseguito lo script potresti ricevere il messaggio che non hai un worker pool per eseguire il deployment. Per crearlo, usa il seguente comando.
$ prefect worker start --pool 'datacamp'
Ora che abbiamo un worker pool, apri un altro terminale ed esegui il deployment. Il comando prefect deployment run richiede come argomento “<flow-function-name>/<deployment-name>”, come mostrato sotto.
$ prefect deployment run 'etl/simple_etl
In seguito all’esecuzione del deployment, riceverai il messaggio che il workflow è in esecuzione. Di solito al flow run creato viene assegnato un nome casuale, nel mio caso witty-lorikeet.
Creating flow run for deployment 'etl/simple_etl'...
Created flow run 'witty-lorikeet'.
└── UUID: 4e0495b0-9c7e-4ed8-b9ab-5160994dc7f0
└── Parameters: {}
└── Job Variables: {}
└── Scheduled start time: 2024-06-22 14:05:01 PKT (now)
└── URL: <no dashboard available>
Per vedere l’intero log, torna alla finestra del terminale in cui hai avviato il worker pool.

Riepilogo esecuzione flow in Prefect.
Per visualizzare l’esecuzione del flow in modo più intuitivo e gestire altri workflow, devi avviare il web server di Prefect.
$ prefect server start
Dopo aver eseguito il comando, dovresti essere reindirizzato alla dashboard di Prefect. In alternativa, vai direttamente all’indirizzo http://127.0.0.1:4200 nel browser.

Interfaccia web server Prefect
La dashboard ti permette di rieseguire il workflow, visualizzare i log, controllare i work pool, impostare notifiche e scegliere altre opzioni avanzate. È una soluzione completa per le moderne esigenze di orchestrazione dei dati.
Per imparare a creare ed eseguire pipeline di machine learning con Prefect, segui il tutorial Usare Prefect per i workflow di machine learning.
2. Dagster
Dasgter è un framework open source progettato per i data engineer per definire, programmare e monitorare pipeline di dati. È altamente scalabile e facilita la collaborazione tra vari team di dati.
Dagster consente di definire gli asset di dati come funzioni Python usando i decorator. Una volta definiti, è possibile eseguirli facilmente tramite schedulazione o trigger basati su eventi.
Rispetto ad Airflow, Dagster permette di sviluppare, testare e revisionare la pipeline in locale, offre un approccio all’orchestrazione basato su asset ed è cloud- e container-native.
Invece di pensare al workflow in termini di step e flow, ho dovuto cambiare prospettiva e costruire una pipeline usando asset di dati. A parte questo, creare ed eseguire una semplice pipeline ETL è stato piuttosto semplice. Inoltre, il web server è relativamente essenziale ma fornisce tutte le informazioni per monitorare asset, esecuzioni e deployment.
Abid Ali Awan, Author
Per iniziare con Dagster
Creeremo una semplice pipeline ETL, la eseguiremo e la visualizzeremo tramite il web server di Dagster. Come la dashboard di Prefect, anche il web server di Dagster offre modalità centralizzate per monitorare più workflow e pianificare esecuzioni e asset.
Iniziamo installando il pacchetto Python.
$ pip install dagster -q
Poi creeremo tre funzioni Python per estrarre, trasformare e caricare i dati. Queste funzioni si chiamano create_dirty_data(), clean_data(), e load_cleaned_data() nel codice. Con il decorator @asset dichiareremo le funzioni come asset di dati in Dagster.
Successivamente, creeremo il job degli asset (variabile job) usando tutti gli asset (variabile all_assets) e poi creeremo la definizione degli asset (variabile defs).
Puoi saltare la parte di definizione degli asset, ma diventa importante se vuoi schedulare l’esecuzione, eseguire più job e impostare sensori.
import pandas as pd
import numpy as np
from dagster import asset, Definitions, define_asset_job, materialize
@asset
def create_dirty_data():
# Create a sample DataFrame with dirty data
data = {
'Name': [' John Doe ', 'Jane Smith', 'Bob Johnson ', ' Alice Brown'],
'Age': [30, np.nan, 40, 35],
'City': ['New York', 'los angeles', 'CHICAGO', 'Houston'],
'Salary': ['50,000', '60000', '75,000', 'invalid']
}
df = pd.DataFrame(data)
# Save the DataFrame to a CSV file
dirty_file_path = 'dag_data/dirty_data.csv'
df.to_csv(dirty_file_path, index=False)
return dirty_file_path
@asset
def clean_data(create_dirty_data):
# Read the dirty CSV file
df = pd.read_csv(create_dirty_data)
# Clean the data
df['Name'] = df['Name'].str.strip()
df['Age'] = pd.to_numeric(df['Age'], errors='coerce').fillna(df['Age'].mean())
df['City'] = df['City'].str.upper()
df['Salary'] = df['Salary'].replace('[\$,]', '', regex=True)
df['Salary'] = pd.to_numeric(df['Salary'], errors='coerce').fillna(0)
# Calculate average salary
avg_salary = df['Salary'].mean()
# Save the cleaned DataFrame to a new CSV file
cleaned_file_path = 'dag_data/cleaned_data.csv'
df.to_csv(cleaned_file_path, index=False)
return {
'cleaned_file_path': cleaned_file_path,
'avg_salary': avg_salary
}
@asset
def load_cleaned_data(clean_data):
cleaned_file_path = clean_data['cleaned_file_path']
avg_salary = clean_data['avg_salary']
# Read the cleaned CSV file to verify
df = pd.read_csv(cleaned_file_path)
print({
'num_rows': len(df),
'num_columns': len(df.columns),
'avg_salary': avg_salary
})
# Define all assets
all_assets = [create_dirty_data, clean_data, load_cleaned_data]
# Create a job that will materialize all assets
job = define_asset_job("all_assets_job", selection=all_assets)
# Create Definitions object
defs = Definitions(
assets=all_assets,
jobs=[job]
)
if __name__ == "__main__":
result = materialize(all_assets)
print("Pipeline execution result:", result.success)
Puoi eseguire il codice in un Jupyter Notebook o creare un file Python ed eseguirlo.
Come risultato dell’esecuzione, otterremo un log completo del workflow.

Web server di Dagster
Per visualizzare gli asset e le esecuzioni dei job, dobbiamo installare ed eseguire il web server di Dagster. Il web server consente di avviare i job, materializzare singoli asset e monitorare più job contemporaneamente.
$ pip install dagster-webserver
Per avviare il server Dagster, useremo la CLI di Daster fornendole il percorso del file Python. In questo caso ho chiamato il file dagster_pipe.py.
$ dagster dev -f dagster_pipe.py
Il comando avvierà automaticamente il web server nel browser. In alternativa, vai direttamente all’indirizzo http://127.0.0.1:3000 nel browser.

Interfaccia web server Dagster.
Finora abbiamo effettuato solo il deployment del job. Per eseguire il workflow, vai alla scheda “Runs” e clicca su “Launch a new run”.
L’esecuzione dovrebbe concludersi con successo! Per vedere i log, clicca sull’ID dell’esecuzione di tuo interesse.

Log delle esecuzioni in Dagster.
3. Mage AI
Mage AI è un framework open source ibrido per l’orchestrazione dei dati. Ibrido significa che offre la flessibilità di un Jupyter Notebook e il controllo di codice modulare.
Chiunque, anche con conoscenze limitate di Python, può creare, eseguire e monitorare pipeline di dati. Invece di scrivere ed eseguire direttamente un file Python, creerai un progetto Mage AI e lo avvierai dalla dashboard, dove potrai costruire, eseguire e gestire le tue pipeline.
Rispetto ad Airflow, Mage AI offre un’interfaccia intuitiva e facilità d’uso, risultando un’ottima scelta per chi è nuovo al data engineering. È progettato per scalare e gestire grandi volumi di dati e strutture di pipeline complesse in modo efficiente.
Mi sono sentito spaesato perché era completamente diverso da ciò a cui sono abituato. Ho dovuto installare e avviare l’interfaccia web di Mage AI. Doveva essere semplice, ma ho trovato difficile costruire ed eseguire la pipeline ETL. D’altra parte, capisco perché questo design unico possa attrarre chi è alle prime armi: è sostanzialmente drag and drop e click su pulsanti.
Abid Ali Awan, Author
Per iniziare con Mage AI
Avviare Mage AI è piuttosto semplice. Dobbiamo solo installare il pacchetto Python di Mage AI.
$ pip install mage-ai
E avviare il progetto Mage AI.
$ mage start mage_ai_etl
Il comando inizializzerà il web server. Come detto, tutta la modifica del codice, l’esecuzione e il monitoraggio dei job avvengono tramite l’interfaccia di Mage AI.

Interfaccia Mage AI.
Clicca su “+ New pipeline” per creare la tua prima pipeline ETL. Io l’ho chiamata “simple_etl”.

Creazione della nuova pipeline in Mage AI.
L’interfaccia ti chiederà poi di aggiungere un modulo per iniziare a scrivere codice. Seleziona il modulo “Data Loader” e inserisci il seguente codice Python.
Qui dichiariamo una funzione create_sample_csv(), che è il primo step della pipeline. Usiamo il decorator @data_loader di Mage AI. Definiamo anche una funzione test_output() che verifica l’esistenza dell’output. Questo aiuta nella gestione delle dipendenze tra task.
import io
import pandas as pd
if 'data_loader' not in globals():
from mage_ai.data_preparation.decorators import data_loader
if 'test' not in globals():
from mage_ai.data_preparation.decorators import test
@data_loader
def create_sample_csv() -> pd.DataFrame:
"""
Create a sample CSV file with duplicates and missing values
"""
csv_data = """
category,product,quantity,price
Electronics,Laptop,5,1000
Electronics,Smartphone,10,500
Clothing,T-shirt,50,20
Clothing,Jeans,30,50
Books,Novel,100,15
Books,Textbook,20,80
Electronics,Laptop,5,1000
Clothing,T-shirt,,20
Electronics,Tablet,,300
Books,Magazine,25,
"""
return pd.read_csv(io.StringIO(csv_data.strip()))
@test
def test_output(df) -> None:
"""
Template code for testing the output of the block.
"""
assert df is not None, 'The output is undefined'

Creare il blocco data loader in Mage AI.
Crea quindi un altro modulo chiamato “Transformer” e aggiungi la funzione clean_data(), come nel codice qui sotto.
Puoi ignorare la funzione test(); devi solo aggiungere la funzione di trasformazione principale, clean_data().
import pandas as pd
if 'transformer' not in globals():
from mage_ai.data_preparation.decorators import transformer
if 'test' not in globals():
from mage_ai.data_preparation.decorators import test
@transformer
def clean_data(df: pd.DataFrame) -> pd.DataFrame:
"""
Clean and transform the data
"""
# Remove duplicates
df = df.drop_duplicates()
# Fill missing values with 0
df = df.fillna(0)
return df
@test
def test_output(df) -> None:
"""
Template code for testing the output of the block.
"""
assert df is not None, 'The output is undefined'
Allo stesso modo, crea un modulo “Data Exporter” e aggiungi il seguente codice. Dichiara una funzione di caricamento dati, export_data_to_csv(), che salva i dati trasformati in un file CSV.
import pandas as pd
if 'data_exporter' not in globals():
from mage_ai.data_preparation.decorators import data_exporter
@data_exporter
def export_data_to_csv(df: pd.DataFrame) -> None:
"""
Export the processed data to a CSV file
"""
df.to_csv('output_data.csv', index=False)
print("Data exported successfully to output_data.csv")
Per eseguire la pipeline, vai alla scheda “Trigger” e clicca su “Run@once”.

Esecuzione della pipeline in Mage AI.
Per visualizzare i log, vai alla scheda “Runs” e clicca sul pulsante “Logs” della pipeline appena eseguita.

Log del flow in Mage AI.
4. Kedro
Kedro è un altro popolare framework open source per l’orchestrazione dei dati, con un approccio leggermente diverso rispetto agli altri strumenti. È stato creato per i machine learning engineer e prende in prestito molti concetti dall’ingegneria del software, applicandoli ai progetti di ML.
Kedro è progettato per essere altamente modulare, il che significa che persino per esportare un dataset devi creare un data catalog che specifichi posizione e tipo dei dati, garantendo una gestione standardizzata ed efficiente lungo l’intera pipeline.
Per capire come Kedro si inserisce nell’ecosistema del machine learning, puoi esplorare vari strumenti MLOps leggendo l’articolo 25 Top MLOps Tools You Need to Know in 2024.
Rispetto ad Airflow, l’API di Kedro è più semplice per costruire una pipeline di dati. È più focalizzato sull’ingegneria del machine learning e offre categorizzazione e versionamento dei dati.
La parte di codice è abbastanza lineare, ma i problemi sorgono quando vuoi eseguire la pipeline. Devi creare un data catalog, registrare la pipeline e capire la struttura del progetto Kedro. Direi che è più impegnativo rispetto a Dagster e Prefect. Tuttavia, capisco perché sia progettato così: per rendere la pipeline affidabile e priva di errori.
Abid Ali Awan, Author
Per iniziare con Kedro
Creare una pipeline dati con Kedro è un gioco diverso. Il framework è modulare e devi comprendere la struttura del progetto e i vari passaggi per eseguire correttamente il workflow.
Inizia installando il pacchetto Python di Kedro.
$ pip install kedro
Inizializza il progetto Kedro.
$ kedro new --name=kedro_etl --tools=none --example=n
Spostati nella directory del progetto.
$ cd kedro-etl
Crea una cartella all’interno della directory pipelines chiamata data_processing.
$ mkdir -p src/kedro_etl/pipelines/data_processing
Crea un file Python chiamato kedro_pipe.py e aprilo nel tuo IDE preferito, ad esempio Visual Studio Code.
$ code src/kedro_etl/pipelines/data_processing/kedro_pipe.py
Lo script Python deve contenere le funzioni di estrazione, trasformazione e caricamento, che sono i nodi della pipeline. In questo caso sono le funzioni create_sample_data(), clean_data(), e load_and_process_data().
Poi uniamo questi nodi usando la classe Pipeline di Kedro all’interno della funzione create_pipeline(). Nella funzione della pipeline definiamo i nodi, e ogni nodo ha inputs, outputs e un name.
import pandas as pd
import numpy as np
from kedro.pipeline import Pipeline, node
def create_sample_data():
data = {
'id': range(1, 101),
'name': [f'Person_{i}' for i in range(1, 101)],
'age': np.random.randint(18, 80, 100),
'salary': np.random.randint(20000, 100000, 100),
'missing_values': [np.nan if i % 10 == 0 else i for i in range(100)]
}
return pd.DataFrame(data)
def clean_data(df: pd.DataFrame):
# Remove rows with missing values
df_cleaned = df.dropna()
# Convert salary to thousands
df_cleaned['salary'] = df_cleaned['salary'] / 1000
# Capitalize names
df_cleaned['name'] = df_cleaned['name'].str.upper()
return df_cleaned
def load_and_process_data(df: pd.DataFrame):
# Calculate average salary
avg_salary = df['salary'].mean()
# Add a new column for salary category
df['salary_category'] = df['salary'].apply(
lambda x: 'High' if x > avg_salary else 'Low')
# Calculate age groups
df['age_group'] = pd.cut(df['age'], bins=[0, 30, 50, 100], labels=[
'Young', 'Middle', 'Senior'])
print(df)
return df
def create_pipeline(**kwargs):
return Pipeline(
[
node(
func=create_sample_data,
inputs=None,
outputs="raw_data",
name="create_sample_data_node",
),
node(
func=clean_data,
inputs="raw_data",
outputs="cleaned_data",
name="clean_data_node",
),
node(
func=load_and_process_data,
inputs="cleaned_data",
outputs="processed_data",
name="load_and_process_data_node",
),
]
)
Se eseguiamo la pipeline senza creare il data catalog, i dati non verranno esportati. Quindi, apri il file conf/base/catalog.yml e modificalo fornendo la configurazione dei dataset.
raw_data:
type: pandas.CSVDataset
filepath: ./data/kedro/sample_data.csv
cleaned_data:
type: pandas.CSVDataset
filepath: ./data/kedro/cleaned_data.csv
processed_data:
type: pandas.CSVDataset
filepath: ./data/kedro/processed_data.csv
Dobbiamo anche includere il nuovo file Python nel registro delle pipeline. Per farlo, vai al file Python src/simple_etl/pipeline_registry.py e inserisci il seguente codice.
"""Project pipelines."""
from __future__ import annotations
from kedro.pipeline import Pipeline
from kedro_etl.pipelines.data_processing import kedro_pipe
def register_pipelines() -> Dict[str, Pipeline]:
data_processing_pipeline = kedro_pipe.create_pipeline()
return {
"__default__": data_processing_pipeline,
"data_processing": data_processing_pipeline,
}
Esegui la pipeline e visualizza i log in tempo reale nel terminale con il seguente comando.
$ kedro run

Log dell’esecuzione della pipeline Kedro.
Dopo l’esecuzione, i file verranno salvati in formato CSV nella posizione definita nel data catalog.

File di output dell’esecuzione della pipeline Kedro.
Se riscontri problemi nell’esecuzione della pipeline, valuta di installare Kedro con tutte le estensioni.
$ pip install "kedro[all]"
Visualizzazione con Kedro
Possiamo visualizzare e condividere le nostre pipeline installando lo strumento kedro-viz.
$ pip install kedro-viz
Eseguendo il comando seguente potremo visualizzare tutte le pipeline e i nodi di dati. Fornisce anche un’opzione per il tracciamento degli esperimenti e la condivisione della visualizzazione.
$ kedro viz run

Visualizzazione della pipeline in Kedro.
5. Luigi
Luigi è un framework open source basato su Python, sviluppato da Spotify, che eccelle nella gestione di processi batch di lunga durata e pipeline di dati complesse. È efficace nella risoluzione delle dipendenze, gestione dei workflow, visualizzazione e ripristino dagli errori, risultando uno strumento potente per orchestrare workflow di dati.
Rispetto ad Airflow, Luigi ha un’API minimale, schedulazione a calendario e una community fedele che può aiutarti con problemi legati all’orchestrazione delle pipeline.
Se sei principiante in Python, potresti trovare difficile costruire ed eseguire le pipeline. Tuttavia, documentazione e guide aiutano a partire velocemente. I log offrono informazioni limitate e la dashboard è perlopiù uno strumento di visualizzazione di DAG e dipendenze.
Abid Ali Awan, Author
Per iniziare con Luigi
Creare una pipeline con Luigi richiede una comprensione della programmazione a oggetti. Iniziamo installando il pacchetto Python di Luigi.
$ pip install luigi
Per sviluppare una semplice pipeline ETL in Luigi, creeremo task interconnessi. Invece di funzioni Python come task, creeremo una classe Python per ciascuno step della pipeline, FetchData, ProcessData e GenerateReport. Ogni classe avrà tre funzioni: requires(), output() e run().
Le funzioni requires() e output() collegano i task, mentre run() esegue il codice di elaborazione. Alla fine costruiremo la pipeline usando l’ultimo task della catena.
import luigi
import pandas as pd
import numpy as np
class FetchData(luigi.Task):
def output(self):
return luigi.LocalTarget('data/fetch_data.csv')
def run(self):
# Simulate fetching data by creating a sample CSV file
data = {
'column1': [1, 2, np.nan, 4],
'column2': ['A', 'B', 'C', np.nan]
}
df = pd.DataFrame(data)
df.to_csv(self.output().path, index=False)
class ProcessData(luigi.Task):
def requires(self):
return FetchData()
def output(self):
return luigi.LocalTarget('data/process_data.csv')
def run(self):
df = pd.read_csv(self.input().path)
# Fill missing values
df['column1'].fillna(df['column1'].mean(), inplace=True)
df['column2'].fillna('B', inplace=True)
df.to_csv(self.output().path, index=False)
class GenerateReport(luigi.Task):
def requires(self):
return ProcessData()
def output(self):
return luigi.LocalTarget('data/generate_report.txt')
def run(self):
df = pd.read_csv(self.input().path)
# Simple data analysis: calculate mean of column1 and value counts of column2
mean_column1 = df['column1'].mean()
value_counts_column2 = df['column2'].value_counts()
with self.output().open('w') as out_file:
out_file.write(f'Mean of column1: {mean_column1}\n')
out_file.write('Value counts of column2:\n')
out_file.write(value_counts_column2.to_string())
if __name__ == '__main__':
luigi.build([GenerateReport()], local_scheduler=True)
Esegui il codice in Jupyter Notebook oppure crea il file Python ed eseguilo dal terminale.

In modo analogo a Luigi, puoi anche imparare a creare una pipeline ETL con Apache Airflow. Il tutorial copre le basi di estrazione, trasformazione e caricamento con Apache Airflow.
Central planner di Luigi
Dobbiamo inizializzare il central planner di Luigi per programmare le esecuzioni o attivarle con un evento.
Avvia lo scheduler digitando il seguente comando nel terminale.
$ luigid
2024-06-22 13:35:18,636 luigi[25056] INFO: logging configured by default settings
2024-06-22 13:35:18,636 luigi.scheduler[25056] INFO: No prior state file exists at /var/lib/luigi-server/state.pickle. Starting with empty state
2024-06-22 13:35:18,640 luigi.server[25056] INFO: Scheduler starting up
Per eseguire la pipeline, apri un nuovo terminale e digita il seguente comando. Il comando di Luigi richiede il nome del file Python e l’ultimo task che vuoi eseguire. In questo caso, il file si chiama luigi_pipe.py e l’ultimo task è GenerateReport.
$ python -m luigi --module luigi_pipe GenerateReport
Se vuoi visualizzare l’esecuzione della pipeline e lo stato dei task, apri semplicemente http://localhost:8082 nel browser.

WebUI del Luigi Central Planner.
Con questo si conclude la nostra panoramica delle 5 migliori alternative ad Airflow! Se vuoi approfondire uno qualsiasi degli esempi presentati nell’articolo, ecco alcune risorse utili:
- Per il codice sorgente e i dati di Prefect, Dagster e Luigi, consulta lo spazio di lavoro DataLab.
- Per il codice sorgente e i dati di Mage AI e Kedro, consulta il repository GitHub.
Considerazioni finali
In questo tutorial abbiamo discusso le principali alternative open source e gratuite ad Airflow. Abbiamo anche visto ciascuno strumento di orchestrazione, costruito ed eseguito una semplice pipeline ETL. Vedere esempi di codice ti aiuterà a decidere quale si adatta meglio al tuo caso d’uso.
Se sei alle prime armi, ti consiglio di iniziare con Prefect o Mage AI: sono user-friendly e con setup semplice. Se invece cerchi strumenti più avanzati e aderenti alle pratiche di ingegneria del software, esplora Dagster, Kedro e Luigi.
Dopo aver letto questo articolo, il passo successivo nel tuo percorso di data engineering è ottenere una certificazione come il percorso di DataCamp Data Engineer in Python per conoscere altri strumenti e costruire una pipeline end-to-end da portare in produzione.
In quanto data scientist certificato, sono appassionato di sfruttare tecnologie all’avanguardia per creare applicazioni di machine learning innovative. Con una solida esperienza in riconoscimento vocale, analisi e reportistica dei dati, MLOps, AI conversazionale e NLP, ho affinato le mie competenze nello sviluppo di sistemi intelligenti in grado di avere un impatto concreto. Oltre alla mia expertise tecnica, sono anche un comunicatore efficace, con il talento di rendere chiari e sintetici concetti complessi. Di conseguenza, sono diventato un blogger molto seguito in ambito data science, condividendo idee ed esperienze con una community in crescita di professionisti dei dati. Attualmente mi concentro sulla creazione e sull’editing di contenuti, lavorando con large language model per sviluppare contenuti potenti e coinvolgenti che possano aiutare aziende e singoli a valorizzare al meglio i propri dati.
