본문으로 바로가기

데이터 오케스트레이션을 위한 Airflow 대안 Top 5 (코드 예제 포함)

간단한 ETL 파이프라인을 구축, 실행, 시각화하기 위한 코드 예제와 함께 Airflow의 다섯 가지 데이터 오케스트레이션 대안을 살펴봅니다.
업데이트됨 2026년 8월 31일  · 13분 읽다

AI로 탐색하기

ChatGPTClaudePerplexity

Airflow 대안을 고르라는 밈 템플릿

이미지: 작성자 제공.

Apache Airflow는 데이터 파이프라인을 구축, 스케줄링, 모니터링하기 위해 설계된 인기 있는 오픈 소스 데이터 오케스트레이션 도구입니다. 워크플로의 상태를 관리하는 데 도움이 되는 대시보드를 제공하여 대부분의 워크플로 요구 사항에 적합한 도구입니다.

하지만 Airflow에는 현대적이고 복잡한 데이터 오케스트레이션 요구에 중요할 수 있는 일부 핵심 기능이 부족합니다.

이 튜토리얼에서는 Airflow의 한계를 보완하고 기능을 강화한 다섯 가지 대안을 살펴봅니다. 아울러 각 도구로 간단한 ETL 파이프라인을 만들고, 실행하고, 대시보드에서 시각화하는 방법도 배웁니다.

왜 Airflow의 대안을 선택해야 할까요? 

Airflow는 다양한 데이터 워크플로에 강력한 도구이지만, 몇 가지 한계로 인해 대안을 고려해야 할 수도 있습니다. 

다음은 대안을 선택할 만한 이유입니다.

  1. 가파른 학습 곡선: 워크플로 관리 도구를 처음 접하는 경우 Airflow를 익히기 어려울 수 있습니다.
  2. 유지보수 부담: 특히 대규모 환경에서는 상당한 유지보수가 필요합니다.
  3. 불충분한 문서화: 사용자들은 문제 해결이나 신규 기능 학습을 어렵게 만드는 문서상의 여러 이슈를 보고했습니다. 
  4. 리소스 집약적: 효율적으로 실행하려면 상당한 컴퓨팅 자원과 메모리가 필요합니다.
  5. 비(非) Python 사용자의 유연성 제한: 코드 기반 워크플로 철학이 Python에 크게 의존하여, 프로그래밍에 익숙하지 않은 도메인 전문가를 배제할 수 있습니다.
  6. 확장성: 대규모 워크플로로 확장하는 데 어려움을 겪는다는 보고가 있습니다.
  7. 실시간 처리의 한계: Airflow는 주로 배치 처리를 위해 설계되었으며, 실시간 스트리밍에는 적합하지 않습니다.

다른 데이터 오케스트레이션 도구의 코딩 파트로 들어가기 전에, 공정한 비교를 위해 Getting Started with Apache Airflow 튜토리얼을 따라 Airflow로 데이터 파이프라인을 작성하는 방법을 먼저 익히세요.

Airflow가 완전히 처음이라면, 짧은 Introduction to Airflow in Python 과정을 수강하여 데이터 파이프라인을 구축하고 스케줄링하는 기초를 배우는 것도 좋습니다.

데이터 오케스트레이션을 위한 Airflow 대안 5가지

이제 Airflow의 상위 5가지 대안을 소개하고, 실용적인 코드 예제와 함께 사용하는 방법을 살펴보겠습니다.

1. Prefect

Prefect는 현대의 데이터 및 머신 러닝 엔지니어를 위해 만들어진 오픈 소스 Python 워크플로 오케스트레이션 도구입니다. 간단한 API로 데이터 파이프라인을 빠르게 구축하고, 인터랙티브 대시보드에서 관리할 수 있습니다. 

Prefect는 하이브리드 실행 모델을 제공하므로, 워크플로를 클라우드에 배포해 거기서 실행하거나 로컬 리포지토리를 사용할 수 있습니다.

Airflow와 비교하면, Prefect는 자동화된 태스크 의존성, 이벤트 기반 트리거, 내장 알림, 워크플로별 인프라, 태스크 간 데이터 공유 등 고급 기능을 제공합니다. 이러한 기능으로 복잡한 워크플로를 효율적이고 효과적으로 관리할 수 있습니다.

Prefect는 단순하면서도 강력한 기능을 갖추고 있습니다. 예제 코드를 실행하는 데 5분도 채 걸리지 않았습니다. 특히 대시보드 UI 디자인, 알림 설정, 파이프라인 재실행, 대시보드에서 모든 것을 관리·모니터링할 수 있는 점이 마음에 듭니다.

Abid Ali AwanAuthor

자세한 비교는 Airflow vs Prefect: 내 데이터 워크플로에 맞는 선택은? 블로그에서 확인하세요. 

