Opas

Dagster, GA4 ja DBT

DBT muuntaa datan. Dagster vastaa siihen, mitä ajetaan, missä järjestyksessä, milloin ja mitä tehdään kun jokin menee rikki.

Päivitetty 2026-09-02 · Lukuaika 13 min

Mitä Dagsterilla tehdään

Dagster on orkestraattori: se ajaa dataputken osat oikeassa järjestyksessä, seuraa mitä onnistui ja mitä ei, ja tarjoaa käyttöliittymän jossa putken tila on nähtävissä. Samaa tehtävää hoitavat Airflow ja yksinkertaisimmillaan pelkkä ajastin.

Dagsterin ero muihin on siinä, mitä se pitää perusyksikkönä. Perinteinen orkestraattori ajaa tehtäviä: "aja tämä skripti, sitten tuo". Dagster kuvaa assetteja eli lopputuloksia: "tämä taulu on olemassa, se syntyy näin, ja se riippuu näistä". Ero kuulostaa akateemiselta, mutta se muuttaa arjen kysymyksen muodon.

Tehtäväpohjaisessa maailmassa kysyt "ajoiko yön putki läpi". Assettipohjaisessa kysyt "onko tämä taulu ajan tasalla, ja jos ei, mikä sen yläpuolella on rikki". Jälkimmäinen on se kysymys, joka oikeasti esitetään silloin kun raportti näyttää väärältä.

Tarvitsetteko Dagsteria?

Sanotaan tämä ensin, koska se säästää eniten rahaa: useimmat GA4-raportointiputket eivät tarvitse Dagsteria. Jos putki on "GA4-vienti saapuu, DBT ajaa, Data Studio lukee", niin Cloud Run -job ja Cloud Scheduler hoitavat sen muutamalla eurolla kuukaudessa ja ilman ylläpidettävää palvelinta.

Dagster alkaa maksaa itsensä takaisin, kun jokin näistä pitää paikkansa:

  • Lähteitä on useita ja niillä on eri aikataulut — GA4-vienti aamulla, CRM-poiminta tunneittain, mainoskulut omalla viiveellään.
  • Osien välillä on aitoja riippuvuuksia: aggregaattia ei ole mitään järkeä laskea ennen kuin kaikki sen lähteet ovat kunnossa.
  • Historiaa joudutaan laskemaan uudelleen: määritelmä muuttui, ja kuusi kuukautta pitää ajaa läpi hallitusti.
  • Joku muu kuin putken rakentaja tarvitsee näkymän siihen, mikä on ajan tasalla ja mikä ei.
  • Vikatilanteessa pitää tietää mikä osa epäonnistui ja voiko sen ajaa uudelleen yksin, ilman että koko yö ajetaan alusta.

Jos yksikään ei päde, jättäkää Dagster väliin ja lukekaa sen sijaan DBT-oppaamme — siellä kuvattu Cloud Run -ratkaisu riittää. Orkestraattorin tuominen liian aikaisin on tavallinen tapa vaihtaa yksinkertainen ongelma monimutkaiseen.

GA4-tapahtumat partitioituna assettina

GA4:n vienti on päiväkohtainen: yksi taulu per vuorokausi. Se sopii Dagsterin partitiointimalliin lähes täydellisesti. Määrittelet assetin, jolla on päiväpartitiot, ja jokainen partitio vastaa yhtä GA4-shardia.

Tämä ratkaisee kertaheitolla sen ongelman, jota DBT-putkessa kierretään liukuvalla uudelleenlaskentaikkunalla: kun partitio on nimetty, sen voi ajaa uudelleen yksin, milloin tahansa, ja Dagster tietää mitkä partitiot on ajettu ja mitkä ei.

assets/ga4.py
from dagster import asset, AssetExecutionContext, DailyPartitionsDefinition
from dagster_gcp import BigQueryResource

daily = DailyPartitionsDefinition(start_date="2026-01-01")


