Sari la conținutul principal

Top 5 alternative la Airflow pentru orchestrarea datelor (cu exemple de cod)

Explorează cinci alternative la Airflow pentru orchestrarea datelor, cu exemple de cod pentru construirea, rularea și vizualizarea unui pipeline ETL simplu.
Actualizat 31 aug. 2026  · 13 min. citire

Explorează cu AI

ChatGPTClaudePerplexity

Choose an Airflow Alternatives meme template

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ă:

  1. 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.
  2. Mentenanță: Necesită mentenanță semnificativă, în special în implementări la scară mare.
  3. Documentație insuficientă: Utilizatorii au raportat multiple probleme de documentație, ceea ce îngreunează depanarea sau învățarea funcțiilor noi. 
  4. Consum mare de resurse: Airflow poate consuma multe resurse, necesitând capacitate de calcul și memorie substanțiale pentru a rula eficient.
  5. 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.
  6. Scalabilitate: Unii utilizatori raportează dificultăți în scalarea Airflow pentru fluxuri mari.
  7. 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 AwanAuthor

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.

Prefect flow run logs

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.

Prefect flow run summary

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.

Prefect web server UI

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 AwanAuthor

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. 

Dagster execution summary

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.

Dagster Web server

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ă.

Dagster runs detailed view

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 AwanAuthor

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.

Mage AI UI

Interfața Mage AI.

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

Creating the new pipeline in Mage AI

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'

Creating the data loader block in Mage AI

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”.

Running the pipeline in Mage AI

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.

Mage AI flow run logs

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 AwanAuthor

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

Logs of Kedro pipeline 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.

Output files of Kedro pipeline run

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

Kedro Visualization

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 AwanAuthor

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. 

Luigi Execution Summary

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.

Luigi Central Planner webUI

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:

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.

Subiecte
Data Engineering
Data Science

Află mai multe despre ingineria datelor cu aceste cursuri!

course

Introducere în Data Engineering

4 oră
129.7K
Descoperă lumea ingineriei datelor în acest curs scurt, cu instrumente și teme precum ETL și cloud computing.
Vezi detaliiRight Arrow
Începeți Cursul
Vezi mai multRight Arrow