Tutorial Dagster: Orkestrasi Data dengan Software-Defined Assets

# Dagster: Orkestrasi Data Modern dengan Software-Defined Assets Dagster adalah orkestrator data yang menyusun pipeline berdasarkan data yang dihasilkan, bukan berdasarkan task yang dijalankan. Tutor...

By Ruby Abdullah · · tutorial
DagsterData OrchestrationData PipelineETLData EngineeringPython

Dagster: Orkestrasi Data Modern dengan Software-Defined Assets

Dagster adalah orkestrator data yang menyusun pipeline berdasarkan data yang dihasilkan, bukan berdasarkan task yang dijalankan. Tutorial ini menelusuri konsep-konsep inti menggunakan pipeline data penjualan sebagai contoh berjalan, mulai dari instalasi hingga deployment. Di akhir tutorial Anda akan memahami asset, resource, IO manager, partisi, schedule, sensor, serta pengecekan kualitas data.

Apa Itu Dagster

Dagster adalah framework orkestrasi untuk membangun, menjalankan, dan mengamati pipeline data. Framework ini dirancang agar pipeline dapat diuji, diamati, dan dipelihara layaknya perangkat lunak. Alih-alih memperlakukan pipeline sebagai graf task yang buram, Dagster mendorong Anda untuk mendeklarasikan asset yang dihasilkan pipeline — tabel, file, model machine learning, dashboard — dan membiarkan framework menentukan urutan eksekusi dari ketergantungan antar asset tersebut.

Sebuah asset adalah objek persisten di penyimpanan yang merepresentasikan pemahaman tertentu tentang dunia. Tabel penjualan harian di sebuah warehouse adalah asset. File Parquet yang sudah dibersihkan adalah asset. Model yang sudah dilatih adalah asset. Ketika Anda menggambarkan pipeline dalam bentuk output-output ini, orkestrator memperoleh model yang lebih kaya tentang sistem Anda, yang pada gilirannya meningkatkan observabilitas dan lineage.

Perbedaan Dagster dengan Airflow dan Prefect

Airflow dan Prefect pada dasarnya bersifat task-centric. Anda mendefinisikan task (atau flow) dan menyusun urutan eksekusinya. Orkestrator tahu bahwa task B berjalan setelah task A, tetapi secara bawaan ia tidak tahu bahwa task A menghasilkan tabel rawsales dan task B mengonsumsinya. Lineage dan kesadaran terhadap data harus ditambahkan secara terpisah.

Dagster bersifat asset-centric. Anda mendeklarasikan asset, dan graf ketergantungan antar asset diturunkan dari cara mereka saling mereferensikan. Pergeseran ini membawa konsekuensi praktis:

  • Lineage sudah terpasang sejak awal. Graf asset adalah graf lineage. Anda dapat melihat asset hulu mana yang menyuplai asset hilir mana langsung di UI.
  • Bernalar dalam kerangka data. "Materialisasi ulang asset dailysalessummary" lebih alami daripada "jalankan ulang task transform lalu task load."
  • Partisi sebagai warga kelas satu. Setiap asset dapat dipartisi (misalnya satu partisi per hari), dan Dagster melacak partisi mana yang sudah dimaterialisasi.
  • Kualitas data bawaan. Asset check memungkinkan Anda melampirkan validasi pada asset dan menampilkan kegagalannya di UI.

Dagster tetap mendukung eksekusi gaya task melalui ops dan jobs (dibahas nanti), jadi Anda tidak dipaksa memakai model asset untuk segala hal. Namun untuk sebagian besar pipeline data, asset adalah titik awal yang disarankan.

Instalasi

Dagster berjalan pada Python 3.9 atau lebih baru. Buat virtual environment lalu pasang paket inti bersama UI web.

python -m venv .venv

source .venv/bin/activate

pip install dagster dagster-webserver

Paket dagster menyediakan pustaka inti dan perkakas baris perintah. Paket dagster-webserver menyediakan UI web lokal yang akan Anda gunakan untuk memeriksa dan menjalankan asset. Paket integrasi seperti dagster-pandas, dagster-duckdb, dan dagster-aws dipasang terpisah sesuai kebutuhan.

