📊 데이터공학

Prefect

Prefect

현대적 워크플로우 오케스트레이션. 파이썬 네이티브.

상세 설명

Prefect는 2018년 설립된 동명의 회사가 개발한 현대적 워크플로우 오케스트레이션 도구입니다. "Negative Engineering"을 최소화한다는 철학으로, 에러 처리, 재시도, 로깅을 자동화해 데이터 엔지니어가 비즈니스 로직에 집중하도록 합니다.

@flow와 @task 데코레이터가 핵심입니다. 일반 Python 함수에 데코레이터를 추가하면 워크플로우로 변환됩니다. Flow는 워크플로우 전체, Task는 개별 단위 작업을 나타내며, Task 간 의존성은 함수 호출 순서로 자동 추론됩니다.

Prefect 2.0(현재 버전)은 Prefect 1.x 대비 크게 단순화되었습니다. DAG를 명시적으로 정의할 필요 없이 일반 Python 코드처럼 작성합니다. 클라우드 UI로 워크플로우 모니터링, 스케줄링, 알림 설정이 가능하며, 로컬 실행도 서버 없이 됩니다.

재시도(retries), 캐싱(cache_key_fn), 타임아웃(timeout_seconds), 동시성 제어(concurrency_limits)를 Task 레벨에서 설정합니다. Dask, Ray 연동으로 병렬 처리하고, Docker, Kubernetes, AWS ECS 등 다양한 환경에 배포할 수 있습니다.

코드 예제

from prefect import flow, task, get_run_logger
from prefect.tasks import task_input_hash
from datetime import timedelta
import httpx
import pandas as pd

# Task 정의 - 재시도, 캐싱 설정
@task(
    retries=3,
    retry_delay_seconds=60,
    cache_key_fn=task_input_hash,
    cache_expiration=timedelta(hours=1)
)
def extract_data(url: str) -> dict:
    """API에서 데이터 추출"""
    logger = get_run_logger()
    logger.info(f"Fetching data from {url}")

    response = httpx.get(url, timeout=30)
    response.raise_for_status()
    return response.json()

@task
def transform_data(raw_data: dict) -> pd.DataFrame:
    """데이터 변환"""
    logger = get_run_logger()
    df = pd.DataFrame(raw_data['results'])

    # 데이터 정제
    df['created_at'] = pd.to_datetime(df['created_at'])
    df['value'] = df['value'].fillna(0)

    logger.info(f"Transformed {len(df)} rows")
    return df

@task(timeout_seconds=300)
def load_data(df: pd.DataFrame, destination: str) -> int:
    """데이터 적재"""
    logger = get_run_logger()

    # 실제로는 DB, S3 등에 저장
    df.to_parquet(destination)
    logger.info(f"Saved {len(df)} rows to {destination}")

    return len(df)

# Flow 정의
@flow(
    name="ETL Pipeline",
    description="Daily data extraction and loading",
    version="1.0.0"
)
def etl_pipeline(
    api_url: str = "https://api.example.com/data",
    output_path: str = "data/output.parquet"
):
    """메인 ETL 워크플로우"""
    logger = get_run_logger()
    logger.info("Starting ETL pipeline")

    # Task 실행 - 의존성 자동 추론
    raw_data = extract_data(api_url)
    transformed_df = transform_data(raw_data)
    rows_loaded = load_data(transformed_df, output_path)

    logger.info(f"Pipeline completed: {rows_loaded} rows processed")
    return rows_loaded

# 로컬 실행
if __name__ == "__main__":
    etl_pipeline()

# 스케줄링 설정 (배포 시)
# from prefect.deployments import Deployment
# from prefect.server.schemas.schedules import CronSchedule
#
# deployment = Deployment.build_from_flow(
#     flow=etl_pipeline,
#     name="daily-etl",
#     schedule=CronSchedule(cron="0 2 * * *"),  # 매일 새벽 2시
#     work_queue_name="default"
# )
# deployment.apply()

# 병렬 처리 (여러 소스 동시 추출)
@flow
def parallel_etl():
    sources = ["https://api1.example.com", "https://api2.example.com"]

    # 병렬 실행
    futures = [extract_data.submit(url) for url in sources]
    results = [f.result() for f in futures]

    return results

# CLI 실행
# prefect deployment run "ETL Pipeline/daily-etl"
# prefect server start  # UI 서버 시작

실무에서 이렇게 말해요

시니어: "Airflow 대신 Prefect 쓰면 DAG 정의 없이 Python 함수 그대로 쓸 수 있어요. 러닝 커브가 훨씬 낮아요."

주니어: "스케줄링은 어떻게 해요?"

시니어: "Deployment 만들 때 CronSchedule이나 IntervalSchedule 붙이면 돼요. Prefect Cloud UI에서 모니터링하고 수동 트리거도 가능해요."

면접관: "데이터 파이프라인 오케스트레이션 경험이 있나요?"

지원자: "Prefect로 일 100개 이상의 ETL Job을 관리했습니다. @task의 retries와 cache_key_fn으로 실패 복구와 중복 실행 방지를 구현하고, Prefect Cloud에서 실패 알림을 Slack으로 연동했습니다. 기존 Cron 스크립트 대비 가시성과 에러 처리가 크게 개선되었습니다."

리뷰어: "Task에 timeout_seconds 추가하세요. API 호출이 무한 대기하면 전체 파이프라인이 멈춰요."

개발자: "네, retries=3에 timeout_seconds=60 설정하겠습니다. 로깅도 get_run_logger()로 바꿀게요."

주의사항

더 배우기