Przejdź do głównej treści

Top 5 alternatyw dla Airflow do orkiestracji danych (z przykładami kodu)

Poznaj pięć alternatyw dla Airflow z przykładami kodu do budowy, uruchamiania i wizualizacji prostego potoku ETL.
Zaktualizowano 31 sie 2026  · 13 min Czytać

Eksploruj z AI

ChatGPTClaudePerplexity

Choose an Airflow Alternatives meme template

Obraz autora.

Apache Airflow to popularne, otwartoźródłowe narzędzie do orkiestracji danych, zaprojektowane do budowania, harmonogramowania i monitorowania potoków danych. Zawiera pulpit do zarządzania stanem przepływów pracy, co czyni je doskonałym narzędziem dla większości potrzeb związanych z workflow.

Jednak Airflow nie ma kilku ważnych funkcji, które mogą być kluczowe przy złożonych, nowoczesnych wymaganiach orkiestracji danych.

W tym poradniku poznamy pięć alternatyw dla Airflow, które oferują rozszerzone możliwości i adresują niektóre z jego ograniczeń. Ponadto zbudujemy prosty potok ETL w każdym narzędziu, uruchomimy go i zwizualizujemy na ich pulpitach.

Dlaczego warto wybrać alternatywę dla Airflow? 

Airflow to potężne narzędzie dla różnych przepływów danych, ale ma kilka ograniczeń, przez które firmy mogą rozważać alternatywy. 

Oto powody, dla których możesz wybrać inne rozwiązanie:

  1. Wysoka bariera wejścia: Airflow może być trudny do nauki, zwłaszcza dla osób początkujących w narzędziach do zarządzania workflow.
  2. Utrzymanie: Wymaga znaczących nakładów na utrzymanie, zwłaszcza przy wdrożeniach na dużą skalę.
  3. Niewystarczająca dokumentacja: Użytkownicy zgłaszają wiele problemów z dokumentacją, co utrudnia rozwiązywanie problemów lub poznawanie nowych funkcji. 
  4. Zasobożerność: Airflow może być zasobożerny, wymagając sporej mocy obliczeniowej i pamięci do efektywnego działania.
  5. Ograniczona elastyczność dla osób niekorzystających z Pythona: Filozofia „workflow jako kod” mocno opiera się na Pythonie, co może wykluczać ekspertów domenowych bez biegłości programistycznej.
  6. Skalowalność: Część użytkowników zgłasza trudności ze skalowaniem Airflow dla dużych workflow.
  7. Ograniczone przetwarzanie w czasie rzeczywistym: Airflow jest przede wszystkim zaprojektowany do przetwarzania wsadowego, a nie strumieni danych w czasie rzeczywistym.

Zanim przejdziemy do części kodowej innych narzędzi do orkiestracji danych, warto nauczyć się, jak pisać potok danych w Apache Airflow, korzystając z poradnika Getting Started with Apache Airflow, aby móc uczciwie porównać alternatywy.

Jeśli dopiero zaczynasz przygodę z Airflow, rozważ krótki kurs Introduction to Airflow in Python, aby poznać podstawy budowania i harmonogramowania potoków danych.

5 najlepszych alternatyw dla Airflow do orkiestracji danych

Przejdźmy teraz do opisu 5 najlepszych alternatyw dla Airflow i pokażmy, jak z nich korzystać na praktycznych przykładach kodu.

1. Prefect

Prefect to otwartoźródłowe narzędzie do orkiestracji workflow w Pythonie, stworzone dla nowoczesnych inżynierów danych i ML. Oferuje prosty interfejs API, który pozwala szybko zbudować potok danych i zarządzać nim przez interaktywny pulpit. 

Prefect oferuje hybrydowy model wykonania, co oznacza, że możesz wdrożyć workflow w chmurze i tam go uruchamiać lub korzystać z lokalnego repozytorium.

W porównaniu z Airflow, Prefect ma zaawansowane funkcje, takie jak automatyczne zależności zadań, wyzwalacze zdarzeń, wbudowane powiadomienia, infrastruktura specyficzna dla workflow oraz współdzielenie danych między zadaniami. Dzięki temu świetnie sprawdza się przy efektywnym zarządzaniu złożonymi przepływami pracy.

Prefect jest prosty i ma potężne funkcje. Uruchomienie przykładowego kodu zajęło mi dosłownie 5 minut. Szczególnie podoba mi się projekt interfejsu pulpitu, możliwość ustawienia powiadomień, ponownego uruchamiania potoków oraz zarządzania i monitorowania wszystkiego przez Dashboard.

Abid Ali AwanAuthor

Przeczytaj wpis Airflow vs Prefect: Deciding Which is Right For Your Data Workflow, aby poznać szczegółowe porównanie tych dwóch narzędzi do orkiestracji danych. 

Pierwsze kroki z Prefect

Zaczniemy projekt Prefect od instalacji pakietu Pythona. Uruchom poniższe polecenie w terminalu.

$ pip install -U prefect

Następnie utworzymy skrypt Pythona o nazwie prefect_etl.py i wpiszemy do niego poniższy kod.

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()

Powyższy kod definiuje funkcje zadań extract_data(), transform_data(), i load_data() oraz wykonuje je sekwencyjnie wewnątrz funkcji flow o nazwie etl(). Funkcje te tworzymy, korzystając z dekoratorów Prefect dla Pythona. 

W skrócie: tworzymy obiekt pandas DataFrame, przekształcamy go, a następnie wyświetlamy wynik za pomocą print. To prosty sposób na zasymulowanie potoku ETL.

Aby wykonać workflow, po prostu uruchom skrypt Pythona poniższym poleceniem.

$ python prefect_etl.py 

Jak widać, nasz przebieg workflow zakończył się powodzeniem.

Prefect flow run logs

Logi przebiegu flow w Prefect.

Wdrażanie flow

Teraz wdrożymy nasz workflow, aby móc uruchamiać go zgodnie z harmonogramem lub wyzwalać zdarzeniami. Wdrożenie flow pozwala też centralnie monitorować i zarządzać wieloma workflow.

Do wdrożenia flow użyjemy Prefect CLI. Funkcja deploy wymaga nazwy pliku Pythona, nazwy funkcji flow w pliku oraz nazwy wdrożenia. W tym przypadku nazywamy to wdrożenie „simple_etl”.

$ prefect deploy prefect_etl.py:etl -n 'simple_etl'

Po uruchomieniu powyższego polecenia w terminalu możesz otrzymać komunikat, że nie masz puli workerów do uruchomienia wdrożenia. Aby utworzyć pulę workerów, użyj poniższego polecenia.

$ prefect worker start --pool 'datacamp'

Skoro mamy już pulę workerów, uruchomimy kolejne okno terminala i wykonamy wdrożenie. Polecenie prefect deployment run wymaga argumentu w formacie „<nazwa-funkcji-flow>/<nazwa-wdrożenia>”, jak w poleceniu poniżej.

$ prefect deployment run 'etl/simple_etl

W wyniku uruchomienia wdrożenia otrzymasz komunikat, że workflow działa. Zwykle utworzony przebieg flow dostaje losową nazwę, w moim przypadku 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>

Aby zobaczyć pełne logi, przełącz się do okna terminala, w którym uruchomiłeś pulę workerów.

Prefect flow run summary

Podsumowanie przebiegu flow w Prefect.

Musisz uruchomić serwer WWW Prefect, aby wygodniej wizualizować przebieg flow i zarządzać innymi workflow.

$ prefect server start 

Po wykonaniu powyższego polecenia powinieneś zostać przekierowany na pulpit Prefect. Alternatywnie możesz przejść bezpośrednio pod adres http://127.0.0.1:4200 w przeglądarce.

Prefect web server UI

Interfejs serwera WWW Prefect

Pulpit pozwala ponownie uruchamiać workflow, przeglądać logi, sprawdzać pule robocze, ustawiać powiadomienia i wybierać inne zaawansowane opcje. To kompletne rozwiązanie dla nowoczesnych potrzeb orkiestracji danych.

Aby dowiedzieć się, jak budować i uruchamiać potoki uczenia maszynowego w Prefect, skorzystaj z poradnika Using Prefect for Machine Learning Workflows.

2. Dagster

Dagster to otwartoźródłowy framework zaprojektowany dla inżynierów danych do definiowania, harmonogramowania i monitorowania potoków danych. Jest wysoce skalowalny i ułatwia współpracę między różnymi zespołami danych. 

Dagster pozwala definiować zasoby danych jako funkcje Pythona za pomocą dekoratorów. Gdy zasoby zostaną zdefiniowane, można je bezproblemowo wykonywać poprzez harmonogram lub wyzwalacze zdarzeń.

W porównaniu z Airflow, Dagster umożliwia lokalne tworzenie, testowanie i przeglądanie potoku, oferuje podejście oparte na zasobach oraz jest natywny dla chmury i kontenerów.

Zamiast myśleć o workflow w kategoriach kroków i przepływów, musiałem zmienić podejście i zbudować potok w oparciu o zasoby danych. Poza tym stworzenie i uruchomienie prostego potoku ETL było dość proste. Serwer WWW jest raczej minimalistyczny, ale dostarcza wszystkich informacji do monitorowania zasobów, przebiegów i wdrożeń.

Abid Ali AwanAuthor

Pierwsze kroki z Dagster

Utworzymy prosty potok ETL, uruchomimy go i zwizualizujemy za pomocą serwera WWW Dagster. Podobnie jak pulpit Prefect, serwer WWW Dagster zapewnia scentralizowane sposoby monitorowania wielu workflow oraz harmonogramowania przebiegów i zasobów.

Zaczniemy od instalacji pakietu Pythona.

$ pip install dagster -q

Następnie utworzymy trzy funkcje Pythona do ekstrakcji, transformacji i ładowania danych. W kodzie noszą one nazwy create_dirty_data()clean_data() oraz load_cleaned_data(). Za pomocą dekoratora @asset zadeklarujemy te funkcje jako zasoby danych w Dagster.

Następnie utworzymy zadanie dla zasobów (zmienna job) z wykorzystaniem wszystkich zasobów (zmienna all_assets), a następnie zdefiniujemy definicję zasobów (zmienna defs). 

Możesz pominąć część z definicją zasobów, ale staje się ona ważna, jeśli chcesz harmonogramować przebiegi, uruchamiać wiele zadań i konfigurować sensory.

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)

Powyższy kod możesz uruchomić w Jupyter Notebook lub utworzyć plik Pythona i go uruchomić. 

W wyniku wykonania kodu otrzymamy pełny log przebiegu workflow. 

Dagster execution summary

Serwer WWW Dagster

Aby wizualizować zasoby i przebiegi zadań, musimy zainstalować i uruchomić serwer WWW Dagster. Serwer pozwala uruchamiać zadania, materializować pojedyncze zasoby i monitorować wiele zadań jednocześnie.

$ pip install dagster-webserver

Aby zainicjować serwer Dagster, użyjemy CLI Dagster i podamy lokalizację pliku Pythona. W tym przypadku nazwałem plik dagster_pipe.py.

$ dagster dev -f dagster_pipe.py  

Powyższe polecenie automatycznie uruchomi serwer WWW w przeglądarce. Alternatywnie możesz przejść bezpośrednio pod adres http://127.0.0.1:3000 w przeglądarce.

Dagster Web server

Interfejs serwera WWW Dagster.

Na razie wdrożyliśmy tylko zadanie. Aby uruchomić workflow, przejdź do zakładki „Runs” i kliknij przycisk „Launch a new run”. 

Przebieg powinien zakończyć się powodzeniem! Aby zobaczyć logi, kliknij identyfikator interesującego cię przebiegu.

Dagster runs detailed view

Logi przebiegu w Dagster.

3. Mage AI

Mage AI to otwartoźródłowy, hybrydowy framework do orkiestracji danych. Hybrydowy oznacza, że łączy elastyczność Jupyter Notebooka z kontrolą modularnego kodu. 

Każdy, nawet z ograniczoną znajomością Pythona, może budować, uruchamiać i monitorować potoki danych. Zamiast pisać i uruchamiać plik Pythona bezpośrednio, utworzysz projekt Mage AI i uruchomisz go w panelu, gdzie zbudujesz, uruchomisz i będziesz zarządzać potokami danych.

W porównaniu z Airflow, Mage AI zapewnia przyjazny interfejs i łatwość użycia, co czyni go świetnym wyborem dla osób nowych w inżynierii danych. Zaprojektowano go z myślą o skalowalności, dzięki czemu sprawnie obsługuje duże wolumeny danych i złożone struktury potoków.

Czułem się dziwnie, bo to było zupełnie inne niż to, do czego jestem przyzwyczajony. Musiałem zainstalować i uruchomić interfejs webowy Mage AI. Miało być łatwo, ale trudno było mi zbudować i uruchomić potok ETL. Z drugiej strony widzę, dlaczego ten unikalny projekt może przyciągać osoby nowe w branży — to w zasadzie przeciągnij i upuść oraz klikanie przycisków.

Abid Ali AwanAuthor

Pierwsze kroki z Mage AI

Uruchomienie Mage AI jest dość proste. Wystarczy zainstalować pakiet Mage AI dla Pythona.

$ pip install mage-ai

I wystartować projekt Mage AI. 

$ mage start mage_ai_etl 

Powyższe polecenie zainicjuje serwer WWW. Jak wspomniano, cały kod edytuje się, uruchamia zadania i monitoruje je przez interfejs Mage AI.

Mage AI UI

Interfejs Mage AI.

Kliknij „+ New pipeline”, aby utworzyć pierwszy potok ETL. Swój nazwałem „simple_etl”.

Creating the new pipeline in Mage AI

Tworzenie nowego potoku w Mage AI.

Następnie interfejs poprosi o dodanie modułu do rozpoczęcia kodowania. Wybierz moduł „Data Loader” i wpisz poniższy kod w Pythonie. 

Tutaj deklarujemy funkcję create_sample_csv(), czyli pierwszy krok w naszym potoku. Używamy dekoratora Mage AI  @data_loader. Definiujemy też funkcję test_output(), która sprawdza, czy wynik istnieje. Pomaga to w zarządzaniu zależnościami zadań.

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

Tworzenie bloku data loader w Mage AI.

Podobnie, utwórz kolejny moduł „Transformer” i dodaj funkcję clean_data(), jak w kodzie poniżej. 

Możesz zignorować funkcję test(); wystarczy dodać główną funkcję transformującą, 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'

Analogicznie utwórz moduł „Data Exporter” i dodaj poniższy kod. Kod deklaruje funkcję ładującą dane, export_data_to_csv(), która zapisuje przekształcone dane do pliku 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")

Aby uruchomić potok, przejdź do zakładki „Trigger” i kliknij „Run@once”.

Running the pipeline in Mage AI

Uruchamianie potoku w Mage AI.

Aby zobaczyć logi przebiegu, przejdź do zakładki „Runs” i kliknij przycisk „Logs” przy ostatnio uruchomionym potoku.

Mage AI flow run logs

Logi przebiegu flow w Mage AI.

4. Kedro

Kedro to kolejne popularne, otwartoźródłowe narzędzie do orkiestracji danych, nieco inne od pozostałych. Powstało z myślą o inżynierach uczenia maszynowego i zapożycza wiele koncepcji z inżynierii oprogramowania, stosując je w projektach ML.

Kedro jest wysoce modułowe, co oznacza, że nawet aby wyeksportować zbiór danych, musisz utworzyć katalog danych określający lokalizację i typ danych, zapewniając standaryzowane i wydajne zarządzanie danymi w całym potoku.

Aby zrozumieć, jak Kedro wpisuje się w ekosystem ML, możesz poznać różne narzędzia MLOps, czytając artykuł 25 Top MLOps Tools You Need to Know in 2024.

W porównaniu z Airflow, API Kedro jest prostsze do budowy potoku danych. Skupia się bardziej na inżynierii ML i oferuje kategoryzację oraz wersjonowanie danych.

Pisanie kodu jest dość proste, ale problemy zaczynają się przy uruchamianiu potoku. Musisz utworzyć katalog danych, zarejestrować potok i zrozumieć strukturę projektu Kedro. Powiedziałbym, że jest to trudniejsze niż w Dagster i Prefect. Rozumiem jednak, dlaczego zaprojektowano to w ten sposób: aby twój potok danych był niezawodny i wolny od błędów.

Abid Ali AwanAuthor

Pierwsze kroki z Kedro

Budowanie potoku danych w Kedro to inna bajka. Framework jest modułowy i musisz zrozumieć strukturę projektu oraz poszczególne kroki, aby pomyślnie wykonać workflow. 

Zacznij od instalacji pakietu Kedro dla Pythona. 

$ pip install kedro

Zainicjuj projekt Kedro. 

$ kedro new --name=kedro_etl --tools=none --example=n 

Przejdź do katalogu projektu. 

$ cd kedro-etl  

Utwórz folder w katalogu pipelines o nazwie data_processing.

$ mkdir -p src/kedro_etl/pipelines/data_processing  

Utwórz plik Pythona o nazwie kedro_pipe.py i otwórz go w ulubionym IDE, np. Visual Studio Code.

$ code src/kedro_etl/pipelines/data_processing/kedro_pipe.py

Skrypt Pythona powinien zawierać funkcje extract, transform i load, które są węzłami w potoku. W tym przypadku są to funkcje create_sample_data()clean_data(), oraz load_and_process_data().

Następnie łączymy te węzły za pomocą klasy Kedro Pipeline wewnątrz funkcji create_pipeline(). W funkcji potoku definiujemy węzły, a każdy węzeł ma inputs, outputs i 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",
            ),
        ]
    )

Jeśli uruchomimy potok bez utworzenia katalogu danych, nie wyeksportuje on danych. Musimy więc przejść do pliku conf/base/catalog.yml i uzupełnić go konfiguracją zbiorów danych.

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

Musimy także dodać nowo utworzony plik Pythona do rejestru potoków. W tym celu przejdź do pliku Pythona src/simple_etl/pipeline_registry.py i dodaj poniższy kod. 

"""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,
    }

Uruchom potok i podglądaj logi na żywo w terminalu, uruchamiając poniższe polecenie.

$ kedro run

Logs of Kedro pipeline run

Logi uruchomienia potoku w Kedro.

Po uruchomieniu potoku twoje pliki zostaną zapisane w formacie CSV w lokalizacjach zdefiniowanych w katalogu danych.

Output files of Kedro pipeline run

Pliki wyjściowe po uruchomieniu potoku w Kedro.

Jeśli napotkasz problemy z uruchomieniem potoku, rozważ instalację Kedro ze wszystkimi rozszerzeniami. 

$ pip install "kedro[all]"

Wizualizacja w Kedro

Możemy wizualizować i udostępniać nasze potoki, instalując narzędzie kedro-viz

$ pip install kedro-viz

Następnie wykonanie poniższego polecenia pozwoli zwizualizować wszystkie potoki i węzły danych. Dostępna jest też opcja śledzenia eksperymentów i udostępniania wizualizacji potoku.

$ kedro viz run

Kedro Visualization

Wizualizacja potoku w Kedro.

5. Luigi

Luigi to otwartoźródłowy framework oparty na Pythonie, opracowany przez Spotify, który świetnie radzi sobie z zarządzaniem długotrwałymi procesami wsadowymi i złożonymi potokami danych. Sprawdza się w rozwiązywaniu zależności, zarządzaniu workflow, wizualizacji i odzyskiwaniu po błędach, co czyni go potężnym narzędziem do orkiestracji przepływów danych. 

W porównaniu z Airflow, Luigi ma minimalne API, harmonogram kalendarzowy oraz lojalną społeczność użytkowników, która pomoże ci w problemach związanych z potokiem orkiestracji danych. 

Jeśli dopiero zaczynasz z Pythonem, możesz uznać budowanie i uruchamianie potoków za trudne. Jednak dokumentacja i przewodniki pomogą szybko wystartować. Logi dają ograniczone informacje, a pulpit to tylko narzędzie do wizualizacji DAG-ów i zależności.

Abid Ali AwanAuthor

Pierwsze kroki z Luigi

Tworzenie potoku danych w Luigi wymaga zrozumienia programowania obiektowego. Zacznijmy od instalacji pakietu Luigi dla Pythona. 

$ pip install luigi

Aby opracować prosty potok ETL w Luigi, utworzymy połączone zadania. Zamiast definiować funkcje Pythona jako zadania, utworzymy klasę Pythona dla każdego kroku w potoku: FetchData, ProcessData i GenerateReport. Każda klasa będzie miała trzy funkcje: requires(), output() i run()

Funkcje requires() i output() połączą zadania, a funkcja run() wykona kod przetwarzający. Na końcu zbudujemy potok, używając ostatniego zadania w potoku. 

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)

Uruchom powyższy kod w Jupyter Notebook lub utwórz plik Pythona i uruchom go w terminalu. 

Luigi Execution Summary

Podobnie jak w Luigi, możesz też nauczyć się, jak zbudować potok ETL w Apache Airflow. Ten poradnik omawia podstawy ekstrakcji, transformacji i ładowania danych w Apache Airflow.

Central Planner Luigi

Musimy zainicjować centralny planer Luigi, aby harmonogramować przebiegi potoków lub wyzwalać je zdarzeniami.

Uruchom scheduler, wpisując w terminalu poniższe polecenie.

$ 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

Aby uruchomić potok, otwórz nowy terminal i wpisz poniższe polecenie. Polecenie Luigi wymaga nazwy pliku Pythona oraz ostatniego zadania, które chcemy wykonać. W tym przypadku nazwa pliku to luigi_pipe.py, a naszym ostatnim zadaniem Luigi jest GenerateReport.

$ python -m luigi --module luigi_pipe GenerateReport

Jeśli chcesz zwizualizować przebieg potoku i status zadań, po prostu przejdź do http://localhost:8082 w przeglądarce.

Luigi Central Planner webUI

Interfejs webowy Luigi Central Planner.

To już wszystko w naszym przeglądzie 5 najlepszych alternatyw dla Airflow! Jeśli chcesz zagłębić się w któryś z przykładów przedstawionych w tym artykule, rozważ te zasoby:

Podsumowanie

W tym poradniku omówiliśmy najlepsze, otwartoźródłowe i darmowe alternatywy dla Airflow. Poznaliśmy każde narzędzie do orkiestracji danych oraz zbudowaliśmy i uruchomiliśmy prosty potok ETL. Zobaczenie przykładów kodu pomoże ci zdecydować, które rozwiązanie najlepiej pasuje do twojego przypadku użycia.

Jeśli jesteś początkujący, sugeruję zacząć od Prefect lub Mage AI — są przyjazne dla użytkownika i mają prostą konfigurację. Jeśli jednak szukasz bardziej zaawansowanych narzędzi, które trzymają się praktyk inżynierii oprogramowania, polecam Dagster, Kedro i Luigi.

Po lekturze tego artykułu naturalnym kolejnym krokiem w twojej ścieżce inżynierii danych jest zdobycie certyfikacji, np. Data Engineer in Python od DataCamp, aby poznać inne narzędzia i zbudować kompleksowy potok danych gotowy do wdrożenia na produkcję.

Tematy
Inżynieria danych
Data Science

Poznaj data engineering z tymi kursami!

course

Wprowadzenie do inżynierii danych

4 godz.
129.7K
Poznaj świat inżynierii danych w tym krótkim kursie, obejmującym narzędzia i zagadnienia, takie jak ETL i cloud computing.
Zobacz szczegółyRight Arrow
Rozpocznij Kurs
Zobacz więcejRight Arrow