Pianificazione data-aware con gli Asset

Creare data pipeline con Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Il problema del timing

 

  • La pipeline A carica dati grezzi da un'API
  • La pipeline B crea una dashboard da quei dati
  • Se entrambe girano su cron, B può partire prima che A finisca
  • Aggiungere un margine di tempo è fragile

Il problema del timing

Creare data pipeline con Airflow

Che cos'è un Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Una riferimento a un dato, identificato da un nome univoco
  • Una URI può essere associata all'asset quando rappresenta un'entità dati concreta
  • Può essere un file, una tabella di database o qualsiasi dato
Creare data pipeline con Airflow

Produttori e consumatori

  • I task possono aggiornare asset, creando un asset event
  • I DAG possono essere attivati sugli aggiornamenti con un'asset schedule

 

Flusso degli asset

Creare data pipeline con Airflow

Produttore: segnalare con outlets

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

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

 

  • Al successo, Airflow registra che l'asset è stato aggiornato
  • Il produttore non deve conoscere alcun consumatore
Creare data pipeline con Airflow

Consumatore: pianificare su un asset

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

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

 

  • Nessuna espressione cron, solo il riferimento all'asset
  • Airflow avvia il consumatore automaticamente dopo l'aggiornamento
Creare data pipeline con Airflow

Pianificazione condizionale

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

# Attendi che ENTRAMBI gli asset si aggiornino @dag(schedule=(sales & inventory)) def full_report(): ...
# Attiva quando QUALSIASI asset si aggiorna @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Creare data pipeline con Airflow

Verifica degli asset con la 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 mostra tutti gli asset registrati
  • airflow assets details mostra i metadati, incluso updated_at
Creare data pipeline con Airflow

Esercitiamoci!

Creare data pipeline con Airflow

Preparing Video For Download...