Verifikasi instalasi:

dagster --version

Tata letak proyek yang rapi untuk contoh-contoh di bawah ini terlihat seperti berikut:

salespipeline/

salespipeline/

init.py

assets.py

resources.py

pyproject.toml

Dekorator @asset dan Software-Defined Assets

Sebuah software-defined asset adalah fungsi Python yang didekorasi dengan @asset. Fungsi tersebut menghitung nilai asset, dan nilai kembaliannya diserahkan ke IO manager untuk disimpan. Nama fungsi menjadi key dari asset tersebut.

from dagster import asset

import pandas as pd

@asset

def rawsales() -> pd.DataFrame:

"""Ekstrak catatan penjualan mentah dari CSV sumber."""

return pd.readcsv("https://example.com/sales.csv")

Ketergantungan antar asset dinyatakan dengan menambahkan parameter yang namanya cocok dengan key asset hulu. Dagster menyuntikkan nilai asset hulu saat runtime.

@asset

def cleanedsales(rawsales: pd.DataFrame) -> pd.DataFrame:

"""Transform: buang nilai null dan normalkan nama kolom."""

df = rawsales.dropna(subset=["orderid", "amount"])

df = df.rename(columns={"amt": "amount", "cust": "customerid"})

df["amount"] = df["amount"].astype(float)

return df

Di sini cleanedsales bergantung pada rawsales semata-mata karena ia mendeklarasikan parameter dengan nama tersebut. Dagster membaca hal ini dan membangun jalur ketergantungan secara otomatis — tidak ada langkah perangkaian terpisah.

Membangun Pipeline ETL Kecil sebagai Asset

Mari membangun pipeline extract-transform-load yang koheren untuk data penjualan. Ketiga asetnya adalah rawsales (extract), cleanedsales (transform), dan dailysalessummary (load/agregasi).

# salespipeline/assets.py

from dagster import asset

import pandas as pd

@asset

def rawsales() -> pd.DataFrame:

"""Extract: baca catatan penjualan mentah."""

data = {

"orderid": [1, 2, 3, 4, None],

"customerid": ["c1", "c2", "c1", "c3", "c2"],

"amount": [120.0, 80.5, None, 45.0, 60.0],

"orderdate": [

"2026-05-01", "2026-05-01", "2026-05-02",

"2026-05-02", "2026-05-02",

],

}

return pd.DataFrame(data)

@asset

def cleanedsales(rawsales: pd.DataFrame) -> pd.DataFrame:

"""Transform: buang baris tidak valid dan terapkan tipe data."""

df = rawsales.dropna(subset=["orderid", "amount"]).copy()

df["orderid"] = df["orderid"].astype(int)

df["amount"] = df["amount"].astype(float)

df["orderdate"] = pd.todatetime(df["orderdate"])

return df

@asset

def dailysalessummary(cleanedsales: pd.DataFrame) -> pd.DataFrame:

"""Load: agregasi pendapatan per hari."""

summary = (

cleanedsales.groupby("orderdate")["amount"]

.agg(["sum", "count"])

.resetindex()

.rename(columns={"sum": "totalrevenue", "count": "ordercount"})

)

return summary

Rantai ketergantungannya adalah rawsales -> cleanedsales -> dailysalessummary. Dagster menyimpulkannya dari parameter fungsi.

Objek Definitions

Sebuah proyek Dagster mengekspos isinya — asset, resource, schedule, sensor, dan job — melalui satu objek Definitions. Inilah titik masuk yang dimuat oleh Dagster.

# salespipeline/init.py

from dagster import Definitions, loadassetsfrommodules

from salespipeline import assets

allassets = loadassetsfrommodules([assets])

defs = Definitions(

assets=allassets,

)

loadassetsfrommodules memindai sebuah modul dan mengumpulkan setiap fungsi berdekorator @asset, sehingga Anda tidak perlu mendaftarkannya satu per satu. Anda juga dapat memberikan daftar asset eksplisit jika menginginkan kontrol yang lebih ketat.

Menjalankan UI dengan dagster dev