@asset(partitions_def=daily, group_name="ga4", kinds={"bigquery"})
def stg_ga4_events(
    context: AssetExecutionContext,
    bigquery: BigQueryResource,
) -> None:
    """GA4:n raakatapahtumat litteäksi tauluksi, yksi päivä kerrallaan."""
    day = context.partition_key              # "2026-09-01"
    shard = day.replace("-", "")             # "20260901"

    sql = f"""
    begin transaction;

    delete from `analytics.stg_ga4_events`
     where event_date = date '{day}';

    insert into `analytics.stg_ga4_events`
    select
      date '{day}'                                  as event_date,
      timestamp_micros(event_timestamp)             as event_at,
      event_name,
      user_pseudo_id,
      user_id,
      (select value.int_value from unnest(event_params)
         where key = 'ga_session_id')               as ga_session_id,
      (select value.string_value from unnest(event_params)
         where key = 'page_location')               as page_location,
      (select value.string_value from unnest(event_params)
         where key = 'source')                      as source,
      (select value.string_value from unnest(event_params)
         where key = 'medium')                      as medium,
      ecommerce.transaction_id,
      ecommerce.purchase_revenue
    from `projekti.analytics_123456789.events_*`
    where _table_suffix = '{shard}';

    commit transaction;
    """

    with bigquery.get_client() as client:
        job = client.query(sql)
        job.result()

    context.add_output_metadata({"shard": shard, "bytes_billed": job.total_bytes_billed})

Delete ja insert samassa transaktiossa tekee ajosta idempotentin: saman partition ajaminen uudelleen tuottaa saman lopputuloksen eikä kahdenna rivejä. Se on partitioidun assetin koko idea — uudelleenajo on turvallinen operaatio, ei riski.

Metadatan kirjaaminen kannattaa tehdä heti alusta. Kun ajon kustannus näkyy Dagsterin käyttöliittymässä partitiokohtaisesti, kalliin mallin löytäminen on minuutin työ eikä laskutusraportin arkeologiaa.

Miten DBT liittyy Dagsteriin

Tämä on se kohta, jossa yhdistelmä alkaa kannattaa. dagster-dbt-kirjasto lukee DBT-projektin manifestin ja luo jokaisesta DBT-mallista oman Dagster-assettinsa. DBT:n sisäinen riippuvuusgraafi ei siis jää DBT:n sisään, vaan sulautuu samaan graafiin muun putken kanssa.

Lopputulos on yksi lineage-näkymä, joka alkaa GA4:n raakaviennistä ja päättyy siihen aggregaattitauluun, jota Data Studio lukee. Kun aggregaatti on väärin, näet yhdellä silmäyksellä mikä sen yläpuolella epäonnistui — riippumatta siitä oliko rikkoutunut osa Python-assetti vai DBT-malli.

assets/dbt.py
from pathlib import Path

from dagster import AssetExecutionContext, AssetKey
from dagster_dbt import (
    DagsterDbtTranslator,
    DbtCliResource,
    DbtProject,
    dbt_assets,
)

dbt_project = DbtProject(
    project_dir=Path(__file__).joinpath("..", "..", "dbt_projekti").resolve()
)
dbt_project.prepare_if_dev()


class Ga4Translator(DagsterDbtTranslator):
    """Nimeää DBT:n lähteet samoiksi asseteiksi kuin Python-puolen assetit.

    Ilman tätä DBT:n source 'analytics.stg_ga4_events' saisi eri asset-avaimen
    kuin yllä määritelty stg_ga4_events, ja graafi katkeaisi juuri siitä
    kohdasta, jonka takia se rakennettiin.
    """

    def get_asset_key(self, dbt_resource_props):
        if dbt_resource_props["resource_type"] == "source":
            return AssetKey(dbt_resource_props["name"])
        return super().get_asset_key(dbt_resource_props)


@dbt_assets(
    manifest=dbt_project.manifest_path,
    dagster_dbt_translator=Ga4Translator(),
)
def dbt_analytics_assets(context: AssetExecutionContext, dbt: DbtCliResource):
    yield from dbt.cli(["build"], context=context).stream()

Huomaa että dbt build ajaa myös DBT:n omat testit. Ne raportoituvat Dagsteriin osana ajoa, joten rikkoutunut oletus näkyy samassa näkymässä kuin kaikki muukin — ei erillisessä lokissa, jota kukaan ei lue.

Työnjako pysyy selvänä: DBT vastaa siitä miten data muunnetaan, Dagster siitä milloin ja missä järjestyksessä se tehdään. Kumpikaan ei yritä tehdä toisen työtä.

Datan yhdistäminen muista lähteistä

Jokainen lähde on oma assettinsa omalla aikataulullaan. Mainoskulut noudetaan rajapinnasta, CRM-data tietokannasta, budjetit ehkä Sheetsistä. Dagster ei vaadi niitä samaan tahtiin — se vaatii vain, että riippuvuudet on kerrottu.

Riippuvuus ilmaistaan yksinkertaisesti: kun assetti ottaa toisen assetin parametrikseen tai listaa sen deps-määrittelyssä, Dagster tietää järjestyksen eikä aja yhdistävää mallia ennen kuin sen lähteet ovat valmiit.

assets/sources.py
from dagster import AssetExecutionContext, asset
from dagster_gcp import BigQueryResource


