전체 그래프
Dagster

Dagster

data-engineeringorchestrationdagster

상위: Data Engineering

요약

Dagster는 Asset 중심(Software-Defined Assets) 오케스트레이터이다. Airflow가 "이 태스크를 언제 실행할까"에 집중한다면, Dagster는 "이 데이터 자산이 최신인가"에 집중한다. 데이터 리니지, 로컬 개발, 테스트 경험이 뛰어나며, "Data as Code" 철학을 따른다.

핵심 개념

Software-Defined Assets (SDA)

Dagster의 핵심 추상화. 테이블, ML 모델, 리포트 같은 데이터 자산을 코드로 선언한다. 자산이 "어떤 상태여야 하는지"를 정의하면, 오케스트레이터가 실행 로직을 알아서 처리한다.

from dagster import asset

@asset
def raw_orders():
    """S3에서 주문 데이터를 읽어온다."""
    return pd.read_parquet("s3://bucket/orders.parquet")

@asset
def cleaned_orders(raw_orders):
    """주문 데이터를 정제한다. raw_orders에 의존."""
    return raw_orders.dropna().drop_duplicates()

@asset
def order_metrics(cleaned_orders):
    """일별 주문 지표를 계산한다."""
    return cleaned_orders.groupby("date").agg(
        total=("amount", "sum"),
        count=("order_id", "count")
    )

함수 인자가 곧 의존성이다. cleaned_ordersraw_orders를 인자로 받으므로 자동으로 DAG가 생성된다.

Airflow와의 근본적 차이

관점AirflowDagster
단위Task (실행 단위)Asset (데이터 단위)
질문"이 작업을 언제 실행?""이 자산이 최신인가?"
스케줄링Cron 기반Freshness Policy (선언적)
리니지태스크 의존성데이터 자산 리니지
패러다임ImperativeDeclarative

Freshness Policy (선언적 스케줄링)

Cron 표현식 대신 "이 자산은 최대 1시간 전 데이터까지 허용"처럼 선언한다.

from dagster import asset, FreshnessPolicy

@asset(freshness_policy=FreshnessPolicy(maximum_lag_minutes=60))
def order_metrics(cleaned_orders):
    ...

Dagster가 알아서 상위 자산부터 순서대로 실행한다. 이미 최신이면 실행하지 않는다.

Partitions (파티션)

시간/카테고리별로 자산을 분할하여 증분 처리한다.

from dagster import asset, DailyPartitionsDefinition

@asset(partitions_def=DailyPartitionsDefinition(start_date="2024-01-01"))
def daily_orders(context):
    date = context.partition_key  # "2024-01-15"
    return query_orders(date)

Resources (외부 시스템 추상화)

DB 연결, API 클라이언트 등을 Resource로 추상화한다. 테스트 시 Mock으로 교체 가능.

from dagster import asset, ConfigurableResource

class DatabaseResource(ConfigurableResource):
    connection_string: str

    def query(self, sql: str):
        ...

@asset
def users(database: DatabaseResource):
    return database.query("SELECT * FROM users")

프로덕션에서는 실제 DB, 테스트에서는 Mock DB를 주입한다.

로컬 개발 & 테스트

Dagster의 가장 큰 강점 중 하나. Docker 없이 로컬에서 전체 파이프라인을 실행/테스트할 수 있다.

# 로컬 개발 서버 시작 (UI 포함)
dagster dev

# 단위 테스트  일반 Python 함수처럼 테스트
def test_cleaned_orders():
    raw = pd.DataFrame({"order_id": [1, 1, 2], "amount": [100, 100, 200]})
    result = cleaned_orders(raw)
    assert len(result) == 2  # 중복 제거 확인

Airflow는 로컬에서 테스트하려면 DB 초기화 + Webserver + Scheduler를 모두 띄워야 한다. Dagster는 dagster dev 한 줄이면 된다.

아키텍처

Dagster Instance (메타데이터 DB)
├── Dagster Daemon (스케줄링, 센서)
├── Dagster Webserver (UI)
└── Code Location (사용자 코드)
    ├── Assets
    ├── Jobs
    ├── Schedules
    └── Sensors
  • Code Location: 사용자 코드가 독립 프로세스로 실행. 코드 변경 시 오케스트레이터 재시작 불필요
  • Sensor: 외부 이벤트(S3 파일 도착, API 응답)를 감지하여 자산 실행 트리거

Dagster Cloud

  • Serverless: 인프라 관리 없이 바로 사용. GitHub PR마다 자동 프리뷰 환경 생성
  • Hybrid: 오케스트레이션은 Cloud, 실행은 자체 인프라
  • Branch Deployment: PR별로 독립된 환경에서 파이프라인 테스트 가능 (Airflow에는 없는 기능)

관련 개념