Server pengembangan lokal menggabungkan UI web dan daemon latar belakang yang diperlukan untuk schedule dan sensor. Jalankan dari root proyek, mengarah ke modul yang mendefinisikan defs.

dagster dev -m salespipeline

Secara bawaan UI disajikan di http://localhost:3000. Buka UI tersebut dan Anda akan melihat graf asset untuk pipeline penjualan. Dari halaman Assets Anda dapat mengklik Materialize all untuk menjalankan seluruh rantai, atau memilih satu asset untuk memateralisasi ulang hanya asset itu beserta asset hulu yang diperlukannya.

Setiap materialisasi dicatat. UI menampilkan riwayat run, log, graf lineage, dan metadata yang dilampirkan pada setiap asset.

Ops dan Jobs (Singkat)

Sebelum software-defined assets, abstraksi inti Dagster adalah ops dan jobs. Sebuah op adalah unit komputasi; sebuah job adalah graf op yang dirangkai bersama. Ops tetap berguna untuk pekerjaan prosedural yang tidak menghasilkan asset persisten — mengirim notifikasi, memanggil API eksternal untuk efek samping, atau menjalankan skrip sembarang.

from dagster import op, job

@op

def fetchcount() -> int:

return 42

@op

def report(count: int) -> None:

print(f"Memproses {count} catatan")

@job

def reportingjob():

report(fetchcount())

Pada proyek modern Anda umumnya bekerja dengan asset, dan Dagster dapat membangun asset job yang memateralisasi sekumpulan asset. Gunakan ops dan jobs ketika tugas memang benar-benar tidak memiliki output berupa asset.

Resource dan Konfigurasi

Resource memodelkan koneksi dan sistem eksternal — basis data, API, penyimpanan cloud. Mendefinisikannya secara terpisah menjaga asset tetap bersih dan memudahkan penggantian saat pengujian. Resource modern adalah subclass dari ConfigurableResource.

# salespipeline/resources.py

from dagster import ConfigurableResource

import duckdb

import pandas as pd

class DuckDBResource(ConfigurableResource):

databasepath: str

def query(self, sql: str) -> pd.DataFrame:

with duckdb.connect(self.databasepath) as conn:

return conn.execute(sql).fetchdf()

def writetable(self, table: str, df: pd.DataFrame) -> None:

with duckdb.connect(self.databasepath) as conn:

conn.register("tmpdf", df)

conn.execute(

f"CREATE OR REPLACE TABLE {table} AS SELECT FROM tmpdf"

)

Sebuah asset mendeklarasikan resource yang dibutuhkannya sebagai parameter bertipe, dan Dagster menyuntikkan instance yang telah dikonfigurasi.

from dagster import asset

import pandas as pd

from salespipeline.resources import DuckDBResource

@asset

def dailysalessummary(

cleanedsales: pd.DataFrame, warehouse: DuckDBResource

) -> pd.DataFrame:

summary = (

cleanedsales.groupby("orderdate")["amount"]

.agg(["sum", "count"])

.resetindex()

.rename(columns={"sum": "totalrevenue", "count": "ordercount"})

)

warehouse.writetable("dailysalessummary", summary)

return summary

Rangkai resource ke dalam Definitions dan berikan konfigurasinya. Membaca path dari variabel lingkungan menjaga rahasia dan nilai khas deployment tetap di luar kode.

from dagster import Definitions, EnvVar, loadassetsfrommodules

from salespipeline import assets

from salespipeline.resources import DuckDBResource

defs = Definitions(

assets=loadassetsfrommodules([assets]),

resources={

"warehouse": DuckDBResource(databasepath=EnvVar("DUCKDBPATH")),

},

)

IO Manager

Sebuah IO manager mengatur cara nilai kembalian asset disimpan dan cara nilai itu dimuat kembali ketika asset hilir membutuhkannya. Ini memisahkan apa yang dihitung asset dari di mana datanya berada. Secara bawaan Dagster melakukan pickle terhadap nilai ke sistem berkas lokal, yang memadai untuk pengembangan tetapi jarang sesuai untuk produksi.

Anda dapat melampirkan IO manager secara global atau per asset. Pengaturan umum memakai IO manager berbasis warehouse sehingga mengembalikan sebuah DataFrame berarti menulis tabel, dan bergantung pada asset itu berarti membaca tabel kembali.

