Перейти к основному контенту

Топ‑5 альтернатив Airflow для оркестрации данных (с примерами кода)

Изучите пять альтернатив Airflow для оркестрации данных с примерами кода по созданию, запуску и визуализации простого ETL‑конвейера.
Обновлено 31 авг. 2026 г.  · 13 мин читать

Изучить с помощью AI

ChatGPTClaudePerplexity

Choose an Airflow Alternatives meme template

Изображение автора.

Apache Airflow — популярный инструмент с открытым исходным кодом для оркестрации данных, предназначенный для создания, планирования и мониторинга конвейеров данных. Он включает панель, помогающую управлять состоянием рабочих процессов, что делает его отличным решением для большинства задач.

Однако Airflow не хватает ряда важных функций, критичных для современных сложных сценариев оркестрации данных.

В этом руководстве мы рассмотрим пять альтернатив Airflow с расширенными возможностями, которые закрывают некоторые его ограничения. Кроме того, мы создадим простой ETL‑конвейер в каждом инструменте, запустим его и визуализируем на их панелях.

Почему стоит выбрать альтернативу Airflow? 

Airflow — мощный инструмент для разных рабочих процессов с данными, но у него есть ограничения, из‑за которых компании могут искать альтернативы. 

Вот несколько причин, почему имеет смысл рассмотреть другой вариант:

  1. Крутая кривая обучения: осваивать Airflow непросто, особенно новичкам в инструментах управления рабочими процессами.
  2. Поддержка: требует значительных усилий по администрированию, особенно при масштабных развёртываниях.
  3. Недостаточная документация: пользователи отмечают проблемы в документации, что усложняет устранение неполадок и изучение новых функций. 
  4. Ресурсоёмкость: для эффективной работы Airflow нужны существенные вычислительные ресурсы и память.
  5. Ограниченная гибкость для тех, кто не использует Python: философия «воркфлоу как код» сильно опирается на Python, что исключает доменных экспертов без навыков программирования.
  6. Масштабируемость: некоторые пользователи сталкиваются со сложностями при масштабировании Airflow для больших рабочих процессов.
  7. Ограниченная обработка в реальном времени: Airflow в первую очередь рассчитан на пакетную обработку, а не на стриминг в реальном времени.

Прежде чем перейти к коду других инструментов оркестрации, полезно научиться писать конвейер данных на Apache Airflow по руководству Getting Started with Apache Airflow — так вы сможете корректно сравнивать альтернативы.

Если вы совсем новичок в Airflow, пройдите краткий курс Introduction to Airflow in Python, чтобы освоить основы построения и планирования конвейеров данных.

5 лучших альтернатив Airflow для оркестрации данных

Теперь опишем пять лучших альтернатив Airflow и покажем, как ими пользоваться на практических примерах кода.

1. Prefect

Prefect — это инструмент оркестрации рабочих процессов на Python с открытым исходным кодом, созданный для современных инженеров по данным и ML. Он предлагает простой API для быстрого создания конвейеров и управления ими через интерактивную панель. 

Prefect предлагает гибридную модель исполнения: вы можете развернуть воркфлоу в облаке и запускать его там или использовать локальный репозиторий.

По сравнению с Airflow, Prefect предоставляет продвинутые возможности: автоматические зависимости задач, триггеры на события, встроенные уведомления, инфраструктуру под конкретные воркфлоу и обмен данными между задачами. Всё это делает его мощным решением для эффективного управления сложными рабочими процессами.

Prefect прост в использовании и при этом очень функционален. На запуск примерного кода ушло буквально 5 минут. Особенно понравился дизайн интерфейса панели: можно настраивать уведомления, перезапускать конвейеры, управлять и мониторить всё через Dashboard.

Abid Ali AwanAuthor

Прочитайте Airflow vs Prefect: как выбрать инструмент для своего рабочего процесса с данными, чтобы узнать подробности сравнения этих инструментов оркестрации. 