@asset(partitions_def=daily, group_name="ads", kinds={"bigquery"})
def stg_google_ads_cost(
    context: AssetExecutionContext,
    bigquery: BigQueryResource,
) -> None:
    """Mainoskulut päivätasolla. Oma lähde, oma viive."""
    ...


@asset(partitions_def=daily, group_name="crm", kinds={"bigquery"})
def stg_crm_orders(
    context: AssetExecutionContext,
    bigquery: BigQueryResource,
) -> None:
    """Tilaukset CRM:stä. Sisältää peruutukset ja katteen, joita GA4 ei tiedä."""
    ...
DBT-malli, joka yhdistää kolme lähdettä
-- models/marts/agg_channel_performance.sql
-- Dagster tietää riippuvuudet ref()- ja source()-kutsuista, joten tämä
-- malli ajetaan vasta kun kaikki kolme lähdettä ovat kunnossa.

with sessions as (
    select * from {{ ref('fct_sessions') }}
),
orders as (
    select * from {{ source('analytics', 'stg_crm_orders') }}
),
cost as (
    select * from {{ source('analytics', 'stg_google_ads_cost') }}
)

select
  s.session_date                         as date,
  s.channel_group,
  count(*)                               as sessions,
  count(distinct o.order_id)             as orders,
  sum(o.margin)                          as margin,
  max(c.ad_cost)                         as ad_cost,
  safe_divide(max(c.ad_cost), sum(o.margin)) as cost_margin_ratio
from sessions s
left join orders o on s.transaction_id = o.order_id
left join cost   c on s.session_date = c.date
                  and s.channel_group = c.channel_group
group by 1, 2

Käytännön hyöty näkyy vikatilanteessa. Jos mainoskulujen rajapinta on nurin, Dagster ajaa GA4-osuuden normaalisti ja jättää vain kuluista riippuvat mallit ajamatta. Ilman riippuvuuksien kuvausta vaihtoehtoja olisi kaksi, ja molemmat huonoja: koko yö kaatuu, tai aggregaatti lasketaan puuttuvilla kuluilla ja raportti näyttää mainonnan ilmaiselta.

Aggregointi Data Studiota varten

Aggregointi itsessään tehdään DBT:llä täsmälleen kuten ilman Dagsteria: valmiiksi lasketut päivätaulut, partitiointi päivämäärän mukaan, klusterointi yleisimmän suodattimen mukaan. Raportti lukee tuhansia rivejä miljoonien sijaan.

Dagster lisää tähän kaksi asiaa. Ensinnäkin aggregaatti ajetaan vasta kun sen lähteet ovat valmiit, joten Data Studio ei koskaan näytä puoliksi laskettua päivää. Toiseksi, kun aggregaatti on assetti, sen tuoreus on mitattavissa ja näkyvissä — kysymykseen "milloin tämä luku on viimeksi päivittynyt" on vastaus, eikä se ole arvaus.

Käytännön nyrkkisääntö pysyy samana kuin ilman orkestraattoria: ota aggregaattiin mukaan ne ulottuvuudet, joilla raportissa oikeasti suodatetaan. Uusi aggregaattitaulu on halvempi kuin yksi hidas yleistaulu — ja Dagsterissa uuden lisääminen on yhden mallin ja yhden riippuvuuden asia.

Missä Dagster ja DBT ajetaan

Tässä on Dagsterin todellinen hinta, ja se kannattaa tietää ennen päätöstä. Cloud Run -job herää, ajaa DBT:n ja sammuu. Dagster ei toimi niin: sensorit ja ajastukset vaativat jatkuvasti käynnissä olevan daemon-prosessin, ja käyttöliittymä oman palvelimensa. Ajojen tila tarvitsee lisäksi Postgres-tietokannan.

Vaihtoehtoja on käytännössä kolme:

  • Dagster+ hybridimallilla — hallintapaneeli on Dagsterin pilvipalvelussa, mutta koodi ja data pysyvät teidän GCP-projektissanne agentin kautta. Vähiten ylläpitoa, kuukausimaksu. Useimmille tämä on oikea valinta, jos Dagsteriin ylipäätään mennään.
  • Itse hostattu GKE:ssä — täysi kontrolli, ei lisenssikuluja, mutta klusteri, daemon, webserver ja Postgres ovat teidän ylläpidettävinänne. Perusteltua jos GKE on jo käytössä ja osaamista löytyy.
  • Itse hostattu Cloud Runissa — mahdollista, mutta daemon vaatii palvelun jossa min-instances on vähintään yksi. Silloin maksatte jatkuvasti käynnissä olevasta instanssista, ja suuri osa Cloud Runin kustannusedusta katoaa.