from dagster import Definitions, FilesystemIOManager, loadassetsfrommodules

from salespipeline import assets

defs = Definitions(

assets=loadassetsfrommodules([assets]),

resources={

"iomanager": FilesystemIOManager(basedir="data/storage"),

},

)

Ide kuncinya: asset berfokus pada logika transformasi dan mengembalikan objek Python biasa, sementara IO manager memegang tanggung jawab penyimpanan. Mengganti penyimpanan (file lokal saat dev, warehouse saat prod) menjadi perubahan konfigurasi, bukan perubahan kode. Paket integrasi menyediakan IO manager siap pakai, misalnya dagster-duckdb-pandas dan dagster-snowflake-pandas.

Schedule dan Sensor

Schedule dan sensor menentukan kapan asset sebaiknya dimaterialisasi.

Schedule

Sebuah schedule memicu job dengan irama mirip cron. Bangun job dari sebuah seleksi asset, lalu lampirkan schedule.

from dagster import (

AssetSelection,

defineassetjob,

ScheduleDefinition,

)

salesjob = defineassetjob(

name="salesjob", selection=AssetSelection.all()

)

dailysalesschedule = ScheduleDefinition(

job=salesjob,

cronschedule="0 6 ", # setiap hari pukul 06:00

)

Daftarkan keduanya di Definitions:

defs = Definitions(

assets=allassets,

jobs=[salesjob],

schedules=[dailysalesschedule],

)

Sensor

Sebuah sensor memantau kondisi eksternal dan meluncurkan run ketika kondisi terpenuhi — misalnya file baru yang tiba di sebuah bucket. Fungsi sensor menghasilkan RunRequest ketika pekerjaan perlu dijalankan.

import os

from dagster import sensor, RunRequest, SkipReason

@sensor(job=salesjob)

def newfilesensor(context):

dropdir = "data/incoming"

files = os.listdir(dropdir) if os.path.isdir(dropdir) else []

if not files:

return SkipReason("Tidak ada file baru ditemukan")

for filename in files:

yield RunRequest(runkey=filename)

runkey memastikan setiap file memicu tepat satu run, bahkan jika sensor melihat file yang sama pada tick berikutnya. Baik schedule maupun sensor membutuhkan daemon, yang dijalankan otomatis oleh dagster dev.

Partisi

Partisi membagi sebuah asset menjadi irisan-irisan yang dapat dimaterialisasi secara terpisah — paling umum satu irisan per hari. Ini memungkinkan Anda melakukan backfill riwayat, memproses ulang satu hari yang bermasalah, dan melacak dengan tepat hari mana yang sudah selesai.

from dagster import asset, DailyPartitionsDefinition

import pandas as pd

dailypartitions = DailyPartitionsDefinition(startdate="2026-05-01")

@asset(partitionsdef=dailypartitions)

def partitionedsales(context) -> pd.DataFrame:

partitiondate = context.partitionkey

# Muat hanya data untuk hari spesifik ini.

df = readsalesfordate(partitiondate)

context.log.info(f"Memuat {len(df)} baris untuk {partitiondate}")

return df

context.partitionkey adalah string tanggal untuk partisi yang sedang dimaterialisasi. Asset hilir yang juga dipartisi memetakan partisinya ke partisi hulu, sehingga dailysalessummary untuk 2026-05-02 bergantung pada partitionedsales untuk 2026-05-02. Di UI Anda mendapatkan grid partisi yang menunjukkan hari mana yang telah dimaterialisasi, hilang, atau gagal, dan Anda dapat meluncurkan backfill untuk rentang partisi dalam satu aksi.

Asset Check untuk Kualitas Data

Asset check melampirkan validasi pada sebuah asset dan melaporkan status lulus/gagal di UI. Mereka mengubah ekspektasi kualitas data menjadi objek kelas satu yang dapat diamati, bukan asersi yang berserakan.

from dagster import assetcheck, AssetCheckResult

import pandas as pd

@assetcheck(asset="cleanedsales")

def nonegativeamounts(cleanedsales: pd.DataFrame) -> AssetCheckResult:

badrows = (cleanedsales["amount"] < 0).sum()

return AssetCheckResult(

passed=bool(badrows == 0),

metadata={"negativerows": int(badrows)},

)

@assetcheck(asset="cleanedsales")

def orderidisunique(cleanedsales: pd.DataFrame) -> AssetCheckResult:

duplicates = int(cleanedsales["orderid"].duplicated().sum())

return AssetCheckResult(

passed=duplicates == 0,

metadata={"duplicateorderids": duplicates},

)

Daftarkan check bersama asset:

from dagster import Definitions, loadassetchecksfrommodules

from salespipeline import assets, checks

defs = Definitions(

assets=loadassetsfrommodules([assets]),

assetchecks=loadassetchecksfrommodules([checks]),

)

Check berjalan setelah asset dimaterialisasi (secara bawaan). Check yang gagal terlihat di UI di samping asetnya, dan Anda dapat mengonfigurasi run agar memblokir materialisasi hilir saat sebuah check gagal.

Menguji Asset di Python

Karena asset adalah fungsi Python biasa, Anda dapat mengujinya langsung dengan memanggilnya menggunakan input biasa — tanpa perlu orkestrator. Ini adalah salah satu manfaat praktis utama dari model asset.

# tests/testassets.py

import pandas as pd

from salespipeline.assets import cleanedsales, dailysalessummary

def testcleanedsalesdropsnulls():

raw = pd.DataFrame(

{

"orderid": [1, None],

"customerid": ["c1", "c2"],

"amount": [10.0, None],

"orderdate": ["2026-05-01", "2026-05-01"],

}

)

result = cleanedsales(raw)

assert len(result) == 1

assert result["orderid"].iloc[0] == 1

def testdailysummaryaggregates():

cleaned = pd.DataFrame(

{

"orderid": [1, 2],

"customerid": ["c1", "c2"],

"amount": [10.0, 20.0],

"orderdate": pd.todatetime(["2026-05-01", "2026-05-01"]),

}

)

summary = dailysalessummary(cleaned)

assert summary["totalrevenue"].iloc[0] == 30.0

assert summary["ordercount"].iloc[0] == 2

Untuk asset yang membutuhkan resource, Dagster menyediakan materializetomemory, yang mengeksekusi sekumpulan asset secara in-process dengan resource pengujian yang diberikan langsung.

from dagster import materializetomemory

from salespipeline.assets import rawsales, cleanedsales

def testpipelineinmemory():

result = materializetomemory([rawsales, cleanedsales])

assert result.success

Jalankan rangkaian pengujian dengan pytest:

pytest tests/

Catatan Deployment

Deployment Dagster di produksi memiliki dua proses yang berjalan terus-menerus:

  • dagster-webserver menyajikan UI dan API GraphQL yang dipakai klien untuk meluncurkan dan memeriksa run.
  • dagster-daemon menjalankan schedule, sensor, antrian run, dan backfill. Tanpa daemon, pemicu berbasis waktu dan berbasis peristiwa tidak akan aktif.

dagster-webserver -h 0.0.0.0 -p 3000 -m salespipeline

dagster-daemon run -m salespipeline

Kedua proses membaca direktori DAGSTERHOME yang sama, yang menyimpan konfigurasi instance (dagster.yaml) serta penyimpanan run/event. Di dagster.yaml Anda mengonfigurasi tempat run dieksekusi (misalnya run launcher Kubernetes), tempat log dan metadata run disimpan (sering kali Postgres), dan batas konkurensi.

Bagi tim yang lebih memilih opsi terkelola, Dagster+ (Dagster Cloud) meng-host control plane — UI web, penyimpanan metadata, schedule, dan daemon — sementara kode Anda berjalan di lingkungan Anda sendiri melalui agent. Ia menambahkan fitur seperti branch deployment untuk menguji pull request, kontrol akses berbasis peran, dan pemberitahuan. Self-hosting dengan webserver dan daemon sumber terbuka tetap sepenuhnya didukung; Dagster+ merupakan kemudahan operasional, bukan keharusan.

Pola deployment ber-container yang umum:

# docker-compose.yaml (disederhanakan)

