Gereksinimler ve Ön Hazırlık
Uygulamaya başlamadan önce sisteminizde aşağıdaki araçların kurulu olduğundan emin olmalısınız. 2026 standartlarına uygun olarak en güncel sürümleri kullanmanız önerilir.
- Node.js (LTS sürümü, v22 veya üzeri)
- Apache Kafka (Docker üzerinde çalıştırılması tavsiye edilir)
- kafkajs (Node.js için modern ve popüler Kafka istemcisi)
- npm veya pnpm paket yöneticisi
Kafka'yı yerel ortamınızda hızlıca ayağa kaldırmak için bir docker-compose.yml dosyası kullanabilirsiniz. Bu, geliştirme sürecinde servislerinizi izole etmenize yardımcı olur.
version: '3.8'
services:
zookeeper:
image: confluentinc/cp-zookeeper:latest
environment:
ZOOKEEPER_CLIENT_PORT: 2181
kafka:
image: confluentinc/cp-kafka:latest
ports:
- "9092:9092"
environment:
KAFKA_ADVERTISED_LISTENERS: PLAINTEXT://localhost:9092
KAFKA_ZOOKEEPER_CONNECT: zookeeper:2181
Kafka İstemcisini Yapılandırma
Node.js projenizde Kafka ile iletişim kurmak için kafkajs kütüphanesini kullanacağız. İlk adım olarak bir Kafka istemcisi oluşturmalı ve bağlantı ayarlarını yapmalısınız. Bu yapılandırma, tüm mikroservisleriniz için merkezi bir nokta görevi görür.
const { Kafka } = require('kafkajs');
const kafka = new Kafka({
clientId: 'siparis-servisi',
brokers: ['localhost:9092'],
retry: {
initialRetryTime: 100,
retries: 8
}
});
module.exports = kafka;
Yukarıdaki kod, Kafka kümesine (cluster) bağlanmak için gerekli istemciyi tanımlar. retry ayarları, ağ kopmaları gibi geçici hatalarda sistemin otomatik olarak kendini toparlamasını sağlar.
Üretici (Producer) Servisi ile Olay Yayınlama
Mikroservis mimarisinde bir servis, bir işlem gerçekleştiğinde (örneğin yeni bir sipariş oluşturulduğunda) bir olay yayınlar. Bu olayı yayınlayan servise "Producer" denir. Aşağıdaki kod, basit bir sipariş oluşturma olayının Kafka'ya nasıl gönderileceğini gösterir.
const kafka = require('./kafka-config');
const producer = kafka.producer();
async function siparisOlustur(siparisVerisi) {
await producer.connect();
await producer.send({
topic: 'siparis-olusturuldu',
messages: [{ value: JSON.stringify(siparisVerisi) }],
});
await producer.disconnect();
}
Bu kod bloğu, siparis-olusturuldu adlı bir "topic" (konu) üzerine mesaj yazar. Veriyi gönderirken JSON formatında serileştirmek, farklı dillerle yazılmış servislerin de veriyi okuyabilmesini sağlar.
Tüketici (Consumer) Servisi ile Olayları Dinleme
Olay tabanlı mimarinin kalbi "Consumer" servisleridir. Bu servisler, belirli bir konu üzerindeki mesajları sürekli dinler ve yeni bir mesaj geldiğinde tanımlanan fonksiyonu tetikler. Örneğin, stok servisi sipariş olayını dinleyerek stok düşümü yapabilir.
const kafka = require('./kafka-config');
const consumer = kafka.consumer({ groupId: 'stok-servis-group' });
async function stokIslemleriniBaslat() {
await consumer.connect();
await consumer.subscribe({ topic: 'siparis-olusturuldu', fromBeginning: true });
await consumer.run({
eachMessage: async ({ topic, partition, message }) => {
const siparis = JSON.parse(message.value.toString());
console.log('Yeni sipariş alındı, stok güncelleniyor:', siparis.id);
},
});
}
Burada groupId kullanımı kritiktir. Aynı gruptaki servisler, mesajları kendi aralarında paylaşarak yük dengelemesi (load balancing) yaparlar. fromBeginning: true ayarı, servisin daha önce kaçırılmış mesajları da okumasını sağlar.
Mikroservisler Arası İletişim Yöntemlerinin Karşılaştırılması
Mikroservisler arasında iletişim kurarken Kafka gibi olay tabanlı sistemler ile geleneksel HTTP (REST/gRPC) yöntemleri arasında seçim yapmanız gerekir.
| Özellik | Olay Tabanlı (Kafka) | Senkron (REST/gRPC) |
|---|---|---|
| Bağımlılık | Düşük (Gevşek bağlı) | Yüksek (Sıkı bağlı) |
| Performans | Yüksek (Asenkron) | Düşük (Bekleme süresi) |
| Hata Yönetimi | Mesaj kuyruğu ile kurtarılabilir | Hemen hata döner |
Güvenlik ve Hata Yönetimi
Kritik Uyarı: Kafka üzerinde taşınan verilerin hassasiyetine göre TLS/SSL şifrelemesi kullanın. Üretim ortamında mesajların kaybolmaması için
acks: 'all'ayarını kullanmayı ve mesajların doğrulanması için Schema Registry gibi araçları tercih etmeyi unutmayın. Kod güvenliği sorumluluğu geliştiriciye aittir; asla ham veriyi doğrulamadan veritabanına işlemeyin.
Hatalı mesajların sisteminizi kilitlememesi için "Dead Letter Queue" (Ölü Mesaj Kuyruğu) mantığını kurmalısınız. İşlenemeyen mesajları ayrı bir konuya (topic) aktararak daha sonra manuel inceleme yapabilirsiniz.
// Hatalı mesajları yakalama örneği
try {
// İşleme mantığı
} catch (error) {
await producer.send({
topic: 'siparis-hatali',
messages: [{ value: JSON.stringify({ error: error.message, data: siparis }) }],
});
}
Sıkça Sorulan Sorular
Kafka'yı neden RabbitMQ yerine tercih etmeliyim?
Kafka, yüksek veri hacimli (high-throughput) ve verinin kalıcı olarak saklanması gereken senaryolar için tasarlanmıştır. RabbitMQ ise daha çok karmaşık yönlendirme mantığı gerektiren mesajlaşma ihtiyaçlarında öne çıkar.
Mesajların sırası önemli mi?
Evet, Kafka aynı partition içinde mesajların sırasını garanti eder. Ancak, paralel işleme yapıyorsanız mesajların işlenme sırası değişebilir. Bu durumda "key" bazlı partition stratejisi kullanmalısınız.
Node.js üzerinde Kafka kullanırken bellek yönetimi nasıl olmalı?
Büyük mesajlar (payload) göndermekten kaçının. Eğer çok büyük veri taşımanız gerekiyorsa, veriyi bir nesne depolama (S3 gibi) alanına yükleyip Kafka üzerinden sadece referans (URL) gönderin.
Kafka'da "Consumer Group" neden önemlidir?
Consumer Group, aynı mantıksal işi yapan servislerin mesajları paylaşarak işlemesini sağlar. Bu sayede uygulamanızı yatayda ölçekleyebilir ve performansı artırabilirsiniz.
Kafka'nın "offset" kavramı nedir?
Offset, bir consumer'ın bir partition içindeki hangi mesajda kaldığını belirten bir indekstir. Bu sayede uygulama çöktüğünde kaldığı yerden devam edebilir.
Kafka Mesajlarını İzleme ve Hata Ayıklama (Debugging)
Dağıtık sistemlerde bir mesajın kaybolması veya yanlış işlenmesi, sistemin bütününde tutarsızlıklara yol açabilir. Kafka üzerinde çalışan Node.js mikroservislerinizde hata ayıklamak için standart loglama yöntemlerinin ötesine geçmeniz gerekir.
Dead Letter Queue (DLQ) Stratejisi
İşlenemeyen mesajları doğrudan silmek yerine, bunları özel bir Dead Letter Queue konusuna (topic) yönlendirmek en iyi uygulamadır. Böylece hatalı mesajları daha sonra analiz edebilir veya manuel olarak tekrar işletebilirsiniz.
// Hatalı mesajı DLQ'ya gönderme örneği
const handleMessage = async (message) => {
try {
await processMessage(message);
} catch (error) {
console.error('Mesaj işlenemedi, DLQ\'ya aktarılıyor:', error);
await producer.send({
topic: 'orders-dead-letter',
messages: [{ value: message.value }],
});
}
};
KafkaJS ile Performans Optimizasyonu
Node.js üzerinde Kafka performansını artırmak için "batching" (toplu işleme) ve "compression" (sıkıştırma) özelliklerini doğru yapılandırmalısınız. Özellikle yüksek trafikli sistemlerde, her mesajı tek tek göndermek yerine belirli bir süre veya boyut eşiğine ulaşıldığında mesajları toplu göndermek ağ yükünü ciddi oranda azaltır.
| Parametre | Açıklama | Öneri |
|---|---|---|
lingerMs |
Mesajları göndermeden önce bekletme süresi (ms). | 5-10ms (Düşük gecikme için) |
compression |
Mesaj sıkıştırma algoritması. | GZIP veya Snappy |
batchSize |
Tek seferde gönderilecek maksimum bayt sayısı. | 16KB - 64KB |
Avro ile Mesaj Şeması Yönetimi (Schema Registry)
Mikroservisler arası iletişimde mesaj yapısının (JSON schema) değişmesi, tüketicilerin hata almasına neden olur. Confluent Schema Registry kullanarak mesajlarınızı Avro formatında serileştirmek, tip güvenliğini sağlar.
Avro kullanımı, mesajın hem üretici hem de tüketici tarafında aynı sözleşmeye (contract) sahip olmasını zorunlu kılar. Bu, "breaking changes" (kırılgan değişiklikler) riskini ortadan kaldırır.
// Avro şeması ile veri doğrulama örneği
const { SchemaRegistry } = require('@kafkajs/confluent-schema-registry');
const registry = new SchemaRegistry({ host: 'http://localhost:8081' });
const schemaId = 1; // Registry üzerindeki şema ID'si
const payload = { id: 1, status: 'CREATED' };
const encoded = await registry.encode(schemaId, payload);
await producer.send({
topic: 'orders',
messages: [{ value: encoded }]
});
Dağıtık İzleme (Distributed Tracing)
Bir kullanıcı isteğinin hangi mikroservislerden geçtiğini takip etmek için OpenTelemetry kütüphanesini Kafka ile entegre edin. Mesaj başlıklarına (headers) ekleyeceğiniz trace-id sayesinde, Kafka üzerinden geçen bir mesajın yaşam döngüsünü Jaeger veya Zipkin gibi araçlarla görselleştirebilirsiniz.
// Kafka mesaj başlığına trace-id ekleme
await producer.send({
topic: 'orders',
messages: [{
value: JSON.stringify(data),
headers: {
'trace-id': currentTraceId
}
}]
});
Sonuç
Node.js ile Kafka kullanarak olay tabanlı bir mikroservis mimarisi kurmak, uygulamanızın ölçeklenebilirliğini ve dayanıklılığını önemli ölçüde artırır. Bu rehberde, temel bağlantıdan mesaj üretimine ve tüketimine kadar olan süreci inceledik. Bir sonraki adım olarak, sisteminize bir "Schema Registry" entegre ederek mesaj yapılarındaki değişiklikleri yönetmeyi ve "Distributed Tracing" (dağıtık izleme) araçları ile servisler arası trafiği görselleştirmeyi deneyebilirsiniz.

Yorumlar (0)
Yorum Yaz