Derinlemesine yazılım eğitimleri için kanalımı takip edebilirsiniz...

.NET’te NATS JetStream İle Çalışmak

.NET'te NATS JetStream İle ÇalışmakMerhaba,


Microservice mimarileri yaygınlaştıkça uygulamalar arasındaki iletişim de giderek daha kritik bir hal almaktadır. Bir sipariş oluşturulduğunda stok servisinin haberdar edilmesi, ödeme tamamlandığında e-posta gönderilmesi veya arka planda çalışan onlarca servisin birbirleriyle güvenilir bir şekilde haberleşmesi artık modern yazılım geliştirmenin temel ihtiyaçlarından birisidir. Bu noktada mesajlaşma sistemleri (message broker) devreye girmektedir. Uzun yıllardır RabbitMQ, Apache Kafka ve Azure Service Bus gibi çözümler bu ihtiyacı karşılıyor olsa da son yıllarda performansı, düşük gecikme süresi ve sade mimarisiyle öne çıkan NATS ve onun kalıcı mesajlaşma katmanı olan JetStream, özellikle microserve tabanlı .NET uygulamalarında güçlü bir alternatif haline gelmektedir. Bu içeriğimizde, NATS ve JetStream‘in ne olduğunu, birbirlerinden hangi noktalarda ayrıldıklarını, hangi problemleri çözdüklerini ve Asp.NET Core uygulamalarında nasıl kullanılabileceklerini adım adım inceleyecek ve ayrıca makalenin sonunda yalnızca temel Publish/Subscribe yapısını değil, kalıcı mesajlar (Persistence), ACK mekanizması, Durable Consumer, Stream yapıları ve güvenilir event işleme gibi production ortamlarında ihtiyaç duyulan birçok özellikleri de uygulamalı olarak irdelemiş olacağız.

NATS Nedir?

NATS, distributed sistemler ve microservice mimariler için geliştirilmiş, yüksek performanslı ve düşük gecikmeli (low latency) bir mesajlaşma sistemidir (ya da bir başka deyişle message broker’dır) Temel amacı, uygulamaların birbirleriyle doğrudan bağlantı kurmasına gerek kalmaksızın güvenilir ve hızlı bir şekilde haberleşebilmelerini sağlamaktır.

NATS’ın Çalışma Mantığı Nasıldır?

NATS’ın temel çalışma prensibi Puclic/Subscribe mantığına dayanmaktadır. Bir subject‘e mesaj yayınlandığı taktirde, o subject’e abone olan/dinleyen herkese bu mesaj gönderilecektir. Eğer subject’i kimse dinlemiyorsa, mesaj kaybolacaktır. Yani anlayacağınız canlı bildirim senaryoları için oldukça uygundur ancak iş kuyruğu dediğimiz (work queue) için hiçte değildir!

Esasında NATS’te queue diye bir kavram yoktur. Onun yerine subject kavramı vardır. Subject, mesajların hangi topic’e ait olduğunu belirten mantıksal kanaldır. Producer mesajını belirli bir subject’e yayınlar, bu subject’i dinleyen tüm subscriber’lar ise mesajı elde edebilirler.

.NET'te NATS JetStream İle Çalışmak.NET'te NATS JetStream İle Çalışmak

JetStream Nedir?

JetStream, NATS’in persistent messaging katmanıdır.

NATS’in varsayılan çalışma mantığında mesajlar yalnızca o anda subscribe olmuş olanlara iletilecektir. Eğer ilgili subscriber çalışmıyorsa veya o subject’i dinleyen hiçbir uygulama yoksa, mesaj kalıcı olarak saklanmayacak ve kaybolacaktır.

Bu davranış gerçek zamanlı (real-time) haberleşme senaryoları için oldukça avantajlıdır. Çünkü NATS’in bu kadar yüksek performans sunmasının en önemli sebeplerinden biri de mesajları varsayılan olarak diske yazmamasıdır.