services:

webserver:

image: my-org/sales-pipeline:latest

command: dagster-webserver -h 0.0.0.0 -p 3000 -m salespipeline

ports:

  • "3000:3000"
environment:

DAGSTERHOME: /opt/dagster/home

DUCKDBPATH: /data/warehouse.duckdb

daemon:

image: my-org/sales-pipeline:latest

command: dagster-daemon run -m salespipeline

environment:

DAGSTERHOME: /opt/dagster/home

DUCKDBPATH: /data/warehouse.duckdb

Praktik Terbaik

  • Modelkan datanya, bukan task-nya. Beri nama asset sesuai tabel dan file yang dihasilkannya. Graf asset semestinya terbaca seperti deskripsi warehouse data Anda.
  • Jaga logika transformasi tetap murni. Biarkan asset mengembalikan objek biasa dan biarkan IO manager menangani penyimpanan, agar fungsi tetap mudah diuji.
  • Letakkan koneksi di resource. Jangan pernah menanam kredensial atau string koneksi langsung di asset. Gunakan ConfigurableResource dengan EnvVar untuk rahasia.
  • Partisi asset deret waktu. Partisi harian (atau per jam) membuat backfill dan pemrosesan ulang menjadi murah dan presisi.
  • Kodekan ekspektasi kualitas sebagai asset check. Check yang berada di samping asetnya jauh lebih tahan lama dibanding asersi dadakan yang terkubur dalam kode transform.
  • Tulis unit test untuk fungsi asset. Panggil langsung dengan input contoh; sisakan materializetomemory untuk alur yang bergantung pada resource.
  • Pisahkan concern berdasarkan modul. Simpan asset, resource, schedule, dan check di file masing-masing lalu rakit dalam satu Definitions.
  • Jalankan daemon di setiap lingkungan yang memakai schedule atau sensor. Daemon yang tidak ada adalah penyebab paling umum pemicu gagal aktif secara diam-diam.

Kesimpulan dan Poin Penting

Model asset-centric Dagster membingkai ulang orkestrasi di sekitar data yang dihasilkan pipeline Anda. Dibandingkan perkakas task-centric, pendekatan ini memberi Anda lineage, observabilitas, dan pelacakan partisi tanpa perlu menempelkannya belakangan.

Poin penting:

  • Sebuah software-defined asset adalah fungsi Python berdekorator; ketergantungan berasal dari nama parameter.
  • Objek Definitions adalah satu titik masuk yang menyatukan asset, resource, schedule, sensor, dan check.
  • Resource dan IO manager memisahkan logika bisnis dari koneksi dan penyimpanan, yang membuat pipeline dapat diuji dan dapat dipindahkan.
  • Partisi, schedule, dan sensor mengontrol kapan dan atas irisan apa asset dimaterialisasi.
  • Asset check menjadikan kualitas data sebagai perhatian kelas satu yang terlihat.
  • Asset adalah fungsi biasa, sehingga Anda dapat menguji sebagian besar pipeline dengan pytest tanpa orkestrator.
  • Di produksi, jalankan dagster-webserver dan dagster-daemon; Dagster+ menawarkan control plane terkelola jika Anda lebih memilih tidak mengoperasikannya sendiri.

Mulailah dari kecil dengan satu atau dua asset, jalankan dagster dev, lalu kembangkan grafnya dari sana. Model asset menskala dengan mulus dari prototipe lokal hingga pipeline warehouse produksi.

Artikel Terkait

Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook

Mage: Bikin Data Pipeline Modern yang Rasanya Kayak Ngoding di Notebook Halo temen-temen, ketemu lagi sama aku Ruby Abdu...

Tutorial dlt: Pipeline Ingestion Data Berbasis Python

Membangun Pipeline EL Berbasis Python dengan dlt (data load tool) Sebagian besar tim data menghabiskan waktu yang tidak ...

Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk ML

Tutorial Lengkap Prefect: Modern Workflow Orchestration untuk ML Prefect adalah platform workflow orchestration modern y...

Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines

Tutorial Lengkap Apache Airflow: Workflow Orchestration untuk Data Pipelines Apache Airflow adalah platform open-source ...