Prefect 시작하기

Python 패키지를 설치하는 것부터 시작합니다. 터미널에서 아래 명령을 실행하세요.

$ pip install -U prefect

그다음 prefect_etl.py라는 Python 스크립트를 생성하고 다음 코드를 작성합니다.

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 플로 실행 로그

Prefect 플로 실행 로그.

플로 배포

이제 워크플로를 배포해 일정에 따라 실행하거나 이벤트 기반으로 트리거할 수 있도록 하겠습니다. 플로를 배포하면 여러 워크플로를 중앙에서 모니터링하고 관리할 수도 있습니다.

플로를 배포하기 위해 Prefect CLI를 사용합니다. deploy 기능에는 Python 파일명, 그 파일 내 플로 함수명, 배포명이 필요합니다. 여기서는 배포명을 “simple_etl”로 지정합니다.

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

터미널에서 위 스크립트를 실행한 뒤, 배포를 실행할 워커 풀(worker pool)이 없다는 메시지가 표시될 수 있습니다. 워커 풀을 만들려면 다음 명령을 사용하세요.

$ prefect worker start --pool 'datacamp'

워커 풀이 준비되었으니, 다른 터미널 창을 열어 배포를 실행합니다. prefect deployment run 명령에는 아래와 같이 “<flow-function-name>/<deployment-name>” 인자가 필요합니다.

$ 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 플로 실행 요약

Prefect 플로 실행 요약.

플로 실행을 더 직관적으로 시각화하고 다른 워크플로를 관리하려면 Prefect 웹 서버를 시작해야 합니다.

$ prefect server start 

위 명령을 실행하면 Prefect 대시보드로 리다이렉트됩니다. 또는 브라우저에서 http://127.0.0.1:4200 주소로 바로 이동할 수도 있습니다.

Prefect 웹 서버 UI

Prefect 웹 서버 UI

대시보드에서는 워크플로를 재실행하고, 로그를 조회하며, 워크 풀을 확인하고, 알림을 설정하는 등 다양한 고급 옵션을 사용할 수 있습니다. 현대적 데이터 오케스트레이션 요구에 대한 완전한 솔루션입니다.

Prefect로 머신 러닝 파이프라인을 구축하고 실행하는 방법은 Using Prefect for Machine Learning Workflows 튜토리얼을 참고하세요.

2. Dagster

Dasgter는 데이터 엔지니어가 데이터 파이프라인을 정의, 스케줄링, 모니터링하도록 설계된 오픈 소스 프레임워크입니다. 고도로 확장 가능하며 여러 데이터 팀 간 협업을 촉진합니다. 

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의 데이터 애셋으로 선언합니다.

다음으로 모든 애셋(all_assets 변수)을 대상으로 애셋 잡(job 변수)을 만들고, 애셋 정의(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)

위 코드는 주피터 노트북에서 실행하거나 Python 파일로 만들어 실행할 수 있습니다. 

코드를 실행하면 워크플로 실행에 대한 전체 로그를 확인할 수 있습니다. 

Dagster 실행 요약

Dagster 웹 서버

애셋과 잡 실행을 시각화하려면 Dagster 웹 서버를 설치하고 실행해야 합니다. 웹 서버에서 잡을 실행하고, 개별 애셋을 머터리얼라이즈하며, 여러 잡을 한 번에 모니터링할 수 있습니다.

$ pip install dagster-webserver

Dagster 서버를 시작하려면 Daster CLI를 사용해 Python 파일 경로를 제공합니다. 여기서는 파일명을 dagster_pipe.py로 지정했습니다.

$ dagster dev -f dagster_pipe.py  

위 명령을 실행하면 브라우저에서 자동으로 웹 서버가 열립니다. 또는 브라우저에서 http://127.0.0.1:3000 주소로 바로 이동할 수도 있습니다.

Dagster 웹 서버

Dagster 웹 서버 UI.

현재는 잡만 배포한 상태입니다. 워크플로를 실행하려면 “Runs” 탭으로 이동해 “Launch a new run” 버튼을 클릭하세요. 

실행이 성공적으로 완료되어야 합니다! 로그를 보려면 관심 있는 실행의 ID를 클릭하세요.

Dagster 실행 상세 보기

Dagster 실행 로그.

3. Mage AI

Mage AI는 오픈 소스 하이브리드 데이터 오케스트레이션 프레임워크입니다. 하이브리드는 Jupyter Notebook의 유연성과 모듈식 코드의 통제력을 모두 제공한다는 의미입니다. 

Python 지식이 많지 않아도 누구나 데이터 파이프라인을 구축, 실행, 모니터링할 수 있습니다. Python 파일을 직접 작성하고 실행하는 대신, Mage AI 프로젝트를 생성해 대시보드에서 파이프라인을 구축하고 실행·관리합니다.