Ancak hak verirsiniz ki her sistem canlı bildirim mantığıyla çalışmamaktadır. Örneğin, bir e-ticaret uygulamasında oluşturulan sipariş, tamamlanan ödeme veya gerçekleştirilen banka transferi gibi olayların kaybolması kabul edilebilir bir durum değildir 🙂 Eğer ilgili servis o anda çalışmıyorsa bile, sistem ayağa kalktığında bu mesajın alınıp işlenebilmesi gerekmektedir.

İşte JetStream, tam olarak bu problemi çözmek için geliştirilmiştir. JetStream, NATS’in üzerine eklenen bir persistence (kalıcı mesajlaşma) katmanıdır. Mesajları yalnızca iletmekle kalmamakta, aynı zamanda belirlenen kurallara göre diskte saklamakta, gerektiğinde ise tekrar tüketilmesini sağlayarak, mesajların güvenilir bir şekilde işlenmesini takip etmektedir.

Böylece NATS, yalnızca anlık mesajlaşma yapan bir message broker olmaktan çıkabilmekte ve güvenilir mesaj teslimi (reliable messaging) yapabilen ve event streaming yeteneğine sahip olan güçlü bir platform haline gelmektedir.

JetStream sayesinde;

  • Mesajlar kalıcı olarak saklanabilmektedir.
  • Daha sonra sisteme katılan consumer’lar geçmiş mesajları okuyabilmektedir.
  • İşlenemeyen mesajlar tekrar işlenebilmektedir (retry)
  • Consumer’ın mesajı başarıyla işlediği ACK mekanizmasıyla doğrulanabilmektedir.
  • Consumer kapansa bile kaldığı yerden devam edebilmektedir (durable consumer)
  • İstenildiğinde geçmiş mesajlar yeniden oynatılabilmektedir (replay)

Kısaca JetStream’in temel amacı, Core NATS’in hızından vazgeçmeden, mesajların güvenilir bir şekilde saklanmasını ve tüketilmesini sağlamaktır.

.NET'te NATS JetStream İle Çalışmak

Kullanımını İnceleyelim…

NATS ve JetStream’in hangi problemleri çözdüğünü ve birbirlerinden hangi noktalarda ayrıldıklarını incelediğimize göre, artık bunları .NET uygulamalarında nasıl kullanabileceğimize geçebiliriz. Bu aşamada, öncelikle NATS sunucusunun kurulumuna odaklanacak, ardından bu sunucuya bağlantı kurarak, mesaj yayınlama ve dinleme işlemlerini ele alıyor olacağız. Devamında ise aynı örneği JetStream desteğiyle geliştirerek mesajların kalıcı olarak saklanmasını sağlayacak ve nihai olarak da ACK mekanizmasını ve güvenilir mesaj tüketimini de uygulamalı olarak inceliyor olacağız.

NATS Server’ın Kurulumu (Docker)

NATS ile ilgili çalışmalara geçmeden önce mesajları yönetecek olan NATS Server’ı ayağa kaldırmamız gerekmektedir. Bunun için Core NATS’ı JetStream desteğiyle birlikte Docker ortamında çalıştırabilmek için aşağıdaki talimattan istifade edebiliriz:

docker run -d --name nats -p 4222:4222 -p 8222:8222 nats:latest -js -m 8222

Bu talimatta kullandığımız -js parametresi ile NATS Server’da JetStream özelliği aktif bir şekilde gelmektedir. -m 8222 parametresi ise NATS Server’ın monitoring dashboard’unu http://localhost:8222/ portunda aktifleştirecektir.

Asp.NET Core Üzerinden Core NATS’a Bağlanma

Asp.NET Core’dan NATS’a bağlantı sağlayabilmek için öncelikle aşağıdaki iki paketin uygulamaya yüklenmesi gerekmektedir:

Ardından aşağıdaki yapılandırmalar eşliğinde servisler uygulamaya eklenmelidir:

