Görevler arasında veri paylaşma

Airflow ile Veri İş Hatları Oluşturma

Volker Janz

Senior Developer Advocate at Astronomer

XCom nedir?

XCom push pull flow

 

  • XCom = görevler arası iletişim
  • Değerler metadata veritabanında saklanır
  • Klasik yaklaşım: açık xcom_push() ve xcom_pull()
Airflow ile Veri İş Hatları Oluşturma

XComArg: TaskFlow yolu

 

  • Bir @task dönüş değeri bir XComArg olur
  • XComArg gerçek veri değil, tembel bir başvurudur
  • Bunu başka bir fonksiyona vermek bağımlılık yaratır
  • Çalışma anında Airflow, XCom'dan çekerek bunu çözer

XComArgs in TaskFlow

Airflow ile Veri İş Hatları Oluşturma

Kodda XComArg kullanımı

@task
def get_config():
    return {"name": "daily_etl"}

@task
def run_pipeline(config):
    print(config["name"])


config = get_config() # XComArg
run_pipeline(config) # bağımlılık + veri
Airflow ile Veri İş Hatları Oluşturma

Klasik ve TaskFlow'u karıştırma

get_timestamp = BashOperator(
    task_id="get_timestamp",
    bash_command="date +%Y-%m-%d")

@task
def log_timestamp(ts_value):
    print(f"Timestamp: {ts_value}")


log_timestamp(get_timestamp.output)
Airflow ile Veri İş Hatları Oluşturma

XCom kısıtları

XCom size limits by backend

 

  • Maksimum boyut, veritabanı arka ucuna bağlıdır
  • JSON serileştirilebilir olmalı (dict, liste, string, sayı)
  • Her değer veritabanına yük ekler
Airflow ile Veri İş Hatları Oluşturma

Özel XCom arka uçları

Custom XCom backend routing

  • Değerleri bir nesne depolamaya yönlendir
  • AIRFLOW__COMMON_IO__XCOM_OBJECTSTORAGE_THRESHOLD: yalnızca eşik aşılırsa XCom'ları nesne depolamada tut
Airflow ile Veri İş Hatları Oluşturma

Hadi pratik yapalım!

Airflow ile Veri İş Hatları Oluşturma

Preparing Video For Download...