Airflow와 비교하면 Mage AI는 사용자 친화적인 인터페이스와 사용 편의성을 제공해 데이터 엔지니어링 초보자에게 특히 적합합니다. 확장성을 염두에 두고 설계되어 대용량 데이터와 복잡한 파이프라인 구조도 효율적으로 처리할 수 있습니다.

제가 익숙한 방식과 완전히 달라서 다소 낯설었습니다. Mage AI 웹 UI를 설치하고 실행해야 했습니다. 쉬울 것이라 생각했지만 ETL 파이프라인을 구축하고 실행하는 데 어려움이 있었습니다. 반면, 드래그 앤 드롭과 버튼 클릭 위주라 이 분야에 새로 입문하는 사람들에게는 이 독특한 설계가 매력적일 수 있겠다는 생각이 들었습니다.

Abid Ali AwanAuthor

Mage AI 시작하기

Mage AI 시작은 매우 간단합니다. Mage AI Python 패키지를 설치하기만 하면 됩니다.

$ pip install mage-ai

그리고 Mage AI 프로젝트를 시작합니다. 

$ mage start mage_ai_etl 

위 명령은 웹 서버를 시작합니다. 앞서 언급했듯이, 코드 편집, 잡 실행, 잡 모니터링은 모두 Mage AI UI에서 이루어집니다.

Mage AI UI

Mage AI UI.

“+ New pipeline”을 클릭해 첫 ETL 파이프라인을 만듭니다. 저는 이름을 “simple_etl”로 지정했습니다.

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'

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” 모듈을 만들고 다음 코드를 추가합니다. 이 코드는 변환된 데이터를 CSV 파일로 저장하는 데이터 로딩 함수 export_data_to_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”를 클릭하세요.

Mage AI에서 파이프라인 실행

Mage AI에서 파이프라인 실행.

실행 로그를 보려면 “Runs” 탭으로 이동해 최근 실행한 파이프라인의 “Logs” 버튼을 클릭하세요.

Mage AI 플로 실행 로그

Mage AI 플로 실행 로그.

4. Kedro

Kedro는 다른 도구들과는 약간 다른, 또 하나의 인기 있는 오픈 소스 데이터 오케스트레이션 프레임워크입니다. 머신 러닝 엔지니어를 위해 만들어졌으며, 소프트웨어 엔지니어링의 많은 개념을 차용해 머신 러닝 프로젝트에 적용합니다.

Kedro는 고도의 모듈식으로 설계되어, 데이터셋을 내보내기 위해서도 데이터 위치와 유형을 지정하는 데이터 카탈로그를 만들어야 합니다. 이를 통해 파이프라인 전반의 데이터 관리를 표준화하고 효율적으로 수행합니다.

Kedro가 머신 러닝 생태계에서 어떤 역할을 하는지 이해하려면, 2024년에 알아야 할 MLOps 도구 25선 글을 참고해 다양한 MLOps 도구를 살펴보세요.

Airflow와 비교하면, Kedro API는 데이터 파이프라인을 구축하기 더 단순합니다. 머신 러닝 엔지니어링에 더 초점을 맞추고 데이터 분류 및 버저닝을 제공합니다.

코딩 자체는 비교적 간단하지만, 파이프라인을 실제로 실행하려면 문제가 생길 수 있습니다. 데이터 카탈로그를 만들고, 파이프라인을 등록하며, Kedro 프로젝트 구조를 파악해야 합니다. Dagster와 Prefect에 비해 더 도전적이라고 말할 수 있습니다. 하지만 이렇게 설계된 이유는 데이터 파이프라인을 신뢰성 있고 오류 없이 만들기 위함이라는 점을 이해합니다.

Abid Ali AwanAuthor

Kedro 시작하기

Kedro 데이터 파이프라인 구축은 조금 다른 접근이 필요합니다. 프레임워크가 모듈식이므로 프로젝트 구조와 워크플로 실행에 필요한 여러 단계를 이해해야 성공적으로 실행할 수 있습니다. 

먼저 Kedro Python 패키지를 설치하세요. 

$ 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  

kedro_pipe.py라는 Python 파일을 생성하고, 선호하는 IDE(예: Visual Studio Code)로 엽니다.

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

Python 스크립트에는 파이프라인의 노드에 해당하는 추출, 변환, 적재 함수가 포함되어야 합니다. 여기서는 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

Kedro 파이프라인 실행 로그

Kedro 파이프라인 실행 로그.

파이프라인 실행 후, 데이터 카탈로그에서 지정한 위치에 CSV 형식으로 파일이 저장됩니다.

Kedro 파이프라인 실행 결과 파일

Kedro 파이프라인 실행 결과 파일.

