Datadriven schemaläggning med Assets

Bygg datapipelines med Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Timingproblemet

 

  • Pipeline A läser in rådata från ett API
  • Pipeline B bygger en dashboard från dessa data
  • Kör båda med cron kan B starta innan A är klar
  • En tidsbuffert är en skör lösning

Timingproblemet

Bygg datapipelines med Airflow

Vad är en Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • En referens till ett dataobjekt, identifierat med ett unikt namn
  • En URI kan kopplas till en asset när den representerar en konkret dataentitet
  • Kan vara en fil, en databastabell eller annan data
Bygg datapipelines med Airflow

Producenter och konsumenter

  • Uppgifter kan uppdatera assets, vilket skapar en assethändelse
  • DAG:ar kan utlösas vid assetuppdateringar med ett asset-schema

 

Asset-flöde

Bygg datapipelines med Airflow

Producent: signalering med outlets

sales_data = Asset("s3://bucket/sales/daily.parquet")

@task(outlets=[sales_data]) def write_sales(): # Write data to S3 ...

 

  • Vid lyckad körning registrerar Airflow att asseten uppdaterades
  • Producenten behöver inte känna till någon konsument
Bygg datapipelines med Airflow

Konsument: schemaläggning på en asset

sales_data = Asset("s3://bucket/sales/daily.parquet")

@dag(schedule=[sales_data]) def build_dashboard(): ...

 

  • Inget cron-uttryck – bara assetreferensen
  • Airflow utlöser konsumenten automatiskt efter uppdateringen
Bygg datapipelines med Airflow

Villkorsstyrd schemaläggning

sales = Asset("s3://bucket/sales.parquet")
inventory = Asset("s3://bucket/inventory.parquet")

# Wait for BOTH assets to update @dag(schedule=(sales & inventory)) def full_report(): ...
# Trigger when ANY asset updates @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Bygg datapipelines med Airflow

Verifiera assets med CLI

$ airflow assets list
name                                  | group    | uri
s3://data-lake/sales/daily.csv        | asset    | s3://data-lake/sales/daily.csv
$ airflow assets details --name "s3://data-lake/sales/daily.csv"
property_name    | property_value
=================+=====================                                         
name             | s3://data-lake/sales/daily.csv
                ...
updated_at       | 2026-04-14T10:00:50.820041Z
  • airflow assets list visar alla registrerade assets
  • airflow assets details visar metadata inklusive updated_at
Bygg datapipelines med Airflow

Nu kör vi en övning!

Bygg datapipelines med Airflow

Preparing Video For Download...