파이프라인에서 작업이 성공했다는 로그만으로는 결과 데이터가 올바른지 알기 어렵다. 어떤 원본으로 집계했는지, 마지막 갱신은 언제인지, 품질 검사가 통과했는지를 데이터와 함께 보고 싶을 때 Dagster의 자산 모델이 도움이 된다.
Dagster는 데이터 자산과 그 의존성을 코드로 정의하는 오케스트레이터다. 이 글은 Dagster 1.x의 @asset, materialize, asset check와 AutomationCondition을 구분해 설명한다. 예제는 Python 3.12, Dagster 1.11 이상, pandas 2.x를 기준으로 한다.
자산은 파일, 테이블, 모델처럼 관리할 가치가 있는 데이터 결과를 나타낸다. @asset 함수는 그 결과를 만드는 계산을 정의한다. 입력 인자로 다른 자산을 받으면 의존성을 표현할 수 있다.
raw_sales 원본 매출 → daily_summary 날짜별 집계
↓
품질 검사 결과 기록
이 선언은 데이터 계보를 이해하는 데 도움이 된다. 하지만 함수 이름만 선언했다고 외부 원본의 변경을 자동으로 감지하거나 데이터가 항상 최신으로 유지되지는 않는다. 언제 실행할지는 별도의 실행·자동화 규칙이다.
아래 예제는 고정 입력 두 건을 날짜별로 합산한다. 문자열인 가게 이름까지 sum()에 포함하지 않고, 금액 열만 선택한다.
python -m pip install "dagster>=1.11,<2" "pandas>=2.2,<3"
python my_assets.py
# my_assets.py
import dagster as dg
import pandas as pd
@dg.asset
def raw_sales() -> pd.DataFrame:
return pd.DataFrame([
{"date": "2025-01-01", "store": "A", "amount": 100},
{"date": "2025-01-01", "store": "B", "amount": 150},
])
@dg.asset
def daily_summary(raw_sales: pd.DataFrame) -> pd.DataFrame:
return raw_sales.groupby("date", as_index=False)["amount"].sum()
if __name__ == "__main__":
result = dg.materialize(
[raw_sales, daily_summary],
resources={"io_manager": dg.mem_io_manager},
)
assert result.success
print(result.output_for_node("daily_summary").to_dict("records"))
애플리케이션 출력은 다음과 같다. Dagster 실행 로그는 별도로 표시된다.
[{'date': '2025-01-01', 'amount': 250}]
이 예제는 메모리 I/O manager를 사용하므로 프로세스가 종료되면 결과가 유지되지 않는다. 운영에서는 DB·오브젝트 스토리지 등에 연결한 I/O manager나 명시적인 저장 로직이 필요하다. “구체화했다”는 이벤트와 실제 데이터의 저장 위치·보존 기간을 함께 이해해야 한다.
materialize([raw_sales, daily_summary])는 제공한 두 자산을 한 번 실행한다. 이 예제에서는 raw_sales를 만든 뒤 daily_summary가 그 값을 사용한다.
하위 자산만 선택해 실행하면 선택하지 않은 상위 자산을 무조건 다시 계산하는 것은 아니다. 필요한 상위 입력을 I/O manager가 저장소에서 불러오는 구성이 필요하다. 메모리에만 있던 이전 실행 결과가 다음 프로세스에서도 자동으로 복원되지는 않는다.
다음 개념을 분리하면 “Dagster가 알아서 최신 데이터로 만든다”는 오해를 줄일 수 있다.
| 개념 | 결정하는 것 |
|---|---|
| 자산 의존성 | 어떤 입력으로 어떤 결과를 만드는가 |
| 실행 대상 선택 | 이번 실행에 어떤 자산을 포함할 것인가 |
| I/O manager | 결과를 어디에 저장하고 입력을 어떻게 읽을 것인가 |
| 자동화 조건 | 어떤 상태 변화에 실행을 요청할 것인가 |
| 품질 검사 | 결과가 정한 조건을 만족하는가 |
자산 함수가 예외 없이 끝났다고 데이터가 유효한 것은 아니다. 행이 비었는지, 금액이 음수인지, 원본 건수와 집계가 맞는지처럼 서비스에 필요한 조건을 코드로 작성한다.
# definitions.py
import dagster as dg
import pandas as pd
from my_assets import raw_sales, daily_summary
@dg.asset_check(asset=daily_summary, blocking=True)
def non_negative(daily_summary: pd.DataFrame) -> dg.AssetCheckResult:
valid = (
not daily_summary.empty
and daily_summary["amount"].notna().all()
and (daily_summary["amount"] >= 0).all()
)
return dg.AssetCheckResult(passed=bool(valid))
defs = dg.Definitions(
assets=[raw_sales, daily_summary],
asset_checks=[non_negative],
jobs=[dg.define_asset_job("sales_checked")],
resources={"io_manager": dg.mem_io_manager},
)
if __name__ == "__main__":
result = defs.resolve_job_def("sales_checked").execute_in_process()
assert result.success
blocking=True인 검사는 실패 시 관련 하위 실행을 막는 데 사용한다. 검사 범위와 심각도, 하위 의존성을 함께 구성해야 하며, 이미 외부에 저장한 데이터를 자동으로 롤백하는 기능은 아니다.
위 검사는 음수가 없는지만 본다. 잘못된 양수 합계나 누락된 날짜까지 찾아 주지는 않는다. 행 수, 키 중복, 원본 대비 합계 같은 규칙을 요구사항에 맞춰 추가한다. 품질 검사를 썼다는 사실보다 어떤 오류를 검출하도록 정의했는지가 중요하다.
자산이 존재한다는 이유로 주기적으로 실행되지는 않는다. 일정에 따라 실행하려면 schedule을, 외부 사건을 관찰하려면 sensor를, 자산 상태를 기준으로 실행하려면 선언적 자동화를 사용할 수 있다.
다음은 원본이 갱신됐을 때 파생 자산의 실행을 요청하는 정의다. 앞 예제와 별도 코드 위치에서 살펴볼 수 있다.
# automation.py
import dagster as dg
import pandas as pd
from my_assets import raw_sales
@dg.asset(automation_condition=dg.AutomationCondition.eager())
def auto_summary(raw_sales: pd.DataFrame) -> pd.DataFrame:
return raw_sales.groupby("date", as_index=False)["amount"].sum()
defs = dg.Definitions(assets=[raw_sales, auto_summary])
정의를 로드한 다음 UI에서 기본 automation condition sensor를 활성화하고, 해당 자동화를 평가하는 daemon이 실행돼 있어야 한다. eager()는 의존 자산의 변경 등에 반응하지만 상위 자산이 없거나 실행 중인 경우처럼 대기하는 조건도 있다. 단순히 Python 파일을 한 번 실행한 것만으로 계속 갱신되지는 않는다.
이 예제는 raw_sales를 주기적으로 수집하는 규칙을 포함하지 않는다. 외부 원본의 변경을 Dagster가 알 수 있도록 수집 실행이나 관찰 이벤트를 연결해야 한다. 하위 자산의 자동화 규칙과 원본 수집 규칙은 다른 책임이다.
Dagster UI에서는 자산 의존성, 구체화 이력, 실행 로그와 검사 결과를 볼 수 있다. 예전 자료의 ‘Dagit’이라는 이름은 현재 문서에서 Dagster webserver·UI로 안내하는 부분과 연결해 읽으면 된다.
웹서버, daemon, 사용자 코드 위치와 실행 프로세스는 서로 역할이 다르다. 웹 화면이 열린다고 자동화와 실행 환경까지 모두 준비된 것은 아니다. 운영에서는 실행·이벤트 기록 저장소, 결과 데이터 저장소, 코드 버전과 의존성을 함께 구성한다.
UI를 로컬에서 사용하려면 dagster-webserver 등 대응 패키지를 설치하고 dagster dev -f definitions.py처럼 코드 위치를 지정한다. 개발용 실행과 운영 배포를 같은 가용성 수준으로 보지 않는다.
날짜별 데이터를 다룬다면 partition을 정의해 처리 단위를 명확하게 할 수 있다. 과거 partition을 다시 처리하는 backfill에서는 같은 구간을 덮어쓸지, 별도 버전으로 보존할지 정한다. 실행을 다시 요청하는 기능은 외부 저장의 멱등성을 대신하지 않는다.
여러 팀이 데이터 자산의 계보·갱신 상태·품질을 함께 이해해야 한다면 Dagster의 모델이 도움이 될 수 있다. 반대로 단순한 스크립트 몇 개만 주기적으로 실행한다면 자산 모델과 운영 구성의 학습 비용을 고려한다.
Airflow나 Prefect에서도 품질 검증과 데이터 중심 운영을 구성할 수 있다. Dagster만 가능한 기능이라고 단정하기보다 각 도구가 기본으로 제공하는 모델, 실행·배포 방식과 팀의 관리 단위를 비교한다.
Dagster를 이해하는 출발점은 자산 선언, 데이터 저장, 품질 검사, 실행 조건이 서로 다른 책임이라는 점이다. 이 경계를 명시하면 실행 성공과 데이터의 신뢰성을 더 정확하게 설명할 수 있다.