Varlıklarla veri farkındalıklı zamanlama

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

Zamanlama sorunu

 

  • Pipeline A, bir API'den ham veri yükler
  • Pipeline B, o veriden bir kontrol paneli oluşturur
  • İkisi de cron ile çalışırsa, B A bitmeden önce çalışabilir
  • Zaman tamponu eklemek kırılgandır

Zamanlama sorunu

Airflow ile Veri İş Hatları Oluşturma

Asset nedir?

from airflow.sdk import Asset
sales_data = Asset("s3://bucket/sales/daily.parquet")
  • Benzersiz bir adla tanımlanan bir veri parçasına referans
  • Somut bir veri varlığını temsil ediyorsa varlığa bir URI eklenebilir
  • Bir dosya, veritabanı tablosu ya da herhangi bir veri olabilir
Airflow ile Veri İş Hatları Oluşturma

Üreticiler ve tüketiciler

  • Görevler varlıkları güncelleyebilir; bu bir varlık olayı oluşturur
  • Dagg'ler, varlık zamanlaması ile varlık güncellemelerinde tetiklenebilir

 

Varlık akışı

Airflow ile Veri İş Hatları Oluşturma

Üretici: outlets ile sinyal verme

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

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

 

  • Başarıyla bittiğinde, Airflow varlığın güncellendiğini kaydeder
  • Üreticinin herhangi bir tüketiciyi bilmesi gerekmez
Airflow ile Veri İş Hatları Oluşturma

Tüketici: bir varlığa göre zamanlama

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

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

 

  • Cron ifadesi yok, sadece varlık referansı
  • Güncellemeden sonra Airflow tüketiciyi otomatik tetikler
Airflow ile Veri İş Hatları Oluşturma

Koşullu zamanlama

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

# İKİ varlığın da güncellenmesini bekle @dag(schedule=(sales & inventory)) def full_report(): ...
# HERHANGİ bir varlık güncellenince tetikle @dag(schedule=(sales | inventory)) def quick_refresh(): ...
Airflow ile Veri İş Hatları Oluşturma

CLI ile varlıkları doğrulama

$ 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 tüm kayıtlı varlıkları gösterir
  • airflow assets details, updated_at dahil metaveriyi gösterir
Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...