course

Imagine de la autor.
Apache Airflow este un instrument open-source popular pentru orchestrarea datelor, conceput pentru a construi, programa și monitoriza pipeline-uri de date. Include un panou de control care te ajută să gestionezi starea fluxurilor de lucru, fiind o alegere excelentă pentru majoritatea nevoilor de workflow.
Totuși, Airflow nu include unele funcționalități importante, esențiale pentru cerințele moderne și complexe de orchestrare a datelor.
În acest tutorial, vom explora cinci alternative la Airflow care oferă capabilități avansate și abordează unele dintre limitările sale. Mai mult, vom învăța să construim un pipeline ETL simplu cu fiecare instrument, să-l rulăm și să-l vizualizăm în panoul lor de control.
De ce să alegi o alternativă la Airflow?
Airflow este un instrument puternic pentru diverse fluxuri de date, dar are câteva limitări care pot determina companiile să ia în considerare alternative.
Iată câteva motive pentru care ai putea alege o alternativă:
- Curbă de învățare abruptă: Airflow poate fi dificil de învățat, mai ales pentru cei noi în instrumentele de gestionare a fluxurilor de lucru.
- Mentenanță: Necesită mentenanță semnificativă, în special în implementări la scară mare.
- Documentație insuficientă: Utilizatorii au raportat multiple probleme de documentație, ceea ce îngreunează depanarea sau învățarea funcțiilor noi.
- Consum mare de resurse: Airflow poate consuma multe resurse, necesitând capacitate de calcul și memorie substanțiale pentru a rula eficient.
- Flexibilitate limitată pentru non-programatori în Python: Filosofia workflow-as-code se bazează puternic pe Python, ceea ce poate exclude experții de domeniu care nu stăpânesc programarea.
- Scalabilitate: Unii utilizatori raportează dificultăți în scalarea Airflow pentru fluxuri mari.
- Procesare în timp real limitată: Airflow este conceput în principal pentru procesare batch, nu pentru fluxuri de date în timp real.
Înainte să trecem la partea de cod pentru alte instrumente de orchestrare a datelor, este important să înveți cum să scrii pipeline-ul de date folosind Apache Airflow urmând tutorialul Primii pași cu Apache Airflow, ca să poți compara corect alternativele.
Dacă ești complet nou în Airflow, ia în considerare cursul scurt Introducere în Airflow în Python pentru a învăța elementele de bază ale construirii și programării pipeline-urilor de date.
Top 5 alternative la Airflow pentru orchestrarea datelor
Acum, să descriem cele mai bune 5 alternative la Airflow și să arătăm cum se folosesc, cu exemple practice de cod.
1. Prefect
Prefect este un instrument open-source de orchestrare a fluxurilor în Python, creat pentru inginerii moderni de date și machine learning. Oferă un API simplu care îți permite să construiești rapid un pipeline de date și să-l gestionezi printr-un panou interactiv.
Prefect oferă un model de execuție hibrid, ceea ce înseamnă că poți implementa workflow-ul în cloud și să îl rulezi acolo sau să folosești depozitul local.
Comparativ cu Airflow, Prefect vine cu funcții avansate, precum dependențe de sarcini automatizate, declanșatoare bazate pe evenimente, notificări integrate, infrastructură specifică workflow-urilor și partajare de date între sarcini. Aceste capabilități îl fac o soluție puternică pentru gestionarea eficientă a fluxurilor complexe.
Prefect e simplu și vine cu funcții puternice. Practic, mi-a luat 5 minute să rulez codul de exemplu. Îmi place în special cum e gândit UI-ul dashboard-ului, cum poți seta notificări, re-rula pipeline-uri, gestiona și monitoriza totul din Dashboard.
Abid Ali Awan, Author
Citește Airflow vs Prefect: Cum alegi instrumentul potrivit pentru workflow-ul tău de date pentru o comparație detaliată între aceste două instrumente de orchestrare a datelor.
Primii pași cu Prefect
Vom începe proiectul Prefect instalând pachetul Python. Rulează următoarea comandă în terminal.
$ pip install -U prefect
După aceea, vom crea un script Python numit prefect_etl.py și vom scrie următorul cod.
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()
Codul de mai sus definește funcțiile task extract_data(), transform_data(), și load_data() și le execută în serie într-o funcție flow numită etl(). Aceste funcții sunt create folosind decoratori Prefect pentru Python.
Pe scurt, creăm un DataFrame pandas, îl transformăm, apoi afișăm rezultatul final cu print. Este o modalitate simplă de a simula un pipeline ETL.
Pentru a executa workflow-ul, rulează scriptul Python folosind comanda de mai jos.
$ python prefect_etl.py
După cum putem vedea, rularea workflow-ului s-a finalizat cu succes.

