Plánování na základě dat (Assets)

Tvorba datových pipeline s Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problém s časováním

 

  • Pipeline A načítá surová data z API
  • Pipeline B sestavuje dashboard z těchto dat
  • Pokud obě běží na cronu, B může spustit dřív, než A skončí
  • Přidání časové rezervy je křehké řešení

Problém s časováním

Tvorba datových pipeline s Airflow

Co je Asset?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Odkaz na kus dat identifikovaný unikátním názvem
  • Pokud asset reprezentuje konkrétní datovou entitu, lze mu přiřadit URI
  • Může jít o soubor, databázovou tabulku nebo libovolná data
Tvorba datových pipeline s Airflow

Producenti a konzumenti

  • Tasky mohou aktualizovat assety, čímž vzniká asset event
  • DAGy lze spouštět při aktualizaci assetů pomocí asset schedule

 

Tok assetů

Tvorba datových pipeline s Airflow

Producent: signalizace pomocí outlets

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

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

 

  • Po úspěšném dokončení Airflow zaznamená, že byl asset aktualizován
  • Producent nemusí vědět nic o konzumentech
Tvorba datových pipeline s Airflow

Konzument: plánování na základě assetu

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

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

 

  • Žádný cron výraz, jen odkaz na asset
  • Airflow spustí konzumenta automaticky po aktualizaci
Tvorba datových pipeline s Airflow

Podmíněné plánování

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(): ...
Tvorba datových pipeline s Airflow

Ověření assetů pomocí 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 zobrazí všechny registrované assety
  • airflow assets details zobrazí metadata včetně updated_at
Tvorba datových pipeline s Airflow

Pojďme cvičit!

Tvorba datových pipeline s Airflow

Preparing Video For Download...