Data-aware scheduling met Assets

Data-pijplijnen bouwen met Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Het timingprobleem

 

  • Pipeline A laadt ruwe data uit een API
  • Pipeline B bouwt op die data een dashboard
  • Als beide op cron draaien, kan B klaar zijn vóór A
  • Een tijdmarge toevoegen is fragiel

Het timingprobleem

Data-pijplijnen bouwen met Airflow

Wat is een Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Een verwijzing naar een datapunt, met een unieke naam
  • Een URI kan aan de asset worden gekoppeld als het een concrete data-entiteit is
  • Kan een bestand, een databasetabel of andere data zijn
Data-pijplijnen bouwen met Airflow

Producers en consumers

  • Taken kunnen assets bijwerken; dat creëert een asset event
  • Dags kun je op asset-updates triggeren met een asset schedule

 

Asset-flow

Data-pijplijnen bouwen met Airflow

Producer: signaleren met outlets

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

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

 

  • Bij succes registreert Airflow dat de asset is bijgewerkt
  • De producer hoeft niets te weten van een consumer
Data-pijplijnen bouwen met Airflow

Consumer: plannen op een asset

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

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

 

  • Geen cron-expressie, alleen de asset-referentie
  • Airflow triggert de consumer automatisch na de update
Data-pijplijnen bouwen met Airflow

Voorwaardelijk plannen

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

# Wacht tot BEIDE assets zijn bijgewerkt @dag(schedule=(sales & inventory)) def full_report(): ...
# Trigger wanneer ÉÉN van de assets bijwerkt @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Data-pijplijnen bouwen met Airflow

Assets controleren met de 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 toont alle geregistreerde assets
  • airflow assets details toont metadata, inclusief updated_at
Data-pijplijnen bouwen met Airflow

Laten we oefenen!

Data-pijplijnen bouwen met Airflow

Preparing Video For Download...