Jurnale de rulare Prefect flow.
Implementarea flow-ului
Acum vom implementa workflow-ul, astfel încât să îl putem rula programat sau declanșa pe baza unui eveniment. Implementarea flow-ului îți permite, de asemenea, să monitorizezi și să gestionezi mai multe workflow-uri într-un mod centralizat.
Pentru a implementa flow-ul, vom folosi Prefect CLI. Funcția deploy necesită numele fișierului Python, numele funcției flow din fișier și numele deployment-ului. În acest caz, numim acest deployment „simple_etl”.
$ prefect deploy prefect_etl.py:etl -n 'simple_etl'
După rularea scriptului de mai sus în terminal, este posibil să primești mesajul că nu ai un worker pool pentru a rula deployment-ul. Pentru a crea worker pool-ul, folosește următoarea comandă.
$ prefect worker start --pool 'datacamp'
Acum că avem un worker pool, vom deschide o nouă fereastră de terminal și vom rula deployment-ul. Comanda prefect deployment run necesită „<flow-function-name>/<deployment-name>” ca argument, așa cum este arătat mai jos.
$ prefect deployment run 'etl/simple_etl
Ca rezultat al rulării deployment-ului, vei primi mesajul că workflow-ul rulează. De obicei, flow run-ul creat primește un nume aleator, în cazul meu 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>
Pentru a vedea jurnalul complet, revino la fereastra de terminal unde ai pornit worker pool-ul.

Rezumatul rularii flow-ului Prefect.
Trebuie să inițializezi serverul web Prefect pentru a vizualiza rularea flow-ului într-un mod mai prietenos și pentru a gestiona alte workflow-uri.
$ prefect server start
După executarea comenzii de mai sus, ar trebui să fii redirecționat către dashboard-ul Prefect. Alternativ, poți accesa direct adresa http://127.0.0.1:4200 în browser.

Interfața serverului web Prefect
Dashboard-ul îți permite să re-rulezi workflow-ul, să vezi jurnalele, să verifici work pool-urile, să setezi notificări și să alegi alte opțiuni avansate. Este o soluție completă pentru nevoile moderne de orchestrare a datelor.
Pentru a învăța cum să construiești și să execuți pipeline-uri de machine learning cu Prefect, poți urma tutorialul Folosirea Prefect pentru workflow-uri de Machine Learning.
2. Dagster
Dasgter este un framework open-source conceput pentru inginerii de date, pentru a defini, programa și monitoriza pipeline-uri de date. Este foarte scalabil și facilitează colaborarea între diverse echipe de date.
Dagster le permite utilizatorilor să își definească activele de date ca funcții Python folosind decoratori. Odată definite, acestea pot fi executate ușor prin programare sau declanșatoare bazate pe evenimente.
Comparativ cu Airflow, Dagster permite dezvoltarea, testarea și revizuirea pipeline-ului local, oferă o abordare bazată pe active pentru orchestrare și este nativ pentru cloud și containere.
În loc să mă gândesc la workflow în termeni de pași și fluxuri, a trebuit să îmi schimb abordarea și să construiesc un pipeline folosind active de date. În rest, construirea și execuția unui pipeline ETL simplu au fost destul de ușoare. De asemenea, serverul web e relativ minimal, dar oferă toate informațiile pentru a monitoriza activele, rularile și deployment-urile.
Abid Ali Awan, Author
Primii pași cu Dagster
Vom crea un pipeline ETL simplu, îl vom executa și îl vom vizualiza folosind serverul web Dagster. Similar cu dashboard-ul Prefect, serverul web Dagster oferă modalități centralizate de a monitoriza mai multe workflow-uri și de a programa rulări și active.
Vom începe prin a instala pachetul Python.
$ pip install dagster -q
Apoi, vom crea trei funcții Python pentru extragerea, transformarea și încărcarea datelor. Aceste funcții se numesc create_dirty_data(), clean_data(), și load_cleaned_data() în cod. Folosind decoratorul @asset, vom declara funcțiile ca active de date în Dagster.
În continuare, vom crea jobul de active (variabila job) folosind toate activele (variabila all_assets) și apoi vom crea definiția de active (variabila defs).
Poți sări peste partea de definire a asset-urilor, dar devine importantă dacă vrei să îți programezi rularile, să rulezi mai multe joburi și să configurezi senzori.
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)
Poți rula codul de mai sus într-un Jupyter Notebook sau poți crea fișierul Python și să-l rulezi.
Ca urmare a executării codului, vom obține un jurnal complet al rularii workflow-ului.

