course

Bild av författaren.
Apache Airflow är ett populärt öppet verktyg för dataorkestrering, utformat för att bygga, schemalägga och övervaka datapipelines. Det har en instrumentpanel som hjälper till att hantera arbetsflödenas status, vilket gör det till ett utmärkt verktyg för de flesta behov kring arbetsflöden.
Airflow saknar dock några viktiga funktioner som kan vara avgörande för komplexa, moderna krav på dataorkestrering.
I den här handledningen utforskar vi fem alternativ till Airflow som erbjuder utökade möjligheter och adresserar vissa av dess begränsningar. Dessutom lär vi oss att bygga en enkel ETL‑pipeline med varje verktyg, köra den och visualisera den i respektive instrumentpanel.
Varför välja ett alternativ till Airflow?
Airflow är kraftfullt för många dataarbetsflöden, men har flera begränsningar som kan få företag att överväga alternativ.
Här är några skäl att välja ett alternativ:
- Brant inlärningskurva: Airflow kan vara svårt att lära sig, särskilt för dem som är nya inom verktyg för arbetsflödeshantering.
- Underhåll: Det kräver mycket underhåll, särskilt vid storskaliga driftsättningar.
- Otillräcklig dokumentation: Användare har rapporterat flera dokumentationsproblem som försvårar felsökning eller inlärning av nya funktioner.
- Resurskrävande: Airflow kan vara resursintensivt och kräver betydande beräkningskraft och minne för att köras effektivt.
- Begränsad flexibilitet för icke‑Pythonanvändare: Filosofin ”workflow‑as‑code” bygger tungt på Python, vilket kan utestänga domänexperter som inte är så bevandrade i programmering.
- Skalbarhet: Vissa användare rapporterar svårigheter att skala Airflow för stora arbetsflöden.
- Begränsad realtidsbearbetning: Airflow är främst utformat för batchbearbetning, inte strömmar av data i realtid.
Innan vi dyker in i kodningen i andra verktyg för dataorkestrering är det viktigt att lära sig hur man skriver datapipelinen med Apache Airflow genom att följa handledningen Kom igång med Apache Airflow, så att du rättvist kan jämföra alternativen.
Om du är helt ny på Airflow kan du överväga att gå den korta kursen Introduction to Airflow in Python för att lära dig grunderna i att bygga och schemalägga datapipelines.
5 bästa alternativen till Airflow för dataorkestrering
Låt oss nu beskriva de 5 främsta alternativen till Airflow och visa hur du använder dem med praktiska kodexempel.
1. Prefect
Prefect är ett öppet Python‑verktyg för orkestrering av arbetsflöden, byggt för moderna data‑ och ML‑ingenjörer. Det erbjuder ett enkelt API som låter dig snabbt bygga en datapipeline och hantera den via en interaktiv instrumentpanel.
Perfect erbjuder en hybrid körningsmodell, vilket innebär att du kan driftsätta arbetsflödet i molnet och köra det där, eller använda det lokala arkivet.
Jämfört med Airflow kommer Prefect med avancerade funktioner som automatiserade uppgiftsberoenden, händelsebaserade triggrar, inbyggda aviseringar, arbetsflödesspecifik infrastruktur och datadelning mellan uppgifter. Dessa möjligheter gör det till en kraftfull lösning för att hantera komplexa arbetsflöden effektivt och ändamålsenligt.
Prefect är enkelt och har kraftfulla funktioner. Det tog mig i princip 5 minuter att köra exempelkoden. Jag gillar särskilt hur instrumentpanelens UI är utformad, hur du kan ställa in aviseringar, köra om pipelines samt hantera och övervaka allt via instrumentpanelen.
Abid Ali Awan, Author
Läs bloggen Airflow vs Prefect: Deciding Which is Right For Your Data Workflow för en detaljerad jämförelse mellan dessa två verktyg för dataorkestrering.
Kom igång med Prefect
Vi börjar vårt Prefect‑projekt med att installera Python‑paketet. Kör följande kommando i en terminal.
$ pip install -U prefect
Därefter skapar vi ett Python‑skript med namnet prefect_etl.py och skriver följande 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()
Koden ovan definierar uppgiftsfunktionerna extract_data(), transform_data(), och load_data() och exekverar dem i serie i en flödesfunktion med namnet etl(). Dessa funktioner skapas med Prefects Python‑dekorerare.
Kort sagt skapar vi en pandas DataFrame, transformerar den och visar sedan slutresultatet med print. Detta är ett enkelt sätt att simulera en ETL‑pipeline.
För att köra arbetsflödet, kör bara Python‑skriptet med följande kommando.
$ python prefect_etl.py
Som vi kan se slutfördes vårt arbetsflöde utan problem.

Loggar för Prefect‑flödeskörning.
Driftsätta flödet
Vi ska nu driftsätta vårt arbetsflöde så att vi kan köra det enligt schema eller trigga det baserat på en händelse. Att driftsätta flödet gör det också möjligt att övervaka och hantera flera arbetsflöden på ett centraliserat sätt.
För att driftsätta flödet använder vi Prefect CLI. Funktionen deploy kräver Python‑filens namn, namnet på flödesfunktionen i filen och driftsättningens namn. I det här fallet kallar vi driftsättningen ”simple_etl”.
$ prefect deploy prefect_etl.py:etl -n 'simple_etl'
Efter att ha kört ovanstående skript i terminalen kan vi få meddelandet att vi inte har någon worker‑pool för att köra driftsättningen. För att skapa worker‑poolen, använd följande kommando.
$ prefect worker start --pool 'datacamp'
Nu när vi har en worker‑pool öppnar vi ett nytt terminalfönster och kör driftsättningen. Kommandot prefect deployment run kräver ”<flow-function-name>/<deployment-name>” som argument, som visas i kommandot nedan.
$ prefect deployment run 'etl/simple_etl
Som resultat av driftsättningen får du ett meddelande om att arbetsflödet körs. Vanligtvis tilldelas flödeskörningen som skapas ett slumpmässigt namn, i mitt fall 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>
För att se hela loggen, växla tillbaka till terminalfönstret där du startade worker‑poolen.

Sammanfattning av Prefect‑flödeskörning.
Du måste starta Prefects webbserver för att visualisera flödeskörningen på ett mer användarvänligt sätt och hantera andra arbetsflöden.
$ prefect server start
Efter att ha kört kommandot ovan bör du omdirigeras till Prefects instrumentpanel. Alternativt kan du gå direkt till adressen http://127.0.0.1:4200 i din webbläsare.

Prefects webbserver‑UI
Instrumentpanelen låter dig köra om arbetsflödet, visa loggar, kontrollera work‑pooler, ställa in aviseringar och välja andra avancerade alternativ. Det är en komplett lösning för moderna behov av dataorkestrering.
För att lära dig hur man bygger och kör maskininlärningspipeliner med Prefect kan du följa handledningen Using Prefect for Machine Learning Workflows.
2. Dagster
Dasgter är ett öppet ramverk utformat för dataingenjörer att definiera, schemalägga och övervaka datapipelines. Det är mycket skalbart och underlättar samarbete mellan olika datateam.
Dagster gör det möjligt för användare att definiera sina dataresurser som Python‑funktioner med dekorerare. När dessa tillgångar är definierade kan de köras sömlöst via schemaläggning eller händelsebaserade triggrar.
Jämfört med Airflow låter Dagster oss utveckla, testa och granska pipelinen lokalt, erbjuder ett tillgångsbaserat angreppssätt för orkestrering och är moln‑ och container‑native.
I stället för att tänka på arbetsflöden i termer av steg och flöden behövde jag ändra mitt tänk och bygga en pipeline med dataresurser. Förutom det var det ganska enkelt att bygga och köra en enkel ETL‑pipeline. Webbservern är också relativt minimalistisk men ger all information för att övervaka tillgångar, körningar och driftsättningar.
Abid Ali Awan, Author
Kom igång med Dagster
Vi skapar en enkel ETL‑pipeline, kör den och visualiserar den via Dagsters webbserver. Precis som i Prefects instrumentpanel erbjuder Dagsters webbserver centraliserade sätt att övervaka flera arbetsflöden samt schemalägga körningar och tillgångar.
Vi börjar med att installera Python‑paketet.
$ pip install dagster -q
Sedan skapar vi tre Python‑funktioner för att extrahera, transformera och ladda data. Dessa funktioner heter create_dirty_data(), clean_data(), och load_cleaned_data() i koden. Med dekoreraren @asset deklarerar vi funktionerna som dataresurser i Dagster.
Därefter skapar vi asset‑jobbet (variabeln job) med alla tillgångar (variabeln all_assets) och skapar sedan asset‑definitionen (variabeln defs).
Du kan hoppa över delen med asset‑definition, men den blir viktig om du vill schemalägga körningar, köra flera jobb och ställa in sensorer.
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)
Du kan köra koden ovan i en Jupyter Notebook eller skapa Python‑filen och köra den.
Som resultat av körningen får vi en komplett logg över arbetsflödet.

Dagsters webbserver
För att visualisera tillgångar och jobbkörningar måste vi installera och köra Dagsters webbserver. Webbservern låter dig köra jobb, materialisera enskilda tillgångar och övervaka flera jobb samtidigt.
$ pip install dagster-webserver
För att starta Dagster‑servern använder vi Daster CLI och anger platsen för Python‑filen. I det här fallet döpte jag filen till dagster_pipe.py.
$ dagster dev -f dagster_pipe.py
Kommandot ovan startar webbservern automatiskt i din webbläsare. Alternativt kan du gå direkt till adressen http://127.0.0.1:3000 i din webbläsare.

Dagsters webbserver‑UI.
Hittills har vi bara driftsatt jobbet. För att köra arbetsflödet går du till fliken ”Runs” och klickar på knappen ”Launch a new run”.
Körningen ska slutföras utan problem! För att se loggarna, klicka på ID:t för körningen du är intresserad av.

Dagster – körloggar.
3. Mage AI
Mage AI är ett öppet hybridramverk för dataorkestrering. Hybrid betyder att du får flexibiliteten från en Jupyter Notebook och kontrollen från modulär kod.
Vem som helst, även med begränsade Python‑kunskaper, kan bygga, köra och övervaka datapipelines. I stället för att skriva och köra en Python‑fil direkt skapar du ett Mage AI‑projekt och startar det i instrumentpanelen, där du kan bygga, köra och hantera dina datapipelines.
Jämfört med Airflow erbjuder Mage AI ett användarvänligt gränssnitt och enkel användning, vilket gör det till ett utmärkt val för dem som är nya inom dataengineering. Det är designat med skalbarhet i åtanke och klarar stora datamängder och komplexa pipelinestrukturer effektivt.
Det kändes ovant eftersom det var helt annorlunda än vad jag är van vid. Jag behövde installera och starta Mage AIs webb‑UI. Det skulle vara enkelt, men jag tyckte att det var svårt att bygga och köra ETL‑pipelinen. Å andra sidan förstår jag varför den här unika designen kan vara attraktiv för nybörjare – det är i princip dra‑och‑släpp och att trycka på knappar.
Abid Ali Awan, Author
Kom igång med Mage AI
Att starta Mage AI är ganska enkelt. Vi behöver bara installera Python‑paketet Mage AI.
$ pip install mage-ai
Och starta Mage AI‑projektet.
$ mage start mage_ai_etl
Kommandot ovan startar webbservern. Som nämnts tidigare sker all kodredigering, körning av jobb och övervakning via Mage AIs UI.

Mage AI‑UI.
Klicka på ”+ New pipeline” för att skapa din första ETL‑pipeline. Jag döpte min till ”simple_etl”.

Skapa ny pipeline i Mage AI.
Sedan ber gränssnittet dig att lägga till en modul för att börja koda. Välj modulen ”Data Loader” och skriv följande Python‑kod.
Här deklarerar vi funktionen create_sample_csv(), som är det första steget i vår pipeline. Vi använder Mage AIs @data_loader‑dekorerare. Vi definierar också funktionen test_output() som kontrollerar att utdata finns. Detta hjälper till med hantering av uppgiftsberoenden.
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'

