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 raw
sales() -> 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, load
assetsfrommodules
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
@asset
check(asset="cleanedsales")
def no
negativeamounts(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 sales
pipeline 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 sales
pipeline.assets import cleanedsales, dailysalessummary
def test
cleanedsalesdropsnulls():
raw = pd.DataFrame(
{
"order
id": [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 sales
pipeline.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-webservermenyajikan UI dan API GraphQL yang dipakai klien untuk meluncurkan dan memeriksa run.dagster-daemonmenjalankan 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
ConfigurableResourcedenganEnvVaruntuk 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
materializetomemoryuntuk 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
Definitionsadalah 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
pytesttanpa orkestrator. - Di produksi, jalankan
dagster-webserverdandagster-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.