Apache Ecosystem

Apache Airflow'u Ustalaşmak: Sıfırdan Ölçeklenebilir Veri Boru Hatları Oluşturma

Veri hacimleri patlarken ve altyapılar giderek daha dağıtık hale gelirken, güvenilir, gözlemlenebilir ve genişletilebilir iş akışı orkestrasyonuna olan ihtiyaç hiç bu kadar büyük olmamıştı. Apache Airflow, bu zorluk için fiili endüstri standardı olarak ortaya çıkmıştır. İş akışlarını Python kodu kullanarak Yönlendirilmiş Döngüsüz Grafikler (DAG'lar) olarak ifade etmesine olanak tanıyan Airflow, veri mühendislerinin kırılgan kabuk betiklerinden (shell scripts) uzaklaşarak sürüm kontrolü yapılan, test edilmiş ve sürdürülebilir veri boru hatlarına yönelmesini sağlar.

Basit geceleyin toplu işleri (batch jobs) otomatikleştiriyor olun veya karmaşık gerçek zamanlı mikro hizmetleri orkestrasyon yapıyor olun, modern bir veri yığını için Airflow'un temel bileşenlerini (DAG'lar, operatörler, sensörler ve yürütme stratejileri) anlamak esastır. Bu kılavuzda, mimariyi inceleyecek, pratik kod örneklerini keşfedecek ve üretim dağıtımları için kritik hususları tartışacağız.

Temel Soyutlama: DAG'lar ve Zamanlama

Airflow'un kalbinde DAG (Yönlendirilmiş Döngüsüz Grafik) bulunur. Bir DAG, bağımlılıkları olan görevlerin, düğümler (görevler) ve kenarlar (bağımlılıklar) olarak temsil edildiği bir koleksiyondur. Statik yapılandırma dosyalarına dayanan geleneksel zamanlama sistemlerinin aksine, Airflow DAG'ları Python'da tanımlar. Bu, görevlerin dinamik olarak oluşturulmasına, koşullu mantığa ve kod tabanınız içinde doğrudan parametrelendirmeye olanak tanır.

Bir DAG'ı Airflow arayüzüne dağıttığınızda, zamanlayıcı motoru ~/.airflow/dags klasörünü tarar. Python dosyasını ayrıştırır, zamanlamayı (ör. cron tarzı) belirler ve her zamanlanmış aralık için "DAG Çalıştırma" (DAG Runs) oluşturur. start_date ve schedule_interval parametrelerini anlamak, veri boru hattınızın zamansal bağlamını tanımladıkları için kritiktir.

Temel Yapı Taşları: Operatörler, Sensörler ve Kancalar (Hooks)

Airflow'daki görevler Operatörler tarafından yürütülür. Bir operatör, boru hattındaki tek bir eylemi veya adımı tanımlayan bir Python sınıfıdır. Airflow, AWS, GCP, DBT, Kubernetes ve SQL gibi hizmetler için geniş bir hazır operatör ekosistemi sunar. Ancak belirli bir mantık gerektirildiğinde, özel operatörler yazabilir veya PythonOperator'ı kullanabilirsiniz.

Çoğu zaman, iş akışları tam olarak zamanlanmış saatte gerçekleşmeyebilecek dış olaylara bağlıdır. Sensörler, bir sonraki görevi tetiklemesinden önce bir koşul sağlanana kadar bir dış sistemi (S3 kovası veya bir veritabanı gibi) sorgulayarak (polling) bunu ele alır. Bu ayrıştırma, dağıtık ortamlarda veri tutarlılığını sağlar ve yarış durumlarını (race conditions) önler.

Kancalar (Hooks), kimlik doğrulama ve API çağrılarını işleyen, ardından daha yüksek düzeyde bir soyutlama için operatörler ve sensörler tarafından sarılan dış hizmetlere düşük düzeyde bir arayüz sağlar.

Pratik Örnek: Bir ETL Boru Hattı

Yaygın bir ETL senaryosunu düşünün: PostgreSQL veritabanından veri çıkarma, Pandas ile dönüştürme ve S3 veri gölüne yükleme.