Начало работы с Prefect

Начнём проект Prefect с установки Python‑пакета. Выполните в терминале команду ниже.

$ pip install -U prefect

Затем создадим скрипт Python с именем prefect_etl.py и напишем следующий код.

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

Код выше определяет функции‑задачи extract_data(), transform_data() и load_data() и последовательно выполняет их в функции потока etl(). Эти функции создаются с помощью декораторов Prefect для Python. 

Кратко: мы создаём pandas DataFrame, трансформируем его и выводим итог с помощью print. Это простой способ смоделировать ETL‑конвейер.

Чтобы выполнить воркфлоу, просто запустите скрипт Python командой ниже.

$ python prefect_etl.py 

Как видно, запуск нашего воркфлоу успешно завершился.

Prefect flow run logs

Журналы запуска потока Prefect.

Развёртывание потока

Теперь развернём наш воркфлоу, чтобы запускать его по расписанию или по событию. Развёртывание также позволяет централизованно мониторить и управлять несколькими воркфлоу.

Для развёртывания используем Prefect CLI. Функции deploy нужны имя Python‑файла, имя функции потока в этом файле и имя развёртывания. В нашем примере развёртывание называется «simple_etl».

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

После выполнения команды в терминале вы можете получить сообщение, что нет пула рабочих для запуска развёртывания. Чтобы создать пул, используйте команду ниже.

$ prefect worker start --pool 'datacamp'

Теперь, когда пул создан, откройте новое окно терминала и запустите развёртывание. Команда prefect deployment run принимает аргумент вида «<имя-функции-потока>/<имя-развёртывания>», как показано ниже.

$ prefect deployment run 'etl/simple_etl

В результате вы получите сообщение о запуске воркфлоу. Как правило, создаваемый запуск получает случайное имя, в моём случае — 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>

Чтобы увидеть полный лог, вернитесь в окно терминала, где запускали пул рабочих.

Prefect flow run summary

Сводка запуска потока Prefect.

Чтобы визуализировать запуск потока и управлять другими воркфлоу, запустите веб‑сервер Prefect.

$ prefect server start 

После выполнения команды вы будете перенаправлены на панель Prefect. Либо перейдите по адресу http://127.0.0.1:4200 в браузере.

Prefect web server UI

Интерфейс веб‑сервера Prefect

Панель позволяет перезапускать воркфлоу, просматривать логи, проверять пулы, настраивать уведомления и выбирать другие расширенные опции. Это полноценное решение для современных задач оркестрации данных.

Чтобы узнать, как строить и запускать ML‑конвейеры в Prefect, воспользуйтесь руководством Using Prefect for Machine Learning Workflows.

2. Dagster

Dagster — это фреймворк с открытым исходным кодом, предназначенный для инженеров по данным для определения, планирования и мониторинга конвейеров. Он хорошо масштабируется и способствует сотрудничеству разных дата‑команд. 

Dagster позволяет определять «активы данных» как функции Python с помощью декораторов. После определения активов их можно запускать по расписанию или по событиям.

По сравнению с Airflow, Dagster позволяет локально разрабатывать, тестировать и ревьюить конвейеры, использует подход на основе активов и «родной» для облака и контейнеров.

Вместо мышления «шагами и потоками» пришлось перестроиться и строить конвейер из активов данных. В остальном создание и запуск простого ETL оказалось довольно простым. Веб‑сервер минималистичный, но даёт всю информацию для мониторинга активов, запусков и развёртываний.

Abid Ali AwanAuthor

Начало работы с Dagster

Мы создадим простой ETL‑конвейер, выполним его и визуализируем с помощью веб‑сервера Dagster. Подобно панели Prefect, веб‑сервер Dagster предоставляет централизованные средства для мониторинга нескольких воркфлоу и планирования запусков и активов.

Начнём с установки Python‑пакета.