definitions.py — kaikki kootaan yhteen
from dagster import Definitions, define_asset_job, AssetSelection
from dagster_dbt import DbtCliResource
from dagster_gcp import BigQueryResource

from .assets.ga4 import stg_ga4_events
from .assets.sources import stg_google_ads_cost, stg_crm_orders
from .assets.dbt import dbt_analytics_assets, dbt_project
from .checks import ga4_event_count_matches_raw
from .sensors import ga4_export_sensor

ga4_daily_job = define_asset_job(
    name="ga4_daily",
    selection=AssetSelection.all(),
    partitions_def=daily,
)

defs = Definitions(
    assets=[
        stg_ga4_events,
        stg_google_ads_cost,
        stg_crm_orders,
        dbt_analytics_assets,
    ],
    asset_checks=[ga4_event_count_matches_raw],
    jobs=[ga4_daily_job],
    sensors=[ga4_export_sensor],
    resources={
        "bigquery": BigQueryResource(project="projekti"),
        "dbt": DbtCliResource(project_dir=dbt_project),
    },
)

Käyttöoikeuksista sama sääntö kuin muuallakin: palvelutili ja Workload Identity, ei ladattua avaintiedostoa. Dagsterin daemon pyörii jatkuvasti, joten sen oikeudet ovat pysyvästi käytössä — sitä suuremmalla syyllä ne kannattaa rajata vain niihin datasetteihin, joita putki oikeasti koskee.

Milloin ajo käynnistyy: sensori, ei kellonaika

DBT-oppaassa todettiin, että kiinteään kellonaikaan sidottu ajo epäonnistuu hiljaa niinä päivinä, joina GA4:n vienti myöhästyy. Dagsterissa tähän on suora ratkaisu: sensori, joka tarkistaa onko vienti saapunut, ja käynnistää ajon vasta kun se on.

Sensori ajetaan säännöllisin väliajoin, mutta se ei käynnistä mitään ennen kuin ehto täyttyy. Jos vienti myöhästyy neljä tuntia, ajo alkaa neljä tuntia myöhemmin — ei jää tekemättä.

sensors.py
from datetime import date, timedelta

from dagster import RunRequest, SkipReason, sensor
from dagster_gcp import BigQueryResource


@sensor(job=ga4_daily_job, minimum_interval_seconds=1800)
def ga4_export_sensor(context, bigquery: BigQueryResource):
    """Käynnistää ajon vasta kun eilisen GA4-shard on olemassa."""
    day = date.today() - timedelta(days=1)
    shard = day.strftime("%Y%m%d")

    sql = f"""
    select 1
      from `projekti.analytics_123456789.INFORMATION_SCHEMA.TABLES`
     where table_name = 'events_{shard}'
    """
    with bigquery.get_client() as client:
        found = list(client.query(sql).result())

    if not found:
        return SkipReason(f"GA4-vientiä {shard} ei ole vielä saapunut.")

    # run_key tekee tästä idempotentin: sama päivä käynnistyy vain kerran,
    # vaikka sensori ehtisi nähdä taulun useasti.
    return RunRequest(run_key=f"ga4-{shard}", partition_key=day.isoformat())

run_key on tässä olennainen. Ilman sitä sensori käynnistäisi saman päivän ajon uudelleen joka kerta kun se herää. Sen kanssa Dagster tunnistaa jo käynnistetyn ajon eikä tee sitä toistamiseen.

Ajastettua ajoa kannattaa silti pitää sensorin rinnalla varmistuksena niille malleille, jotka eivät riipu GA4:stä. Sensori vastaa siitä, milloin GA4-haara alkaa; ajastus siitä, ettei muu putki jää kiinni yhdestä myöhästyvästä lähteestä.

Backfillit: uudelleenlaskenta ilman kikkailua

Tämä on yksittäisistä ominaisuuksista se, joka useimmin ratkaisee valinnan Dagsterin hyväksi. Kun mittarin määritelmä muuttuu, historia pitää laskea uudelleen. Ilman partitiointia se tarkoittaa käsin kirjoitettua skriptiä, joka ajaa päivät silmukassa ja jonka etenemistä seurataan lokista.

Partitioidussa assetissa sama asia on backfill: valitset partitiovälin, Dagster ajaa ne, ja käyttöliittymässä näkyy mitkä ovat valmiit, mikä on kesken ja mikä epäonnistui. Epäonnistuneet voi ajaa uudelleen yksin koskematta onnistuneisiin.

