Planificarea bazată pe date cu Assets

Construirea pipeline-urilor de date cu Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problema de sincronizare

 

  • Pipeline A încarcă date brute dintr-un API
  • Pipeline B construiește un dashboard din acele date
  • Dacă ambele rulează pe cron, B poate porni înainte ca A să termine
  • Adăugarea unui buffer de timp este fragilă

Problema de sincronizare

Construirea pipeline-urilor de date cu Airflow

Ce este un Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • O referință la o bucată de date, identificată printr-un nume unic
  • Un URI poate fi atașat asset-ului când reprezintă o entitate de date concretă
  • Poate fi un fișier, un tabel dintr-o bază de date sau orice date
Construirea pipeline-urilor de date cu Airflow

Producători și consumatori

  • Taskurile pot actualiza asset-uri, generând un eveniment de asset
  • DAG-urile pot fi declanșate la actualizarea asset-urilor printr-un asset schedule

 

Fluxul asset-urilor

Construirea pipeline-urilor de date cu Airflow

Producător: semnalizare prin outlets

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

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

 

  • La succes, Airflow înregistrează că asset-ul a fost actualizat
  • Producătorul nu trebuie să cunoască consumatorul
Construirea pipeline-urilor de date cu Airflow

Consumator: planificare pe un asset

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

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

 

  • Fără expresie cron – doar referința la asset
  • Airflow declanșează consumatorul automat după actualizare
Construirea pipeline-urilor de date cu Airflow

Planificare condiționată

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(): ...
Construirea pipeline-urilor de date cu Airflow

Verificarea asset-urilor cu 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 afișează toate asset-urile înregistrate
  • airflow assets details afișează metadate, inclusiv updated_at
Construirea pipeline-urilor de date cu Airflow

Hai să exersăm!

Construirea pipeline-urilor de date cu Airflow

Preparing Video For Download...