from airflow import DAG
from airflow.providers.postgres.operators.postgres import PostgresOperator
from airflow.providers.apache.hive.operators.hive import HiveOperator
from airflow.operators.python import PythonOperator
from airflow.providers.aws.operators.s3 import S3Operator
from airflow.sensors.postgres import PostgresSensor
from datetime import datetime, timedelta

default_args = {
    'owner': 'data-team',
    'retries': 3,
    'retry_delay': timedelta(minutes=5)
}

def extract_custom_logic(**context):
    # Dönüşüm için özel Python mantığı
    data = context['task_instance'].xcom_pull(task_ids='extract_sql')
    processed_data = transform(data) 
    return processed_data

def load_to_s3(processed_data):
    # S3'e yükleme mantığı
    pass

with DAG(
    dag_id='etl_sales_pipeline',
    default_args=default_args,
    description='ETL pipeline for sales data',
    schedule_interval='@daily',
    start_date=datetime(2023, 1, 1),
    catchup=False,
    tags=['etl', 'sales'],
) as dag:

    wait_for_data = PostgresSensor(
        task_id='wait_for_new_data',
        sql="SELECT 1 FROM sales WHERE updated_at > now() - interval '1 day'",
        postgres_conn_id='prod_postgres',
        poke_interval=60,
    )

    extract = PostgresOperator(
        task_id='extract_sql',
        sql="SELECT * FROM sales WHERE updated_at > now() - interval '1 day'",
        postgres_conn_id='prod_postgres',
    )

    transform = PythonOperator(
        task_id='transform_data',
        python_callable=extract_custom_logic,
    )

    load = S3Operator(
        task_id='load_to_s3',
        s3_key='data/sales/etl_output.parquet',
        local_location='/tmp/sales_output.parquet',
    )

    wait_for_data >> extract >> transform >> load

Bu örnekte, veri kullanılabilirliğini sağlamak için yeniden denemeler ve sensör içeren sağlam bir boru hattı tanımlıyoruz. xcom_pull, Python operatörünün önceki SQL görevinden veri almasını sağlar ve durumun adımlar arasında nasıl aktarıldığını gösterir.

Üretim Dağıtımı En İyi Uygulamaları

Airflow'u üretime taşımak, sadece yazılımı yüklemekten fazlasını gerektirir. Temel en iyi uygulamalar şunları içerir:

  • Güçlerin Ayrılması: DAG tanımlarını, Airflow işçi (worker) örneklerinden ayrı olarak Git'te tutun.
  • Yürütücü Seçimi: Yatay ölçekleme için Celery veya Kubernetes yürütücüsünü kullanın ve üretim iş yükleri için Tek Süreç (Single Process) yürütücüsünden kaçının.
  • Bağlantı Yönetimi: Gizli bilgileri (secrets) ortam değişkenlerinde veya bir gizli bilgi yöneticisinde (AWS Secrets Manager gibi) saklayın ve kimlik bilgilerini sabit kodlamak yerine connections üzerinden referans verin.
  • İzleme: Görev süresini, hata oranlarını ve kaynak kullanımını izlemek için Airflow'u Prometheus, Grafana veya Datadog gibi izleme araçlarıyla entegre edin.
  • Test: Dağıtım öncesi DAG mantığını birim test etmek için pytest gibi çerçeveleri ve entegrasyon testleri için airflow dags test komutunu kullanın.

Sonuç

Apache Airflow, veri iş akışlarını orkestrasyon için güçlü, esnek ve genişletilebilir bir çerçeve sunar. Bildirimci (declarative) DAG sözdizimi, zengin operatör ekosistemi ve sağlam zamanlama yeteneklerinden yararlanarak veri mühendisliği ekipleri, yalnızca verimli değil, aynı zamanda dayanıklı ve sürdürülebilir boru hatları inşa edebilir. Veri altyapınız büyüdükçe, uygun Airflow tasarımı ve üretim sertleştirmesine (hardening) yapılan yatırımlar, güvenilirlik ve operasyonel verimlilik açısından karşılığını verecektir.

Share: