Modern veri mühendisliği ortamında güvenilirlik her şeyden önemlidir. Birçok orkestrasyon aracı yalnızca görev yürütmesine odaklanırken, Kestra iş akışı durumunu ve veri kökenini birinci sınıf varlıklar olarak ele alarak kendini ayırt eder. Üretim seviyesi Veri Topla, Dönüştür ve Yükle (ETY) hatları için verinin sisteminizden nasıl aktığını anlamak ve iş akışı yürütme durumunu yönetmek, hata ayıklama, denetim ve veri bütünlüğünü koruma açısından kritik öneme sahiptir. Bu yazı, dayanıklı, gözlemlenebilir ve durum farkında ETY hatları oluşturmak için Kestra'nın yerleşik yeteneklerinden nasıl yararlanılacağını incelemektedir.
Temel Zorluk: Dağıtık Sistemlerde Durum ve Köken
Geleneksel ETY hatları genellikle "siyah kutu" sendromundan muzdarip olur. Bir iş başarısız olduğunda, mühendisler sorunun kaynak veride mi, dönüştürme mantığında mı yoksa hedef yapılandırmasında mı olduğunu belirlemek için mücadele eder. Ayrıca, net bir köken olmadan etki analizi manuel ve hataya yatkın bir süreç haline gelir. Kestra, görevler arasında her girdi ve çıktı değişkenini izlemek için yerleşik destek sağlayarak bu zorlukları ele alır; aynı zamanda iş akışı yürütme durumunun yaşam döngüsünü de yönetir.
Tam Granüler Veri Kökeni Uygulama
Kestra, değişkenlerin kökenini otomatik olarak izler. Ancak, bu özelliği en üst düzeye çıkarmak için akışlarınızı açık girdi ve çıktı tanımlarıyla tasarlamalısınız. Veritabanından veri çıkarıp Python kullanarak dönüştürdüğünüz ve bir veri ambarına yüklediğiniz tipik bir E senaryosunu ele alalım.
Kestra'nın triggers (tetikleyicileri) ve task outputs (görev çıktıları) kullanarak verinin hatlarımızdan nasıl aktığını tam olarak görselleştirebiliriz. Açık değişken aktarımını gösteren pratik bir akış örneği aşağıdadır:
id: production_etl_pipeline
namespace: company.data
tasks:
- id: extract
type: io.kestra.plugin.jdbc.postgresql.Query
url: "{{ secrets.POSTGRES_URL }}"
sql: "SELECT * FROM raw_events WHERE created_at > '{{ prevRunDate.date('yyyy-MM-dd') }}'"
- id: transform
type: io.kestra.plugin.scripts.python.Script
runner: DOCKER
script: |
import json
import sys
# Girdi verisini çıkar
data = json.loads(inputs.raw_data)
# Dönüştürme mantığı
cleaned = [x for x in data if x['value'] is not None]
# Köken izleme için yeni bir değişken olarak çıktı
outputs.cleaned_data = json.dumps(cleaned)
inputs:
raw_data: "{{ taskrun.extract.outputs.rows }}"
- id: load
type: io.kestra.plugin.jdbc.postgresql.BulkOutput
url: "{{ secrets.POSTGRES_URL }}"
schema: public
table: cleaned_events
columns:
- id
- value
- created_at
from: "{{ outputs.transform.cleaned_data }}"
Bu örnekte Kestra, extract görevinden transform görevine ve nihayetinde load görevine kadar kökeni otomatik olarak izler. Bu sayede kullanıcı arayüzünde tıklayarak belirli bir çıktıyı etkileyen satırların tam olarak hangileri olduğunu görebilirsiniz; bu da veri tutarsızlıklarını hata ayıklamak için paha biçilmezdir.
Gelişmiş Durum Yönetimi ve Yeniden Deneme Mantığı
Üretim hatları geçici başarısızlıklara maruz kalır: ağ zaman aşımı, API hız sınırları veya geçici veritabanı kilitleri. Kestra'nın durum yönetimi, ana iş akışı mantığını karmaşıklaştırmadan granüler yeniden deneme politikaları ve hata işleme stratejileri tanımlamanıza olanak tanır.
Durum kalıcılığı için Airflow gibi harici araçlara güvenmek yerine, Kestra yürütme durumunu yerleşik depolama arka ucunda (Postgres veya S3) saklar. Başarısızlıkları zarifçe ele almak için errors bloklarından yararlanabilirsiniz:
- id: load_retry
type: io.kestra.plugin.jdbc.postgresql.BulkOutput
# ... yapılandırma ...
retry:
type: FIXED
interval: 30s
maxDuration: 5m
maxAttempt: 3
errors:
- type: io.kestra.plugin.core.log.Log
message: "Son yükleme denemesi başarısız oldu. Slack'e bildirim gönderiliyor."
Bu yaklaşım, hattınızın varsayılan olarak dayanıklı olmasını sağlar. Durum atomik olarak güncellenir; yani bir iş akışını yürütme sırasında duraklatabilir, devam ettirebilir veya sonlandırabilirsiniz ve Kestra, neyin işlendiğine dair tutarlı bir kayıt tutmaya devam eder.
Sonuç
Kestra'da karmaşık veri kökeni ve durum yönetimi uygulamak, ETY hatlarınızı kırılgan betiklerden sağlam, kurumsal seviyeli sistemlere dönüştürür. Girdileri açıkça tanımlayarak, otomatik köken izlemeyi kullanarak ve granüler yeniden deneme politikalarını yapılandırarak, veri altyapınıza olan güveni korumak için gereken gözlemlenebilirliğe sahip olursunuz. Veri ihtiyaçlarınız büyüdükçe, Kestra'nın bildirimsel YAML yaklaşımı ve güçlü eklenti ekosistemi, iş akışlarınızın sürdürülebilir ve ölçeklenebilir kalmasını sağlar.