$ pip install dagster -q

Затем создадим три функции Python для извлечения, преобразования и загрузки данных. Эти функции называются create_dirty_data(), clean_data() и load_cleaned_data(). С помощью декоратора @asset объявим их как активы данных в Dagster.

Далее создадим задание по активам (переменная job) на основе всех активов (переменная all_assets), а затем определим объект активов (переменная defs). 

Определение активов можно опустить, но оно важно, если вы хотите планировать запуски, запускать несколько заданий и настраивать сенсоры.

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)

Код можно запустить в Jupyter Notebook или сохранить в файл Python и выполнить. 

В результате выполнения мы получим полный лог запуска конвейера. 

Dagster execution summary

Веб‑сервер Dagster

Чтобы визуализировать активы и запуски заданий, необходимо установить и запустить веб‑сервер Dagster. Он позволяет запускать задания, материализовывать отдельные активы и мониторить несколько заданий одновременно.

$ pip install dagster-webserver

Чтобы запустить сервер Dagster, используем CLI Dagster и передадим путь к Python‑файлу. В данном случае файл называется dagster_pipe.py.

$ dagster dev -f dagster_pipe.py  

Команда автоматически откроет веб‑сервер в браузере. Либо перейдите напрямую по адресу http://127.0.0.1:3000 в браузере.

Dagster Web server

Интерфейс веб‑сервера Dagster.

Пока мы развернули только задание. Чтобы запустить воркфлоу, перейдите на вкладку «Runs» и нажмите кнопку «Launch a new run». 

Запуск должен успешно завершиться. Чтобы посмотреть логи, кликните по ID нужного запуска.

Dagster runs detailed view

Журналы запусков Dagster.

3. Mage AI

Mage AI — гибридный фреймворк оркестрации данных с открытым исходным кодом. Гибридный — значит, он сочетает гибкость Jupyter Notebook и управляемость модульного кода. 

Любой пользователь, даже с минимальными знаниями Python, может собирать, запускать и мониторить конвейеры данных. Вместо того чтобы писать и запускать файл Python напрямую, вы создаёте проект Mage AI и работаете с ним в панели, где строите, запускаете и управляете конвейерами.

По сравнению с Airflow, Mage AI предлагает удобный интерфейс и простоту использования — отличный выбор для новичков в инженерии данных. Он спроектирован с учётом масштабируемости и способен эффективно обрабатывать большие объёмы данных и сложные структуры конвейеров.

Ощущения были непривычными, потому что подход совсем иной. Нужно установить и запустить веб‑интерфейс Mage AI. По идее всё просто, но мне показалось не так легко собрать и запустить ETL‑конвейер. С другой стороны, понятно, почему такой подход может нравиться новичкам: по сути, это перетаскивание модулей и нажатие кнопок.

Abid Ali AwanAuthor

Начало работы с Mage AI

Запуск Mage AI довольно прост. Нужно лишь установить Python‑пакет Mage AI.

$ pip install mage-ai

И стартовать проект Mage AI. 

$ mage start mage_ai_etl 

Эта команда запустит веб‑сервер. Как уже говорилось, редактирование кода, запуск и мониторинг заданий выполняются через интерфейс Mage AI.

Mage AI UI

Интерфейс Mage AI.

Нажмите «+ New pipeline», чтобы создать свой первый ETL‑конвейер. Я назвал его «simple_etl».

Creating the new pipeline in Mage AI

Создание нового конвейера в Mage AI.

Затем интерфейс предложит добавить модуль для начала работы. Выберите модуль «Data Loader» и вставьте следующий код на Python. 

Здесь мы объявляем функцию create_sample_csv() — это первый шаг конвейера. Используем декоратор Mage AI  @data_loader. Также определим функцию test_output(), которая проверяет наличие результата. Это помогает управлять зависимостями задач.

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

Создание блока загрузчика данных в Mage AI.

