Görev Gruplarıyla karmaşık Dag'leri düzenleme

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

Karmaşıklık sorunu

Karmaşık Dag

  • Büyük Dag'ler Grafik görünümünde gezmesi zor hale gelir
  • Görev adları birbirine karışır
  • Yeni ekip üyeleri kendi sorumluluklarını bulmakta zorlanır
Airflow ile Veri İş Hatları Oluşturma

@task_group

from airflow.sdk import dag, task, task_group

@task_group(
    group_id="ingest_orders",
    default_args={"retries": 3},
)

def process_orders(): @task def extract_orders(): return [{"id": 1, "amount": 99.99}] @task def transform_orders(orders): return [{"id": o["id"], "total": o["amount"] * 1.08} for o in orders] return transform_orders(extract_orders())
  • group_id bir özel tanımlayıcı ayarlar, varsayılanı fonksiyon adıdır
  • default_args, gruptaki tüm görevlere uygulanır; tekrar eden yapılandırma kodunu önler
  • Görev grupları UI'da genişletilebilir bloklar olarak görünür
Airflow ile Veri İş Hatları Oluşturma

Airflow UI'da nasıl görünür

Görev grupları olmadan

Bir hizada üç görevi gösteren basit grafik: extract_orders, transform_orders, load; hepsi aynı seviyede

  • Tüm görevler aynı seviyede

Görev gruplarıyla

process_orders adlı, içinde extract ve transform olan daraltılmış bir blok ve ardından grubun dışında bir load görevi gösteren grafik

  • UI'da daraltılabilir bloklar
Airflow ile Veri İş Hatları Oluşturma

İç içe kullanım ve özel görünen adlar

@task_group(group_display_name="Process All")
def process_all():

    @task_group(group_display_name="Ingest Orders")
    def orders():
        return transform(extract())

    @task_group(group_display_name="Process Returns")
    def returns():
        return transform(extract())

    return {
        "orders": orders(),
        "returns": returns(),
    }
  • Daha karmaşık veri hatları için görev grupları iç içe kullanılabilir
  • group_display_name, UI'da insan tarafından okunabilir bir etiket ayarlar (emoji de destekler)
Airflow ile Veri İş Hatları Oluşturma

Factory deseni

@task_group
def process_source(source_name, source_path):

    @task
    def extract():
        return read_data(source_path)

    @task
    def transform(data):
        return clean_data(data)

    return transform(extract())

# Farklı kaynaklar için aynı deseni yeniden kullan orders = process_source("orders", "/data/orders.csv") returns = process_source("returns", "/data/returns.csv") events = process_source("events", "/data/events.csv")
  • Görev grupları dekoratörlü Python fonksiyonlarıdır
  • Farklı parametrelerle birden çok kez çağırabilirsin
  • Bu factory deseni, boru hattı mantığını yeniden kullanmanı sağlar
Airflow ile Veri İş Hatları Oluşturma

Gruplama için yönergeler

Alana ya da konuya göre grupla

  • Alan ya da konuya göre grupla; operatör tipine göre değil
  • Net programatik adlar için group_id, UI etiketleri için group_display_name kullan
  • Bir gruptaki tüm görevlerde yeniden denemeler gibi yapılandırmayı paylaşmak için default_args kullan
  • Miller Yasası: 7'den fazla üst düzey öğe, bir görev grubuna ihtiyacın olduğunun işareti olabilir
Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...