Skapa data loader‑block i Mage AI.
Skapa därefter en annan modul kallad ”Transformer” och lägg till funktionen clean_data() enligt koden nedan.
Du kan bortse från funktionen test(); du behöver bara lägga till huvudfunktionen för transformering, 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'
Skapa på liknande sätt en modul ”Data Exporter” och lägg till följande kod. Koden deklarerar en dataladdningsfunktion, export_data_to_csv(), som sparar den transformerade datan i en CSV‑fil.
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")
För att köra pipelinen, gå till fliken ”Trigger” och klicka på ”Run@once”.

Köra pipelinen i Mage AI.
För att visa körloggarna går du till fliken ”Runs” och klickar på knappen ”Logs” på den nyligen körda pipelinen.

Mage AI – loggar för flödeskörning.
4. Kedro
Kedro är ett annat populärt öppet ramverk för dataorkestrering som skiljer sig något från de andra verktygen. Det skapades för ML‑ingenjörer och lånar många koncept från mjukvaruingenjörskap som tillämpas på maskininlärningsprojekt.
Kedro är utformat för att vara mycket modulärt, vilket innebär att även för att exportera en datamängd måste du skapa en datakatalog som anger plats och typ av data, vilket säkerställer standardiserad och effektiv datahantering genom hela pipelinen.
För att förstå hur Kedro passar in i ML‑ekosystemet kan du utforska olika MLOps‑verktyg genom att läsa artikeln 25 Top MLOps Tools You Need to Know in 2024.
Jämfört med Airflow är Kedros API enklare för att bygga en datapipeline. Det fokuserar mer på maskininlärningsengineering och erbjuder datakategorisering och versionering.
Kodningsdelen är ganska rättfram, men problem uppstår när du vill köra din pipeline. Du måste skapa en datakatalog, registrera pipelinen och förstå Kedros projektstruktur. Jag skulle säga att det är mer utmanande jämfört med Dagster och Prefect. Men jag förstår varför det är designat så här: för att göra din datapipeline tillförlitlig och felfri.
Abid Ali Awan, Author
Kom igång med Kedro
Att bygga en Kedro‑datapipeline är en annan femma. Ramverket är modulärt och du behöver förstå projektstrukturen och de olika stegen som krävs för att köra arbetsflödet framgångsrikt.
Börja med att installera Python‑paketet Kedro.
$ pip install kedro
Initiera Kedro‑projektet.
$ kedro new --name=kedro_etl --tools=none --example=n
Gå till projektkatalogen.
$ cd kedro-etl
Skapa en mapp i mappen pipelines som heter data_processing.
$ mkdir -p src/kedro_etl/pipelines/data_processing
Skapa en Python‑fil som heter kedro_pipe.py och öppna den i din favorit‑IDE, till exempel Visual Studio Code.
$ code src/kedro_etl/pipelines/data_processing/kedro_pipe.py
Python‑skriptet ska innehålla funktionerna för extract, transform och load, som är noder i pipelinen. I det här fallet är det funktionerna create_sample_data(), clean_data(), och load_and_process_data().
Sedan kopplar vi ihop dessa noder med Kedros klassen Pipeline i funktionen create_pipeline(). I pipelinefunktionen definierar vi noder, och varje nod har inputs, outputs och ett nod‑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",
),
]
)
Om vi kör pipelinen utan att skapa datakatalogen exporteras inte vår data. Vi behöver därför gå till filen conf/base/catalog.yml och redigera den genom att ange dataset‑konfigurationen.
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
Vi måste också inkludera vår nyskapade Python‑fil i pipeline‑registret. För att göra det, gå till Python‑filen src/simple_etl/pipeline_registry.py och lägg in följande 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,
}
Kör pipelinen och visa live‑loggar i terminalen genom att köra följande kommando.
$ kedro run

Loggar från körning av Kedro‑pipeline.
Efter att ha kört pipelinen lagras dina filer i CSV‑format på platsen som definierats i datakatalogen.

Utdatafiler från Kedro‑pipeline.
Om du stöter på problem med att köra pipelinen kan du överväga att installera Kedro med alla tillägg.
$ pip install "kedro[all]"
Kedro‑visualisering
Vi kan visualisera och dela våra pipelines genom att installera verktyget kedro-viz.
$ pip install kedro-viz
Att sedan köra följande kommando låter oss visualisera alla datapipelines och datanoder. Det ger också möjlighet till spårning av experiment och att dela pipeline‑visualiseringen.
$ kedro viz run

