Programlama ve Yazılım Geliştirme

Go Kullanarak Oturum Sırasını Koruyan bir Kafka İşleme Hattı Nasıl Oluşturulur

Joshua Oluikpe, her oturum içinde mesajların katı bir sırayla işlenmesini sağlayan ve bağımsız oturumların paralel çalışmaya devam etmesine olanak tanıyan pratik bir Kafka mesaj işleme tasarımını açıklıyor. Tasarım; tutarlı hashing, oturum başına bir Go worker'ı, aynı akış içinde yeniden denemeyi ve çökme sonrasında mesaj kaybını önlemek için bitişik commit noktalarının yönetimini temel alıyor.

2026-10-07
4 dk okuma
1 görüntülenme
certi.news Editorial Team
Go Kullanarak Oturum Sırasını Koruyan bir Kafka İşleme Hattı Nasıl Oluşturulur

Binlerce oturum aynı partition'ı paylaştığında yalnızca Kafka'nın partition içindeki sırasına güvenmek yeterli değildir. Sohbet ve sıralı görev sistemlerinde, «hedefi Paris yap» gibi düzeltici bir mesaj «Londra'ya uçuş rezerve et» mesajından sonra gelebilir; ancak ilk mesajdan önce işlenmesi bağlamı bozar. Joshua Oluikpe'nin InfoQ'ta açıkladığı çözüm, sıralama sorumluluğunun bir bölümünü Go ile yazılmış uygulama katmanına aktarır.

Oturum içinde sıralama ve oturumlar arasında paralellik

Tasarım iki worker katmanına dayanır. Dağıtım katmanı Kafka kayıtlarını alır ve tutarlı hashing kullanarak oturum tanımlayıcısına göre yönlendirir. Böylece aynı oturumun mesajları aynı worker'a ulaşırken farklı oturumlar birden fazla worker'a dağıtılabilir. Ardından her etkin oturum, mesajları her seferinde bir tane ve geliş sırasına göre işleyen özel bir goroutine alır.

Bu sayede yavaş bir oturum, düz bir worker havuzunda olduğu gibi diğer oturumların önünde engel oluşturmaz; böyle bir havuzda tek bir oturumun yeniden denenmesi worker'ı tamamen meşgul edebilir. Goroutine'ler ihtiyaç halinde oluşturulur ve bir süre etkinlik olmadığında kaldırılır. Böylece kaynak tüketimi, mümkün olan toplam oturum sayısı yerine etkin oturumların sayısıyla ilişkili olur.

Mesajların önüne geçmeden yeniden deneme

Yeniden deneme, oturumun kendi goroutine'i içinde gerçekleştirilir. Bağımlı bir hizmetin zaman aşımına uğraması gibi geçici bir hata oluştuğunda worker, bir miktar rastgelelik içeren üstel geri çekilme kullanarak bekler ve sonraki mesajlara geçmeden aynı mesajı yeniden işler. Sonraki mesajlar kanalda takılan mesajın arkasında kaldığından dördüncü mesaj üçüncü mesajın önüne geçemez.

Geçersiz veriler gibi düzeltilemeyen hatalar ise doğrudan dead-letter queue'ya (DLQ) gönderilir. Geçici ve nihai hatalar arasındaki bu ayrım, yeniden denemeyi yönetmek için ayrı bir durum makinesine duyulan ihtiyacı ortadan kaldırır; ancak aynı zamanda belirli bir oturumun mesajlarının geri çekilme süresi boyunca birikmesine yol açabilir.

Güvenli commit ve çökme sonrasında kurtarma

Birden fazla oturum mesajlarını paralel olarak işlerken yalnızca tamamlanmış en yüksek offset'e commit yapılamaz. 104 numaralı offset, 102 numaralı offset'ten önce tamamlanır ve 104'e commit edilirse tüketicinin çökmesi 102'nin kalıcı olarak atlanmasına yol açabilir. Bu nedenle sistem, her partition için işlenmekte olan ve tamamlanmış offset'leri izler ve commit noktasını yalnızca bitişik biçimde tamamlanmış en yüksek offset'e kadar ilerletir.

Sonuç olarak yeniden başlatma daha eski bir offset'ten başlayabilir ve bazı mesajları yeniden işleyebilir. Tasarım, yinelenenleri ayıklamaya yardımcı olmak için sabit bir olay kimliği kullanır; ancak «tam olarak bir kez» yürütme sağladığını iddia etmez. Dış etkiler hâlâ bağımlı hizmetlerin yinelenenleri idempotent biçimde işleyebilmesine bağlıdır. Nadir durumlarda, partition ilerlemesini yeniden sağlamak için belirli bir zaman aşımı dolduktan sonra takılı bir boşluk atlanabilir; bu durum operasyonel bir hata olarak kaydedilir ve mesaj DLQ'ya yönlendirilir. Bu, süreklilik ile işlemenin eksiksizliği arasında açık bir tercihtir.

Gerçek üretim ortamı için neler gerekir?

Sıralama güvenceleri yalnızca temel mantıkla tamamlanmaz. Partition'lar yeniden dağıtılırken tüketici yeni kayıtları almayı durdurur, devam eden işlerin tamamlanmasını bekler ve sahipliği devretmeden önce son bitişik watermark'a commit eder. Worker kanalları dolduğunda mesajları düşürmek yerine getirmeyi geçici olarak durdurur; önceden çekilmiş kayıtlar ise daha yeni mesajların önüne geçmemeleri için partition'a özel depolarda tutulur.

Sistemin ayrıca takılı offset'leri tek tek tespit etmesi, Kafka'nın asenkron sınırları boyunca izlemeyi ilişkilendirmesi, getirme durdurulurken otomatik ölçeklendirme için gereken gecikme sinyallerini koruması, kapanmadan önce işleri boşaltması ve DLQ'ya yazması gerekir. Ürün, asenkron üretim sırasında mesaj sırasını güvence altına almak için gönderim tarafında da tutarlı hashing uygular.

Sonuçlar ve sınırlamalar

10 mesajın %10'una geçici hata enjekte edilen, 50.000 mesaj ve 10 oturumdan oluşan sentetik bir testte sistem, herhangi bir sıralama ihlali olmadan saniyede 14.027 mesaj ve 11,7 milisaniyelik medyan gecikme kaydetti. 1.000 oturumla yapılan yük testinde gerçek gönderim hızı, 50.000 hedefinin 48.805 mesaj/saniye seviyesine ulaştı ve hiçbir gönderim hatası görülmedi; ancak kuyruk gecikmesi, oturum sıralama mekanizması nedeniyle değil, API Gateway'deki mesaj hızı kotası nedeniyle 200 milisaniyelik bütçeyi aştı.

Yazar, sistemin gözlemlenebilir bir sıralama ihlali olmadan 40 milyondan fazla üretim mesajını işlediğini ve yalnızca 460 mesajı DLQ'ya yönlendirdiğini, bunun da yaklaşık %0,002'ye karşılık geldiğini belirtiyor. Ancak bu sayı, her mesajın dizisinin kapsamlı biçimde incelenmesine değil operasyonel gözleme dayanıyor; ayrıca daha yüksek yük testi her oturum için sıralamayı bağımsız olarak doğrulamadı. Bu nedenle tasarımın pratik değeri, sıralamayı Kafka düzeyindeki bir varsayımdan uygulamanın zorunlu kıldığı bir özelliğe dönüştürmesinde yatıyor; bununla birlikte açık sıralama testlerine ve bağımlı hizmetlerde idempotency güvencelerine duyulan ihtiyaç devam ediyor.

Haber kaynağı
InfoQ - Architecture Articles
Özgün kaynağı aç ↗
c
Yazar

certi.news Editorial Team

Aynı kategoride

Bunlar da ilginizi çekebilir

Tüm haberleri gör