Java ile NiFi Custom Processor Geliştirme

Java ile NiFi Custom Processor Geliştirme

Java ile geliştirilen bir NiFi custom processor için yaşam döngüsü, property yönetimi, session semantiği, hata yönlendirme ve thread güvenliğini açıklar. İşlemci tasarımı veri akışı ve backpressure bağlamında incelenir.

Apache NiFi içinde özel bir işlemci geliştirmek, "AbstractProcessor" sınıfını genişletip "onTrigger()" metoduna birkaç satır kod yazmaktan ibaret değildir. İşlemci, NiFi'nin zamanlama, transaction, içerik deposu, kuyruk baskısı, küme çalışması ve veri kökeni izleme mekanizmalarına katılır. Kodun doğruluğu kadar bu mekanizmalarla uyumu da önemlidir.

Bir FlowFile iki ayrı parçadan oluşur. Content gerçek veri gövdesini, attributes ise veriye ilişkin bağlamı taşır. "filename", "path" veya "absolute.path" gibi alanlar veri içeriği değildir. Bunlar yönlendirme ve depolama kararlarında kullanılan özniteliklerdir. Processor ise FlowFile'ı doğrudan değiştirmez. Bütün değişiklikler "ProcessSession" üzerinden yeni bir FlowFile sürümü oluşturarak yapılır.

Kurum içi veri akışlarında geliştirdiğim işlemcilerde en belirleyici konu, algoritmayı NiFi'nin çalışma modeli içine doğru yerleştirmekti. Kaynak dizin tarama, dosya seçme, içerik doğrulama, veritabanına aktarma ve arşivleme gibi işlemler tek bir metoda yığıldığında kod çalışsa bile hata sınırları belirsizleşiyordu. Sağlam tasarım, her processor'a tek ve ölçülebilir bir sorumluluk vermeyi gerektiriyordu.

Bileşen sözleşmesi

Custom Processor, kullanıcı arayüzüne yalnızca bir Java sınıfı olarak çıkmaz. Adı, açıklaması, özellikleri, ilişkileri, giriş gereksinimi ve çalışma davranışı birlikte bir bileşen sözleşmesi oluşturur.

"PropertyDescriptor" nesneleri işlemcinin yapılandırma alanlarını tanımlar. Kaynak dizin, zaman aşımı, toplu işlem boyutu, doğrulama seçeneği ve Controller Service bağlantıları burada belirtilir. Her özellik uygun validator ile sınırlandırılmalıdır. Tek tek geçerli görünen fakat birlikte kullanılamayan ayarlar için "customValidate()" gerekir. NiFi, yapılandırması geçersiz olan bir processor'ın çalıştırılmasına izin vermez.

İlişkiler, işlem sonucunu akışın sonraki katmanına taşır. "success" ve "failure" çoğu işlemci için yeterli görünse de hata sınıfları farklı davranış gerektiriyorsa ayrı ilişkiler daha uygundur. Örneğin tekrar denenebilir ağ hatası ile biçimsel olarak geçersiz veri aynı kuyruğa gönderilmemelidir. İlki gecikmeli yeniden denemeye, ikincisi karantina veya inceleme akışına yönlendirilebilir.

Kaynak processor'larda "@InputRequirement(INPUT_FORBIDDEN)" kullanmak, yanlışlıkla giriş bağlantısı oluşturulmasını engeller. Mevcut FlowFile üzerinde çalışan dönüştürücülerde ise "INPUT_REQUIRED" daha doğru sözleşmedir. NiFi bu anotasyonları yalnız belgelemek için değil, bileşenin geçerli olup olmadığını ve arayüzde hangi bağlantıların kurulabileceğini belirlemek için kullanır.

Processor ve bağımlılıkları NAR paketi içinde dağıtılır. NAR, bileşenleri ve bağımlılıklarını ayrı bir classloader altında yalıtır. Bu nedenle sıradan bir JAR dosyasını NiFi lib dizinine kopyalamak kalıcı bir eklenti mimarisi sayılmaz. Ortak Controller Service API'leri ayrı bir üst NAR içinde tutulmalı, implementasyon ve processor NAR'ları aynı API tanımına bağımlı olmalıdır. Aksi halde sınıf adı aynı görünse bile classloader kimliği farklı olduğu için tip uyuşmazlıkları oluşabilir.

Yaşam döngüsü

NiFi processor yaşam döngüsü, pahalı hazırlık işleri ile sıcak veri yolunu ayırmak için kullanılmalıdır.

"init()" veya "initialize()" bileşen oluşturulurken çalışır. Processor kimliği gibi yaşam boyunca değişmeyen bilgiler burada alınabilir. Çalışma ayarları ise henüz kesinleşmemiş olabilir. Yapılandırma değerlerini bu aşamada okumak doğru değildir.

"@OnScheduled", processor her çalıştırıldığında veya yeniden zamanlandığında çağrılır. Kaynak dizinin mutlak yolunu çözme, düzenli ifade derleme, sabit tablolar oluşturma, immutable dosya listesi hazırlama ve bağlantı havuzu nesnesini başlatma gibi işler burada yapılabilir. Framework, bu aşamada processor kodunu aynı anda çalıştıran başka bir thread bulunmamasını garanti eder. Hazırlanan yapı immutable bir nesneye dönüştürülüp "volatile" referans üzerinden yayımlandığında "onTrigger()" çağrıları kilitsiz biçimde okuyabilir.

Kaynak klasör, FlowFile'ın "absolute.path" özniteliğinden alınmamalıdır. Kaynak processor henüz FlowFile üretmemiştir. Klasör yolu bir property üzerinden okunmalı, normalize edilmeli ve gerekiyorsa mutlak yola çevrilmelidir. "filename" ile "relative.path" daha sonra üretilen FlowFile'ın kaynak içindeki konumunu taşır. Yapılandırma ile veri özniteliğini birbirine karıştırmak, yeniden başlatma ve farklı çalışma dizinlerinde belirsiz sonuç üretir.

"onTrigger()" asıl veri işleme yoludur. Burada uzun süren hazırlık, sınıf tarama veya her çağrıda değişmeyen yapıların yeniden oluşturulması gecikmeyi artırır. Kaynak processor yeni FlowFile üretir. Dönüştürücü processor ise "session.get()" ile giriş alır ve veri yoksa hemen döner.

"@OnUnscheduled" zamanlamanın kaldırılmasını, "@OnStopped" etkin çağrıların tamamlanmasından sonraki temizliği temsil eder. Soket, istemci, executor veya yerel havuz gibi kaynaklar burada kapatılmalıdır. "@OnShutdown" ek güvenlik sağlar, fakat işletim sistemi veya JVM zorla kapatıldığında çağrılacağı garanti edilemez. Kalıcı doğruluk yalnız shutdown metoduna bırakılamaz.

FlowFile ve transaction modeli

FlowFile nesnesi immutable olduğu için her değiştirme çağrısının döndürdüğü referans kullanılmalıdır:

flowFile = session.putAttribute(flowFile, "result", "success");
flowFile = session.write(flowFile, callback);
session.transfer(flowFile, SUCCESS);

Eski referansla işleme devam etmek, session içinde geçersiz sürüm kullanımı veya beklenmeyen yönlendirme hatası doğurabilir.

"ProcessSession", FlowFile oluşturma, okuma, yazma, klonlama, silme ve aktarma işlemlerini tek bir atomik çalışma birimi içinde toplar. Session commit edilirse değişiklikler kalıcı akış durumuna geçer. Hata halinde rollback uygulanabilir. "ProcessSession" thread-safe değildir ve "onTrigger()" dışındaki worker thread'lere aktarılmamalıdır. Harici paralellik gerekiyorsa session dışı veri önce bağımsız yapılara dönüştürülmeli, sonuçlar yine çağrıyı yürüten thread üzerinde session'a uygulanmalıdır.

İçerik işlemede temel tercih stream tabanlı çalışmadır. "session.read()", "session.write()" ve "StreamCallback", büyük veriyi tamamını heap içine almadan işleme imkanı verir. NiFi işlemcileri birden fazla concurrent task ile çalışabildiği için her FlowFile'ın bütünü belleğe alınırsa yaklaşık bellek ihtiyacı şu hale gelir:

anlık bellek = concurrent task sayısı x ortalama içerik boyutu

Sekiz eş zamanlı görevde 500 MB büyüklüğünde dosyaların byte dizisine alınması yalnız veri gövdeleri için yaklaşık 4 GB heap ister. Kopyalar ve parser nesneleri bu değere dahil değildir. Resmi geliştirici kılavuzu da içerik boyutu kesin olarak sınırlı değilse bütün FlowFile'ın belleğe alınmamasını önerir.

Harici sistemle transaction sınırı daha zordur. Kaynak dosya NiFi'ye alındıktan sonra dış kaynaktan silinecekse önce session güvenli biçimde commit edilmelidir. Dosya önce silinir ve NiFi commit öncesinde durursa veri kaybolur. Commit sonrasında, dış kaynaktaki silme işleminden önce durursa aynı veri yeniden alınabilir. Bu aralıkta tam olarak bir kez işleme garantisi yoktur. Çoğu veri toplama sisteminde kontrollü tekrar, veri kaybından daha kabul edilebilir bir sonuçtur.

Bu nedenle kaynak kimliği, boyut, son değiştirme zamanı ve gerekiyorsa içerik özeti üzerinden yinelenen kayıt kontrolü yapılmalıdır. MD5 hızlıdır ancak güvenlik veya kasıtlı çakışma riski bulunan doğrulamalarda kullanılmamalıdır. SHA-256 daha güçlü bir içerik kimliği sağlar, fakat büyük dosyanın tamamını ikinci kez okumak I/O maliyeti doğurur. Dosya adı ve metadata yeterliyse her akışta hash hesaplamak gereksizdir.

Eş zamanlılık ve bellek yönetimi

NiFi aynı processor örneğinin "onTrigger()" metodunu birden fazla thread ile çağırabilir. "@TriggerSerially" concurrent task sayısını bire indirir, ancak her çağrının aynı thread üzerinde çalışacağını garanti etmez. Processor alanlarında tutulan bütün değiştirilebilir durum yine thread-safe olmalıdır.

Ortak kaynak dizininden sıralı öğe seçimi yapılacaksa düz bir "int" sayaç ve değiştirilebilir liste kullanılmamalıdır. Sabit bir kaynak dizisi immutable veya "volatile" referansla yayımlanabilir. Round-robin indeks için "AtomicInteger" yeterlidir. Öğelerin tüketilip sonra havuza geri döndüğü bir modelde sınırlı "BlockingQueue" daha uygundur.

Ortak tek bir "byte[]" dizisi paralel "onTrigger()" çağrılarında veri yarışına yol açar. Orta boy ve sabit tamponlarda "ThreadLocal<byte[]>", tekrar eden allocation maliyetini azaltabilir. Büyük tamponlarda thread sayısıyla orantısız bellek büyümesini önlemek için sınırlı bir buffer pool kullanılmalıdır. Pool boşaldığında sınırsız yeni tampon üretmek yerine çağrı bekletilmeli, işlem küçük parçalarla sürdürülmeli veya kontrollü biçimde başarısız olmalıdır.

"ThreadLocal" kullanımı da sınırsız değildir. NiFi scheduler thread'leri uzun ömürlüdür. Çok büyük tamponlar ThreadLocal içinde tutulursa processor boşta kalsa bile heap'e bağlı kalabilir. Boyut ve concurrent task sayısı birlikte hesaplanmalıdır.

Processor kendi executor'ını kurmamalıdır. Harici thread havuzu NiFi scheduler'ın concurrent task, backpressure ve durdurma modelini aşabilir. Zorunlu bir protokol istemcisi kendi worker thread'lerini kullanıyorsa yaşam döngüsü "@OnScheduled" ve "@OnStopped" arasında açıkça yönetilmelidir.

Küme ve hata davranışı

Tek düğümde doğru çalışan kaynak processor, cluster içinde aynı dosyayı her düğümde okuyabilir. Uzak veya ortak dosya sistemi tarayan işlemcilerde "@PrimaryNodeOnly" bu tekrarları engelleyebilir. Daha yüksek paralellik gerekiyorsa dosya sahipliği, atomik claim veya cluster genelinde koordinasyon tasarlanmalıdır.

"StateManager", son işlenen kimlik, zaman damgası veya küçük bir checkpoint için uygundur. "Scope.LOCAL" her düğümde farklı durum tutar. "Scope.CLUSTER" bütün düğümlerin ortak durum görmesini sağlar. State Manager basit anahtar-değer verileri içindir ve yüksek hacimli kayıt deposu olarak kullanılmamalıdır. Cluster state haritası için boyut sınırı vardır. Büyük veya sık güncellenen durum harici bir veri deposuna ya da ortak Controller Service katmanına taşınmalıdır.

Backpressure yalnız akış tasarımcısının ayarı değildir. Processor'ın hata davranışıyla birlikte çalışır. Çıkış kuyruğu doluysa NiFi varsayılan olarak processor'ı çalıştırmaz ve baskıyı yukarı doğru yayar. Processor bunu kendi thread veya sınırsız dahili kuyruğuyla aşarsa NiFi'nin koruma mekanizması etkisiz kalır.

Penalize işlemi belirli bir FlowFile'ın yeniden erişilebilir olmasını geciktirir. Yield ise processor'ın tamamının kısa süre yeniden zamanlanmasını önler. Veri kaynaklı geçici hata için FlowFile penalize edilebilir. Uzak sistemin bütünü erişilemiyorsa "context.yield()" daha doğru olabilir. Beklenen hata "ProcessException" olarak dışarı çıkarsa framework session'ı geri alır ve ilgili FlowFile'ları penalize eder. Beklenmeyen runtime hatalarında idari yield de uygulanır.

Retry sınırsız olmamalıdır. Deneme sayısı attribute veya kalıcı durumla izlenebilir. Kısa ve deterministik timeout, artan bekleme ve üst sınır birlikte kullanılmalıdır. Sürekli erişilemeyen dış kaynak için circuit breaker, her tetiklemede aynı pahalı bağlantı girişiminin yapılmasını önler.

Test ve canlıya geçiş

NiFi Mock Framework, processor ve Controller Service testlerini gerçek yaşam döngüsüne yakın biçimde yürütür. "TestRunner" property yapılandırmasını, FlowFile kuyruğunu, ilişkileri, içerikleri, attribute değerlerini ve StateManager davranışını doğrulayabilir. "@OnScheduled", "onTrigger", "@OnUnscheduled" ve "@OnStopped" çağrıları test akışına dahil edilebilir.

Birim testleri şu durumları kapsamalıdır:

  • Boş giriş
  • Kısmi ve bozuk içerik
  • Aynı dosyanın yeniden görülmesi
  • Çıkış kuyruğu doluluğu
  • Uzak sistem zaman aşımı
  • Session rollback
  • Birden fazla concurrent task
  • Local ve cluster state ayrımı
  • Processor durdurulurken açık kaynakların kapanması
  • NFS veya dosya sistemi yazma başarısızlığı

Canlıya geçişten önce uzun süreli yük testi gerekir. Ortalama süre tek başına yeterli değildir. P95 ve P99 işlem süresi, FlowFile kuyruk uzunluğu, repository I/O, heap kullanımı, GC duraklaması, açık dosya tanıtıcısı ve tekrar deneme oranı birlikte izlenmelidir.

İşlemci etkinliğini "ComponentLog" üzerinden kaydetmeli ve dış kaynaktan veri alma veya dış hedefe gönderme işlemlerini "ProvenanceReporter" ile bildirmelidir. Provenance yalnız hata ayıklama kaydı değildir. Verinin hangi kaynaktan geldiğini, hangi dönüşümlerden geçtiğini ve nereye gönderildiğini gösteren zincirin parçasıdır.

Custom Processor kritik sistemde canlıya alınabilir. Bunun koşulu, yalnız normal akışta çalışması değildir. Aynı girdide deterministik sonuç üretmeli, thread yarışına açık durum taşımamalı, büyük veriyi stream olarak işlemeli ve hata halinde veri kaybetmemelidir. Kaynak toplama, içerik işleme, veritabanına aktarma ve arşivleme ayrı processor'lara bölündüğünde her aşamanın retry, backpressure ve gözlemlenebilirlik sınırı bağımsız yönetilebilir.

NiFi'nin değeri, özel kodu ortadan kaldırması değildir. Özel kodu transaction, kuyruk, provenance ve yaşam döngüsü kuralları içinde çalıştırmasıdır. İyi yazılmış bir Custom Processor bu kuralları aşmaya çalışmaz. Algoritmasını NiFi'nin veri akışı modeliyle uyumlu hale getirir.

Bu sayfanın QR kodu