Sama koskee myöhässä saapuvaa dataa. DBT-oppaan liukuva kolmen vuorokauden ikkuna on hyvä ratkaisu ilman orkestraattoria, mutta se on silti kiertotie: se laskee joka yö uudelleen päiviä, jotka ovat lähes aina jo valmiita. Dagsterissa voi sen sijaan ajaa uudelleen täsmälleen ne partitiot, joita GA4 on korjannut.

Asset checkit: laadun valvonta samassa graafissa

DBT:n testit kattavat muunnokset. Asset checkit kattavat sen, mitä DBT ei näe — esimerkiksi täsmääkö mallinnettu tapahtumamäärä raakadataan. Tämä on se testi, joka huomaa hiljaa rivejä pudottavan purkulogiikan.

checks.py
from dagster import AssetCheckResult, asset_check
from dagster_gcp import BigQueryResource


@asset_check(asset=stg_ga4_events, blocking=True)
def ga4_event_count_matches_raw(context, bigquery: BigQueryResource):
    """Mallinnetun ja raakadatan rivimäärän on täsmättävä päivätasolla."""
    sql = """
    with raw as (
      select parse_date('%Y%m%d', event_date) as d, count(*) as n
        from `projekti.analytics_123456789.events_*`
       where _table_suffix = format_date('%Y%m%d', current_date() - 1)
       group by d
    ),
    modelled as (
      select event_date as d, count(*) as n
        from `analytics.stg_ga4_events`
       where event_date = current_date() - 1
       group by d
    )
    select raw.n as raw_n, modelled.n as modelled_n
      from raw join modelled using (d)
    """
    with bigquery.get_client() as client:
        row = list(client.query(sql).result())[0]

    return AssetCheckResult(
        passed=row.raw_n == row.modelled_n,
        metadata={"raakadata": row.raw_n, "mallinnettu": row.modelled_n},
    )

blocking=True tarkoittaa, ettei alavirran malleja ajeta jos tarkistus epäonnistuu. Se on oikea oletus raportointiputkessa: mieluummin eilinen luku raportissa kuin tämän päivän väärä luku.

Yleisimmät virheet

  • Dagster otetaan käyttöön putkeen, jossa on yksi lähde ja yksi ajo. Ylläpidettävää tulee lisää, hyötyä ei.
  • Assetit ilman partitiointia. Silloin uudelleenlaskenta on taas kaikki tai ei mitään, ja puolet Dagsterin hyödystä jää käyttämättä.
  • DBT ajetaan yhtenä läpinäkymättömänä komentona dagster-dbt-integraation sijaan. Graafiin jää musta laatikko juuri siihen kohtaan, jossa suurin osa logiikasta on.
  • Sensori ilman run_keytä: sama ajo käynnistyy uudelleen joka herätyksellä.
  • Ajo ei ole idempotentti. Partitio ajetaan uudelleen ja rivit kahdentuvat — jolloin uudelleenajo, koko mallin tärkein ominaisuus, muuttuu riskiksi.
  • Daemonin unohtaminen itse hostatussa ympäristössä. Ilman jatkuvasti käynnissä olevaa daemonia sensorit ja ajastukset eivät laukea, ja putki näyttää toimivalta koska mikään ei kaadu.
  • Asset checkit jätetään pois, koska DBT:ssä on jo testit. Ne kattavat eri asiat: DBT testaa muunnoksen, check testaa lähteen ja lopputuloksen suhteen.

Versioista

Dagsterin rajapinta on kehittynyt nopeasti, ja tämän oppaan esimerkit nojaavat vakiintuneisiin osiin: @asset, partitiot, dagster-dbt ja sensorit. Uudempia tapoja kuvata automaatiota ja projektirakennetta on tullut lisää, ja ne muuttuvat edelleen.

Käytännön suositus: kiinnitä versiot ja lue oman versiosi dokumentaatio ennen kuin kopioit esimerkkejä mistään — tämä opas mukaan lukien. Dagster-projektissa versionsa kertova requirements-tiedosto säästää enemmän aikaa kuin mikään yksittäinen ominaisuus.

Mietittekö tarvitsetteko orkestraattorin?

Käydään putkenne läpi ja sanomme suoraan kumpi riittää — ajastettu DBT vai orkestroitu kokonaisuus. Kartoituskeskustelu on maksuton.

Katso data engineering -palvelumme

Aloitetaan maksuttomalla kartoituksella

Kerro tilanteesi, niin sanomme suoraan mitä kannattaa tehdä, missä järjestyksessä ja mitä se maksaisi. Ensimmäinen keskustelu ei sido mihinkään.