builder.Services.AddNatsClient(natsBuilder =>
    natsBuilder.ConfigureConnection(natsConnection => new NatsConnection(NatsOpts.Default with { Url = "nats://localhost:4222" })));

builder.Services.AddSingleton(serviceProvider =>
    serviceProvider.GetRequiredService<INatsConnection>().CreateJetStreamContext());

Burada AddNatsClient metodu ile NATS’a ait servisler uygulamanın IoC container’ına yüklenmekte ve bir yandan da NATS Server’a bağlantı gerçekleştirilmektedir. CreateJetStreamContext metodu ise bu NATS bağlantısı üzerinden JetStream özelliklerine erişim sağlamaktadır.

Mesaj Yayınlama ve Subscriber Oluşturma

Asp.NET Core üzerinden mesaj yayınlamayı örneklendirebilmek için aşağıdaki gibi basit bir endpoint oluşturulması yeterli olacaktır kanaatindeyim:

app.MapGet("/create-job", async (INatsJSContext natsJSContext, CancellationToken cancellationToken) =>
{
    PubAckResponse pubAckResponse = await natsJSContext.PublishAsync(
                                            subject: Subjects.JobsSubject,
                                            data: new Message(Text: "Naber düyna!", Date: DateTime.UtcNow),
                                            cancellationToken: cancellationToken);
    pubAckResponse.EnsureSuccess();

    return TypedResults.Ok();
});

Ve bu mesajı yine Asp.NET Core’da tüketebilmek için de aşağıdaki gibi bir background service tasarlanabilir:

    public class JowWorkerBackgroundService(INatsJSContext natsJSContext) : BackgroundService
    {
        protected override async Task ExecuteAsync(CancellationToken stoppingToken)
        {
            await natsJSContext.CreateStreamAsync(new StreamConfig("JOBS", ["jobs-subject"])
            {
                Retention = StreamConfigRetention.Workqueue,
                Storage = StreamConfigStorage.File
            });

            var consumer = await natsJSContext.CreateOrUpdateConsumerAsync("JOBS", new ConsumerConfig("JOBS-CONSUMER")
            {
                AckWait = TimeSpan.FromSeconds(30),
                AckPolicy = ConsumerConfigAckPolicy.Explicit,
                MaxDeliver = 5
            }, cancellationToken: stoppingToken);

            await foreach (var message in consumer.ConsumeAsync<Message>(cancellationToken: stoppingToken))
            {
                await Console.Out.WriteLineAsync($"Message : {message.Data.Text}");
                await Console.Out.WriteLineAsync($"Date : {message.Data.Date}");
                await Console.Out.WriteLineAsync($"****************");
                await message.AckAsync();
            }
        }
    }

Burada önemli gördüğüm bazı noktaları izah etmekte fayda görmekteyim; CreateStreamAsync metodu ile mesajları kalıcı olarak saklayabilmek için bir stream oluşturulmaktadır. Normal şartlarda NATS, fire-and-forget davranışı sergilediği, yani alıcı yoksa mesajlar kaybolacağı için bu özellik tam da bu noktada JetStream ile devreye sokulmuş olacaktır. Dikkat ederseniz Storage property’sine ‘File’ değeri verilmiştir. Böylece mesajlar disk’te saklanacak ve sunucu kapansa dahi korunacaktırlar. Eğer bu property’e ‘Memory’ değeri verilmiş olsaydı mesajlar bellekte saklanacaktı ve sunucu kapandığı taktirde doğal olarak kaybolacaklardı…

.NET'te NATS JetStream İle ÇalışmakEe peki hoca, mesajlar disk’te kalıcı hale getirildiğinde nerede tutuluyorlar? sorunuzu duyar gibiyim… Bunun için NATS’in dashboard’undan JetStream verilerine http://localhost:8222/jsz adresinden bakmamız yeterli olacaktır.

Dikkat ederseniz yandaki görselde de vurgulandığı gibi store_dir alanı mesajların kalıcı olarak saklandığı path’i bizlere vermektedir.

Retention property’sine gelirsek eğer mesajların ne zaman silineceğini belirleyen bir özelliktir. Burada ‘Workqueue’ diyerek ilgili mesajın ACK edildiği taktirde silineceğini ifade etmiş oluyoruz. Ayrıca ‘Limits’ diyerek belirleyebileceğimiz boyut, sayı veya süre limitine göre silmeyi gerçekleştirtebiliriz ya da Interest diyerek de hiç consumer kalmadığı taktirde mesajları sildiretebiliriz.

Diyelim ki, örnekte olduğu gibi bu alan ‘Workqueue’ olarak yapılandırılmış olmasına rağmen, foreach döngüsü içerisinde AckAsync çağrısı yapılmamış olsaydı, mesaj AckWait süresi dolduktan sonra yeniden subscriber tarafından teslim alınacaktı.

CreateOrUpdateConsumerAsync metoduna gelirsek eğer bunla da adı üzerinde oluşturulan stream’i okuyan/tüketen bir consumer oluşturulmaktadır. Burada da akıllara muhtemelen birden fazla consumer oluşturursak ne olur? sorusu gelecektir. Bu sorunun iki türlü cevabını verebiliriz;

  • Aynı isimde farklı makinelerde consumer oluşturulursa…
    // Makine 1
    new ConsumerConfig("JOBS-CONSUMER")
    
    // Makine 2
    new ConsumerConfig("JOBS-CONSUMER")  // AYNI İSİM
    
    // Makine 3
    new ConsumerConfig("JOBS-CONSUMER")  // AYNI İSİM
    

    Bu durumda mesajlar makineler arasında paylaştırılacak ve böylece load balancing davranışı sergilenmiş olacaktır. Her mesaj sadece bir worker’a gönderilecektir…

    Mesaj 1 → Makine 1
    Mesaj 2 → Makine 2
    Mesaj 3 → Makine 1
  • Farklı isimlerde consumer’lar oluşturulursa…
    // Worker 1
    new ConsumerConfig("JOBS-CONSUMER-1")
    
    // Worker 2
    new ConsumerConfig("JOBS-CONSUMER-2")  // FARKLI İSİM
    
    // Worker 3
    new ConsumerConfig("JOBS-CONSUMER-3")  // FARKLI İSİM
    

    Bu durumda da birbirlerinden gerçekten farklı consumer’lar üretilmiş olacaktır. Dolayısıyla mesajlar hepsine gönderilmiş olacaktır.

Console Application’dan Subscriber Oluşturma

Ayrıca Asp.NET Core’daki background service’in dışında bir console uygulamasından da nasıl subscriber oluşturabileceğimizi ele almakta fayda görmekteyim. Bunun için ilgili console uygulamasına NATS.Net kütüphanesi yüklenmeli ve aşağıdaki gibi bir çalışma gerçekleştirilmelidir:

using NATS.JetStream.Demo.Shared;
using NATS.Net;

await using var natsClient = new NatsClient();

await foreach (var message in natsClient.Connection.SubscribeAsync<Message>(subject: Subjects.JobsSubject))
{
    await Console.Out.WriteLineAsync($"Message : {message.Data.Text}");
    await Console.Out.WriteLineAsync($"Date : {message.Data.Date}");
    await Console.Out.WriteLineAsync($"****************");
}

İşlem Tamamlandıktan Sonra ACK Gönderin!