Аналогично создайте модуль «Transformer» и добавьте функцию clean_data(), как показано ниже. 

Функцию test() можно опустить; достаточно добавить основную функцию трансформации — 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'

Точно так же создайте модуль «Data Exporter» и добавьте код ниже. В нём объявлена функция загрузки данных export_data_to_csv(), которая сохраняет преобразованные данные в 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")

Чтобы запустить конвейер, перейдите на вкладку «Trigger» и нажмите «Run@once».

Running the pipeline in Mage AI

Запуск конвейера в Mage AI.

Чтобы посмотреть логи запуска, откройте вкладку «Runs» и нажмите «Logs» у последнего запуска конвейера.

Mage AI flow run logs

Журналы запуска потока в Mage AI.

4. Kedro

Kedro — ещё один популярный фреймворк оркестрации данных с открытым исходным кодом, немного отличающийся от остальных. Он создан для ML‑инженеров и перенимает многие практики разработки ПО, применяя их к ML‑проектам.

Kedro спроектирован максимально модульно: даже для экспорта набора данных нужно создать каталог данных, указав расположение и тип данных, что обеспечивает стандартизированное и эффективное управление во всём конвейере.

Чтобы понять место Kedro в ML‑экосистеме, изучите инструменты MLOps в статье 25 Top MLOps Tools You Need to Know in 2024.

По сравнению с Airflow, API Kedro проще для сборки конвейеров данных. Он больше ориентирован на ML‑инженерию и предлагает категоризацию и версионирование данных.

Код писать довольно просто, но сложности начинаются при выполнении конвейера. Нужно создать каталог данных, зарегистрировать конвейер и разобраться со структурой проекта Kedro. По сравнению с Dagster и Prefect это, на мой взгляд, сложнее. Однако понятна задумка: сделать конвейер надёжным и безошибочным.

Abid Ali AwanAuthor

Начало работы с Kedro

Построение конвейера данных в Kedro — это другой подход. Фреймворк модульный, и важно понимать структуру проекта и шаги для успешного выполнения воркфлоу. 

Начните с установки Python‑пакета Kedro. 

$ pip install kedro

Инициализируйте проект Kedro. 

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

Перейдите в каталог проекта. 

$ cd kedro-etl  

Создайте папку внутри каталога pipelines с именем data_processing.

$ mkdir -p src/kedro_etl/pipelines/data_processing  

Создайте файл Python с именем kedro_pipe.py и откройте его в любимой IDE, например, в Visual Studio Code.

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

Скрипт должен содержать функции извлечения, преобразования и загрузки — это узлы конвейера. В нашем случае это функции create_sample_data(), clean_data(), и load_and_process_data().

Затем соединим эти узлы с помощью класса Kedro Pipeline внутри функции create_pipeline(). В этой функции мы определяем узлы, и у каждого узла есть inputs, outputs и 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",
            ),
        ]
    )

Если запустить конвейер без создания каталога данных, экспорт не произойдёт. Поэтому откройте файл conf/base/catalog.yml и укажите конфигурацию наборов данных.

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

Также нужно добавить созданный Python‑файл в реестр конвейеров. Для этого откройте файл src/simple_etl/pipeline_registry.py и добавьте код ниже. 

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

Запустите конвейер и просматривайте живые логи в терминале командой ниже.

$ kedro run

Logs of Kedro pipeline run

Журналы запуска конвейера Kedro.

После выполнения конвейера файлы будут сохранены в формате CSV по путям, указанным в каталоге данных.

Output files of Kedro pipeline run

Выходные файлы после запуска конвейера Kedro.

Если возникают проблемы с запуском, попробуйте установить Kedro со всеми расширениями. 

$ pip install "kedro[all]"

Визуализация Kedro

Мы можем визуализировать и делиться своими конвейерами, установив инструмент kedro-viz

$ pip install kedro-viz