파이프라인 실행에 문제가 있다면 Kedro를 모든 확장과 함께 설치해 보세요. 

$ pip install "kedro[all]"

Kedro 시각화

파이프라인을 시각화하고 공유하려면 kedro-viz 도구를 설치합니다. 
$ pip install kedro-viz

그다음 아래 명령을 실행하면 모든 데이터 파이프라인과 데이터 노드를 시각화할 수 있습니다. 실험 추적 옵션과 파이프라인 시각화 공유 기능도 제공합니다.

$ kedro viz run

Kedro 시각화

Kedro 파이프라인 시각화.

5. Luigi

Luigi는 Spotify에서 개발한 오픈 소스 Python 기반 프레임워크로, 장시간 실행되는 배치 처리와 복잡한 데이터 파이프라인 관리에 강점을 보입니다. 의존성 해석, 워크플로 관리, 시각화, 장애 복구에 능해 데이터 워크플로 오케스트레이션에 유용합니다. 

Airflow와 비교하면 Luigi는 미니멀한 API, 캘린더 기반 스케줄링을 제공하며, 데이터 오케스트레이션 파이프라인 관련 이슈를 도와줄 충성도 높은 사용자 커뮤니티가 있습니다. 

Python 초보자라면 파이프라인을 구축하고 실행하는 일이 어렵게 느껴질 수 있습니다. 하지만 문서와 가이드를 참고하면 빠르게 시작할 수 있습니다. 로그는 제공 정보가 제한적이며, 대시보드는 DAG와 의존성의 시각화 도구에 가깝습니다.

Abid Ali AwanAuthor

Luigi 시작하기

Luigi 데이터 파이프라인을 만들려면 객체지향 프로그래밍에 대한 이해가 필요합니다. 먼저 Luigi Python 패키지를 설치해 보겠습니다. 

$ pip install luigi

Luigi에서 간단한 ETL 파이프라인을 개발하기 위해 상호 연결된 태스크를 만들겠습니다. Python 함수를 태스크로 만드는 대신, 파이프라인의 각 단계에 대해 FetchData, ProcessData, GenerateReport와 같은 Python 클래스를 만듭니다. 각 클래스에는 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)

위 코드를 주피터 노트북에서 실행하거나 Python 파일로 만들어 터미널에서 실행하세요. 

Luigi 실행 요약

Luigi와 유사하게, Apache Airflow로 ETL 파이프라인 구축 방법도 배울 수 있습니다. 이 튜토리얼은 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이고, 마지막 Luigi 태스크는 GenerateReport입니다.

$ python -m luigi --module luigi_pipe GenerateReport

파이프라인 실행과 태스크 상태를 시각화하려면 브라우저에서 http://localhost:8082로 이동하면 됩니다.

Luigi Central Planner 웹 UI

Luigi Central Planner 웹 UI.

이로써 Airflow의 다섯 가지 대안을 살펴보는 과정을 마쳤습니다! 이 글의 예제를 더 깊이 탐구하고 싶다면 다음 자료를 참고하세요.

  • Prefect, Dagster, Luigi의 소스 코드와 데이터는 DataLab 작업공간를 참고하세요.
  • Mage AI와 Kedro의 코드 소스 및 데이터는 GitHub 저장소를 참고하세요.

마무리 생각

이 튜토리얼에서는 오픈 소스이자 무료로 사용할 수 있는 Airflow 대안을 살펴보았습니다. 또한 각 데이터 오케스트레이션 도구를 알아보고, 간단한 ETL 파이프라인을 구축하고 실행해 보았습니다. 코드 예제를 통해 어떤 도구가 귀하의 사용 사례에 가장 적합한지 판단하는 데 도움이 될 것입니다.

초보자라면 사용자 친화적이며 설정이 간단한 Prefect나 Mage AI부터 시작하는 것을 권합니다. 반면 소프트웨어 엔지니어링 관행을 따르는 더 고급 도구를 찾는다면 Dagster, Kedro, Luigi를 살펴보세요.

이 글을 살펴본 뒤 데이터 엔지니어링 여정의 다음 단계로, DataCamp의 Data Engineer in Python 같은 인증을 취득해 다른 도구들을 배우고, 프로덕션에 배포할 수 있는 엔드 투 엔드 데이터 파이프라인을 구축해 보시기 바랍니다.

주제
데이터 엔지니어링
데이터 사이언스

이 강의들로 데이터 엔지니어링을 더 깊게 배워보세요!

courses

데이터 엔지니어링 입문

4
129.7K
이 단기 과정을 통해 ETL 및 클라우드 컴퓨팅과 같은 도구와 주제를 다루는 데이터 엔지니어링의 세계를 알아보세요.
자세히 보기Right Arrow
강좌 시작
더 보기Right Arrow