JetStream ile çalışan consumer uygulamalarında dikkat edilmesi gereken en önemli husus, ACK (Acknowledgement) mesajının iş mantığı başarıyla tamamlandıktan sonra gönderilmesidir. Bunun nedeni, JetStream’in sunduğu at-least-once delivery (en az bir kez teslimat) garantisine davranışsal olarak riayet edebilmektir. Bu garanti sayesinde bir consumer, kendisine iletilen mesajları işlerken beklenmedik bir şekilde kapanabilir yahut ACK göndermeden önce hata alabilir. Böyle bir durumda JetStream, mesajların kaybolmasını önlemek için belirli bir süre sonra aynı mesaj(lar)ı yeniden ilgili consumer’a teslim etmelidir. Aksi taktirde mesaja dair işlem tamamlanmadan önce ACK gönderilirse JetStream, ilgili mesajın başarıyla işlendiğini varsayacak ve böylece bahsedilen bu riskler ciddi olarak göze alınmış olacaktır. Bu nedenle mesaj alındıktan sonra işlenmesi beklenmeli ve başarıyla işlendiyse ACK gönderilmelidir.

Consumer’ları Idempotent Olarak Ayarlayın!

Bununla birlikte at-least-once delivery modelinin doğal bir sonucu olarak aynı mesajın birden fazla kez işlenmesi de söz konusu olabilir. Yani consumer bir mesajı başarıyla işlemiş ancak ACK gönderilmediği için JetStream tarafından aynı mesaj (türlü nedenlerden dolayı) tekrar gönderilmiş olabilir. Evet… JetStream, mesajların consumer’lar tarafından işlenip işlenmediğini bilememektedir. Bu nedenle consumer uygulamalarının da idempotent olarak tasarlanması büyük önem taşımaktadır. Başka bir ifadeyle, aynı mesaj birden fazla kez alınsa bile sistem her seferinde aynı sonucu üretmeli ve aynı işlemi tekrar tekrar gerçekleştirmemelidir.

Nihai olarak;
Bu makale boyunca NATS ve JetStream’in temel çalışma mantığını, birbirlerinden hangi noktalarda ayrıldıklarını ve .NET uygulamalarında nasıl kullanılabileceklerini adım adım incelemiş olduk. Görüldüğü üzere NATS, düşük gecikme süreleriyle işlevsellik gösteren, son derece hızlı bir message broker’dır. Varsayılan olarak mesajları kalıcı şekilde saklamadığı için canlı bildirimler ve gerçek zamanlı haberleşmeler gibi senaryolara oldukça uyum göstermektedir. Ancak mesaj kaybının kabul edilemeyeceği iş akışlarında yeterli olmayacağı için JetStream imdadına yetişmektedir. JetStream ile NATS’in üzerine kalıcı mesajlaşma, ACK mekanizması, retry, replay ve durable consumer gibi güvenilir bir mesajlaşma altyapısı sunan özellikler eklenmektedir.

.NET tarafından olaya bakarsak eğer NATS.NET paketi sayesinde hem Core NATS’in hem de JetStream yapılandırmasının oldukça sade ve geliştirici dostu olduğunu görmekteyiz. Birkaç satır kodla mesaj yayınlayabilmekte, farklı servislerde bu mesajları dinleyebilmekte ve JetStream sayesinde production ortamlarında ihtiyaç duyulan güvenilir mesaj işleme senaryolarını kolaylıkla hayata geçirebilmekteyiz.

Umarım bu içerik, NATS ve JetStream’e başlangıç yapmanız için faydalı bir kaynak olmuştur.

İlgilenenlerin faydalanması dileğiyle…
Sonraki yazılarımda görüşmek üzere…
İyi çalışmalar…

Örnek çalışmaya aşağıdaki GitHub reposundan erişebilirsiniz.
https://github.com/gncyyldz/NATS.JetStream.Demo.API

Bu repository, ilgili konuya dair örnek çalışmanın kaynak kodlarını ve mimari yapısını içermektedir. Detaylar için GitHub üzerinden incelemede bulunabilirsiniz.


GitHub’da Görüntüle →

Bunlar da hoşunuza gidebilir...

Bir yanıt yazın

E-posta adresiniz yayınlanmayacak. Gerekli alanlar * ile işaretlenmişlerdir