Serverul web Dagster
Pentru a vizualiza activele și rulările jobului, trebuie să instalăm și să rulăm serverul web Dagster. Serverul web îți permite să rulezi joburile, să materializezi active individuale și să monitorizezi mai multe joburi simultan.
$ pip install dagster-webserver
Pentru a iniția serverul Dagster, vom folosi CLI-ul Dagster și îi vom furniza locația fișierului Python. În acest caz, am numit fișierul dagster_pipe.py.
$ dagster dev -f dagster_pipe.py
Comanda de mai sus va porni automat serverul web în browser. Alternativ, poți accesa direct adresa http://127.0.0.1:3000 în browser.

Interfața serverului web Dagster.
Până acum am implementat doar jobul. Pentru a rula workflow-ul, mergi la fila „Runs” și apasă pe butonul „Launch a new run”.
Rularea ar trebui să se încheie cu succes! Pentru a vedea jurnalele, dă clic pe ID-ul rularii care te interesează.

Jurnalele rularilor Dagster.
3. Mage AI
Mage AI este un framework open-source hibrid pentru orchestrarea datelor. Hibrid înseamnă că ai flexibilitatea unui Jupyter Notebook și controlul unui cod modular.
Oricine, chiar și cu cunoștințe limitate de Python, poate construi, rula și monitoriza pipeline-uri de date. În loc să scrii și să rulezi direct un fișier Python, vei crea un proiect Mage AI și îl vei lansa în dashboard, unde poți construi, rula și gestiona pipeline-urile tale.
Comparativ cu Airflow, Mage AI oferă o interfață prietenoasă și ușurință în utilizare, fiind o alegere excelentă pentru cei noi în ingineria de date. Este conceput cu scalabilitatea în minte și poate gestiona eficient volume mari de date și structuri complexe de pipeline-uri.
M-am simțit ciudat pentru că era complet diferit față de ce sunt obișnuit. A trebuit să instalez și să lansez interfața web a Mage AI. Ar fi trebuit să fie ușor, dar mi s-a părut dificil să construiesc și să rulez pipeline-ul ETL. Pe de altă parte, înțeleg de ce acest design unic poate fi atractiv pentru începători, fiind în esență drag and drop și apăsat pe butoane.
Abid Ali Awan, Author
Primii pași cu Mage AI
Pornirea Mage AI e destul de simplă. Trebuie doar să instalăm pachetul Python Mage AI.
$ pip install mage-ai
Și să pornim proiectul Mage AI.
$ mage start mage_ai_etl
Comanda de mai sus va iniția serverul web. După cum am menționat, toată editarea codului, rularea joburilor și monitorizarea lor se fac prin interfața Mage AI.

Interfața Mage AI.
Apasă pe „+ New pipeline” pentru a crea primul tău pipeline ETL. Eu l-am numit „simple_etl”.

Crearea noului pipeline în Mage AI.
Apoi, interfața îți va cere să adaugi un modul pentru a începe să scrii cod. Selectează modulul „Data Loader” și scrie următorul cod Python.
Aici, declarăm o funcție create_sample_csv(), care e primul pas în pipeline. Folosim decoratorul Mage AI @data_loader. Definim și o funcție test_output() care verifică dacă există output-ul. Asta ajută la gestionarea dependențelor dintre sarcini.
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'

Crearea blocului data loader în Mage AI.
În continuare, creează un alt modul numit „Transformer” și adaugă funcția clean_data(), așa cum este în codul de mai jos.
Poți ignora funcția test(); trebuie doar să adaugi funcția principală de transformare, 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'
La fel, creează un modul „Data Exporter” și adaugă următorul cod. Codul declară o funcție de încărcare a datelor, export_data_to_csv(), care salvează datele transformate într-un fișier 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")
Pentru a rula pipeline-ul, mergi la fila „Trigger” și apasă pe „Run@once”.

