상위: 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_orders가 raw_orders를 인자로 받으므로 자동으로 DAG가 생성된다.
Airflow와의 근본적 차이
| 관점 | Airflow | Dagster |
|---|---|---|
| 단위 | Task (실행 단위) | Asset (데이터 단위) |
| 질문 | "이 작업을 언제 실행?" | "이 자산이 최신인가?" |
| 스케줄링 | Cron 기반 | Freshness Policy (선언적) |
| 리니지 | 태스크 의존성 | 데이터 자산 리니지 |
| 패러다임 | Imperative | Declarative |
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에는 없는 기능)
관련 개념
- Airflow: 가장 널리 쓰이는 워크플로우 오케스트레이터
- Prefect: Python-native 오케스트레이터
- Dagster vs Prefect vs Airflow: 3대 오케스트레이터 비교