Harmonogramowanie oparte na zasobach

Budowanie potoków danych z Airflow

Volker Janz

Senior Developer Advocate at Astronomer

Problem z synchronizacją czasową

 

  • Pipeline A ładuje surowe dane z API
  • Pipeline B buduje dashboard na podstawie tych danych
  • Jeśli oba działają na cronie, B może ruszyć przed zakończeniem A
  • Dodanie bufora czasowego to rozwiązanie kruche

Problem z synchronizacją czasową

Budowanie potoków danych z Airflow

Czym jest zasób (Asset)?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Odwołanie do fragmentu danych, identyfikowane przez unikalną nazwę
  • Do zasobu można dołączyć URI, gdy reprezentuje konkretny obiekt danych
  • Może to być plik, tabela w bazie danych lub dowolne dane
Budowanie potoków danych z Airflow

Producenci i konsumenci

  • Taski mogą aktualizować zasoby, co tworzy zdarzenie zasobu
  • DAG-i można wyzwalać przy aktualizacjach zasobów za pomocą harmonogramu opartego na zasobach

 

Przepływ zasobów

Budowanie potoków danych z Airflow

Producent: sygnalizowanie przez outlets

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

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

 

  • Po pomyślnym wykonaniu Airflow rejestruje aktualizację zasobu
  • Producent nie musi wiedzieć nic o konsumentach
Budowanie potoków danych z Airflow

Konsument: harmonogramowanie na podstawie zasobu

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

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

 

  • Zamiast wyrażenia cron — tylko odwołanie do zasobu
  • Airflow wyzwala konsumenta automatycznie po aktualizacji
Budowanie potoków danych z Airflow

Harmonogramowanie warunkowe

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

# Czekaj na aktualizację OBU zasobów @dag(schedule=(sales & inventory)) def full_report(): ...
# Wyzwól, gdy zaktualizuje się DOWOLNY zasób @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Budowanie potoków danych z Airflow

Weryfikacja zasobów przez 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 wyświetla wszystkie zarejestrowane zasoby
  • airflow assets details pokazuje metadane, w tym updated_at
Budowanie potoków danych z Airflow

Czas na praktykę!

Budowanie potoków danych z Airflow

Preparing Video For Download...