Visualisering av Kedro‑pipeline.
5. Luigi
Luigi är ett öppet, Python‑baserat ramverk utvecklat av Spotify som utmärker sig i att hantera långvariga batchprocesser och komplexa datapipelines. Det är bra på beroendeupplösning, arbetsflödeshantering, visualisering och återhämtning vid fel, vilket gör det till ett kraftfullt verktyg för orkestrering av dataarbetsflöden.
Jämfört med Airflow har Luigi ett minimalt API, kalenderschemaläggning och en lojal användarbas som hjälper dig med eventuella problem relaterade till dataorkestreringspipelinen.
Om du är nybörjare i Python kan du tycka att det är svårt att bygga och köra pipelines. Dokumentation och guider kan dock hjälpa dig att komma igång snabbt. Loggarna ger begränsad information och instrumentpanelen är bara ett visualiseringsverktyg för DAG:ar och beroenden.
Abid Ali Awan, Author
Kom igång med Luigi
Att skapa en Luigi‑datapipeline kräver förståelse för objektorienterad programmering. Låt oss börja med att installera Python‑paketet Luigi.
$ pip install luigi
För att utveckla en enkel ETL‑pipeline i Luigi skapar vi sammankopplade uppgifter. I stället för att skapa Python‑funktioner som uppgifter skapar vi en Python‑klass för varje steg i pipelinen, FetchData, ProcessData och GenerateReport. Varje klass har tre funktioner: requires(), output() och run().
Funktionerna requires() och output() kopplar ihop uppgifterna, och funktionen run() exekverar bearbetningskoden. I slutet bygger vi pipelinen med den sista uppgiften i pipelinen.
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)
Kör koden ovan i Jupyter Notebook eller skapa Python‑filen och kör den via terminalen.

På liknande sätt som för Luigi kan du också lära dig att bygga en ETL‑pipeline med Apache Airflow. Handledningen täcker grunderna i att extrahera, transformera och ladda data med Apache Airflow.
Luigis centrala planerare
Vi behöver starta Luigis centrala planerare för att schemalägga pipelinekörningar eller trigga dem med en händelse.
Starta schemaläggaren genom att skriva följande kommando i terminalen.
$ 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
För att köra pipelinen, öppna en ny terminal och skriv följande kommando. Luigi‑kommandot kräver ett Python‑filnamn och den sista uppgiften vi vill köra. I det här fallet heter filen luigi_pipe.py och vår sista Luigi‑uppgift är GenerateReport.
$ python -m luigi --module luigi_pipe GenerateReport
Om du vill visualisera pipelinekörningen och uppgiftsstatusen kan du helt enkelt gå till http://localhost:8082 i din webbläsare.

Luigi Central Planner webUI.
Där avslutar vi vår genomgång av de 5 bästa alternativen till Airflow! Om du vill fördjupa dig i något av exemplen i artikeln finns här några resurser att överväga:
- För källkod och data till Prefect, Dagster och Luigi, se DataLab‑arbetsytan.
- För källkod och data till Mage AI och Kedro, se GitHub‑arkivet.
Avslutande tankar
I den här handledningen har vi gått igenom de främsta öppna och kostnadsfria alternativen till Airflow. Vi har också bekantat oss med varje verktyg för dataorkestrering samt byggt och kört en enkel ETL‑pipeline. Att se kodexempel hjälper dig att avgöra vilket som passar bäst för ditt användningsfall.
Om du är nybörjare rekommenderar jag att börja med Prefect eller Mage AI, eftersom de är användarvänliga och enkla att komma igång med. Om du däremot letar efter mer avancerade verktyg som följer mjukvaruingenjörspraxis, rekommenderar jag att du utforskar Dagster, Kedro och Luigi.
Efter den här artikeln är nästa naturliga steg i din resa inom dataengineering att skaffa en certifiering, som DataCamps Data Engineer in Python, för att lära dig fler verktyg och bygga en end‑to‑end‑datapipeline som du kan driftsätta i produktion.