После выполнения следующей команды мы увидим все конвейеры данных и узлы. Также доступны трекинг экспериментов и совместное использование визуализаций.

$ kedro viz run

Kedro Visualization

Визуализация конвейера Kedro.

5. Luigi

Luigi — Python‑фреймворк с открытым исходным кодом, разработанный Spotify. Он отлично подходит для управления длительными пакетными процессами и сложными конвейерами данных. Преуспевает в разрешении зависимостей, управлении воркфлоу, визуализации и восстановлении после сбоев, что делает его мощным инструментом оркестрации.

По сравнению с Airflow, у Luigi минималистичный API, календарное планирование и преданная пользовательская база, готовая помочь с любыми вопросами по конвейерам оркестрации. 

Если вы новичок в Python, сборка и запуск конвейеров могут показаться сложными. Однако документация и руководства помогают быстро стартовать. Логи дают ограниченную информацию, а панель — это в основном визуализация DAG и зависимостей.

Abid Ali AwanAuthor

Начало работы с Luigi

Создание конвейера данных на Luigi требует понимания объектно‑ориентированного программирования. Начнём с установки Python‑пакета Luigi. 

$ pip install luigi

Чтобы разработать простой ETL‑конвейер в Luigi, создадим взаимосвязанные задачи. Вместо функций‑задач мы создадим класс Python для каждого шага конвейера — FetchData, ProcessData и GenerateReport. Каждый класс содержит три метода: requires(), output() и run()

Методы requires() и output() связывают задачи, а метод run() выполняет код обработки. В конце мы построим конвейер, вызвав последнюю задачу в цепочке. 

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)

Запустите код в Jupyter Notebook или сохраните в файл и выполните в терминале. 

Luigi Execution Summary

Аналогично Luigi, вы также можете узнать, как построить ETL‑конвейер на Apache Airflow. В руководстве рассматриваются основы извлечения, преобразования и загрузки данных с Apache Airflow.

Центральный планировщик Luigi

Нужно инициализировать центральный планировщик Luigi, чтобы планировать запуски конвейеров или триггерить их по событиям.

Запустите планировщик, выполнив в терминале команду ниже.

$ 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

Чтобы запустить конвейер, откройте новый терминал и выполните команду. Команде Luigi нужны имя Python‑файла и последняя задача для выполнения. В нашем случае файл — luigi_pipe.py, а последняя задача — GenerateReport.

$ python -m luigi --module luigi_pipe GenerateReport

Чтобы визуализировать запуск и статус задач, просто откройте http://localhost:8082 в браузере.

Luigi Central Planner webUI

Веб‑интерфейс Luigi Central Planner.

На этом завершаем обзор 5 лучших альтернатив Airflow. Если хотите глубже изучить примеры из статьи, обратите внимание на ресурсы ниже:

Итоговые мысли

В этом руководстве мы обсудили лучшие бесплатные open‑source альтернативы Airflow. Мы также познакомились с каждым инструментом оркестрации данных, собрали и запустили простой ETL‑конвейер. Примеры кода помогут понять, что лучше подойдёт вашему кейсу.

Если вы начинаете, рекомендуем Prefect или Mage AI — они дружелюбны к пользователю и просты в настройке. Если же нужны более продвинутые инструменты с практиками инженерии ПО, присмотритесь к Dagster, Kedro и Luigi.

После этой статьи логичный следующий шаг в вашем пути инженера данных — получить сертификацию, например, Data Engineer in Python от DataCamp, чтобы изучить другие инструменты и собрать сквозной конвейер данных для продакшна.

Темы
Дата-инжиниринг
Data Science

Узнайте больше об инженерии данных на этих курсах!

Course

Введение в дата-инжиниринг

4 ч
129.7K
Изучите мир data engineering в этом коротком курсе: инструменты и темы, такие как ETL и cloud computing.
ПодробнееRight Arrow
Начать Курс
Смотрите большеRight Arrow