Rularea pipeline-ului în Mage AI.
Pentru a vedea jurnalele rularii, mergi la fila „Runs” și apasă pe butonul „Logs” la pipeline-ul rulat recent.

Jurnalele rularii flow-ului în Mage AI.
4. Kedro
Kedro este un alt framework open-source popular pentru orchestrarea datelor, ușor diferit de celelalte instrumente. A fost creat pentru inginerii de machine learning și împrumută multe concepte din ingineria software, aplicându-le proiectelor de ML.
Kedro este conceput să fie foarte modular, ceea ce înseamnă că, chiar și pentru a exporta un set de date, trebuie să creezi un catalog de date care specifică locația și tipul datelor, asigurând o gestionare standardizată și eficientă pe tot parcursul pipeline-ului.
Pentru a înțelege cum se potrivește Kedro în ecosistemul ML, poți explora diverse instrumente MLOps citind articolul 25 de instrumente MLOps esențiale pe care trebuie să le știi în 2024.
Comparativ cu Airflow, API-ul Kedro este mai simplu pentru construirea unui pipeline de date. Se concentrează mai mult pe ingineria ML și oferă categorizarea și versionarea datelor.
Partea de cod este destul de directă, dar apar probleme când vrei să îți execuți pipeline-ul. Trebuie să creezi un catalog de date, să înregistrezi pipeline-ul și să înțelegi structura proiectului Kedro. Aș spune că e mai provocator decât Dagster și Prefect. Totuși, înțeleg de ce e proiectat așa: pentru a face pipeline-ul tău de date fiabil și fără erori.
Abid Ali Awan, Author
Primii pași cu Kedro
Construirea unui pipeline de date cu Kedro e un joc diferit. Framework-ul este modular și trebuie să înțelegi structura proiectului și pașii implicați pentru a executa cu succes workflow-ul.
Începe prin a instala pachetul Python Kedro.
$ pip install kedro
Inițializează proiectul Kedro.
$ kedro new --name=kedro_etl --tools=none --example=n
Mergi în directorul proiectului.
$ cd kedro-etl
Creează un folder în directorul pipelines numit data_processing.
$ mkdir -p src/kedro_etl/pipelines/data_processing
Creează un fișier Python numit kedro_pipe.py și deschide-l în IDE-ul preferat, de exemplu Visual Studio Code.
$ code src/kedro_etl/pipelines/data_processing/kedro_pipe.py
Scriptul Python ar trebui să conțină funcțiile de extragere, transformare și încărcare, care sunt noduri în pipeline. În acest caz, acestea sunt funcțiile create_sample_data(), clean_data(), și load_and_process_data().
Apoi, unim aceste noduri folosind clasa Kedro Pipeline în cadrul funcției create_pipeline(). În funcția pipeline, definim noduri, iar fiecare nod are inputs, outputs și un name pentru nod.
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",
),
]
)
Dacă rulăm pipeline-ul fără a crea catalogul de date, nu ne va exporta datele. Așadar, trebuie să mergem la fișierul conf/base/catalog.yml și să-l edităm, oferind configurația dataset-urilor.
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
Trebuie, de asemenea, să includem noul fișier Python în registrul pipeline-ului. Pentru asta, mergi la fișierul Python src/simple_etl/pipeline_registry.py și include următorul cod.
"""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,
}
Rulează pipeline-ul și vezi jurnalele live în terminal rulând comanda de mai jos.
$ kedro run

Jurnalele rularii pipeline-ului Kedro.
După rularea pipeline-ului, fișierele tale vor fi stocate în format CSV în locația definită în catalogul de date.

Fișierele rezultate în urma rularii pipeline-ului Kedro.
Dacă întâmpini probleme la rularea pipeline-ului, ia în considerare instalarea Kedro cu toate extensiile.
$ pip install "kedro[all]"
Vizualizarea în Kedro
Putem vizualiza și partaja pipeline-urile instalând instrumentul kedro-viz.
$ pip install kedro-viz
Apoi, executarea comenzii următoare ne va permite să vizualizăm toate pipeline-urile de date și nodurile. Oferă și opțiuni pentru urmărirea experimentelor și posibilitatea de a partaja vizualizarea pipeline-ului.
$ kedro viz run

Vizualizarea pipeline-ului Kedro.
5. Luigi
Luigi este un framework open-source, bazat pe Python, dezvoltat de Spotify, care excelează în gestionarea proceselor batch de lungă durată și a pipeline-urilor complexe de date. Este bun la rezolvarea dependențelor, managementul workflow-urilor, vizualizare și recuperare după eșecuri, fiind un instrument puternic pentru orchestrarea fluxurilor de date.
Comparativ cu Airflow, Luigi are un API minimalist, programare pe bază de calendar și o comunitate loială care te va ajuta cu orice problemă legată de pipeline-ul de orchestrare a datelor.
Dacă ești începător în Python, s-ar putea să ți se pară dificil să construiești și să rulezi pipeline-uri. Totuși, documentația și ghidurile te pot ajuta să începi rapid. Jurnalele oferă informații limitate, iar dashboard-ul este doar un instrument de vizualizare pentru DAG-uri și dependențe.
Abid Ali Awan, Author
Primii pași cu Luigi
Crearea unui pipeline de date în Luigi necesită înțelegerea programării orientate pe obiecte. Să începem prin a instala pachetul Python Luigi.
$ pip install luigi
Pentru a dezvolta un pipeline ETL simplu în Luigi, vom crea taskuri interconectate. În loc să creăm funcții Python ca taskuri, vom crea o clasă Python pentru fiecare pas din pipeline, FetchData, ProcessData și GenerateReport. Fiecare clasă va avea trei funcții numite: requires(), output() și run().
Funcțiile requires() și output() vor conecta taskurile, iar funcția run() va executa codul de procesare. La final, vom construi pipeline-ul folosind ultimul task din pipeline.
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)
Rulează codul de mai sus în Jupyter Notebook sau creează fișierul Python și rulează-l din terminal.

La fel ca în cazul Luigi, poți învăța și cum să construiești un pipeline ETL cu Apache Airflow. Tutorialul acoperă elementele de bază ale extragerii, transformării și încărcării datelor cu Apache Airflow.
Planificatorul central Luigi
Trebuie să inițializăm planificatorul central Luigi pentru a programa rularile pipeline-ului sau a le declanșa printr-un eveniment.
Pornește scheduler-ul tastând următoarea comandă în terminal.
$ 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
Pentru a rula pipeline-ul, deschide un nou terminal și tastează comanda următoare. Comanda Luigi necesită numele fișierului Python și ultimul task pe care vrem să îl executăm. În acest caz, numele fișierului este luigi_pipe.py, iar ultimul nostru task Luigi este GenerateReport.
$ python -m luigi --module luigi_pipe GenerateReport
Dacă vrei să vizualizezi rularea pipeline-ului și statusul taskurilor, poți accesa simplu http://localhost:8082 în browser.

Interfața web Luigi Central Planner.
Asta încheie prezentarea noastră a celor 5 cele mai bune alternative la Airflow! Dacă vrei să aprofundezi oricare dintre exemplele prezentate în acest articol, iată câteva resurse utile:
- Pentru codul sursă și datele Prefect, Dagster și Luigi, consultă workspace-ul DataLab.
- Pentru codul sursă și datele Mage AI și Kedro, consultă repository-ul GitHub.
Gânduri finale
În acest tutorial, am discutat despre cele mai bune alternative open-source, gratuite la Airflow. De asemenea, am învățat despre fiecare instrument de orchestrare a datelor și am construit și executat un pipeline ETL simplu. Văzând exemple de cod, îți va fi mai ușor să decizi care se potrivește cel mai bine cazului tău de utilizare.
Dacă ești începător, îți sugerez să începi cu Prefect sau Mage AI, deoarece sunt prietenoase și au o configurare simplă. Totuși, dacă cauți instrumente mai avansate care respectă practicile de inginerie software, îți recomand să explorezi Dagster, Kedro și Luigi.
După ce ai parcurs acest articol, pasul firesc în drumul tău în ingineria datelor este să obții o certificare precum Data Engineer in Python de la DataCamp, pentru a învăța despre alte instrumente și a construi un pipeline de date cap-coadă pe care îl poți implementa în producție.