Spring Boot uygulamalarında Kafka consumer lag'ini ölçerek azaltma, cooperative rebalance ayarı, partition temelli concurrency hesabı ve gecikmeli retry-DLQ akışını somut yapılandırmalarla ele alır.
Spring Boot'ta Kafka Consumer Lag: Rebalance, Concurrency ve DLQ Tasarımı
Spring Boot'ta Kafka consumer lag'ini doğru teşhis etmek
Consumer lag, tek başına "consumer yavaş" anlamına gelmez. Lag'in artış hızı, üretilen kayıt hızı ile commit edilen offset hızı arasındaki farktır. İlk teşhiste grup seviyesinde toplam lag yerine partition seviyesindeki dağılıma bakın. Bir partition'da 500.000, diğerlerinde 0 lag varsa consumer concurrency artırmak sorunu çözmez; büyük olasılıkla key dağılımı bozuktur veya ilgili partition'da zehirli bir kayıt tekrar deneniyordur. Kafka'nın kendi aracıyla başlangıç fotoğrafını alın:
kafka-consumer-groups.sh --bootstrap-server kafka-1:9092 --describe --group payment-worker
# PARTITION, CURRENT-OFFSET, LOG-END-OFFSET ve LAG alanlarini saklayinBu teşhis, bir spring framework eğitimi veya spring boot eğitimi sırasında genellikle atlanan uygulama sınırını görünür kılar: @KafkaListener çağrısı HTTP request thread'i gibi sınırsız paralel değildir. Her atanan partition, consumer thread'i üzerinde sıralı işlenir. java spring eğitimi kapsamında spring mvc ile yazılan bir spring rest api siparişi Kafka'ya yayınlıyor olabilir; fakat API'nin p95 süresi ile consumer'ın kayıt işleme süresi ayrı ölçülmelidir. Producer tarafında request rate'i, consumer tarafında ise JMX altındaki kafka.consumer:type=consumer-fetch-manager-metrics,client-id=... MBean'inden records-lag-max, records-consumed-rate ve fetch-latency-avg değerlerini toplayın. Bu MBean'lerin Prometheus isimleri kullandığınız JMX exporter kuralına göre değişir; dashboard sorgusunu exporter çıktısına bakmadan varsaymayın.
Ölçüm için değişiklik öncesinde aynı yük altında en az 15 dakika referans alın. Prometheus'ta exporter'ın ürettiği grup lag metriği için örnek yaklaşım şudur:
# Metrik adini kendi exporter'iniza gore dogrulayin
max_over_time(kafka_consumergroup_lag{group="payment-worker"}[15m])
# Lag artma hizi: pozitif ve surekli ise yetisme problemi vardir
deriv(kafka_consumergroup_lag{group="payment-worker"}[10m]) Aynı pencerede rebalance sayısını, kayıt başına dış servis süresini ve hata oranını kaydedin. Sonraki ayarın başarılı sayılması için yalnızca lag'in düşmesi yetmez; deploy sırasında rebalance süresi uzamamalı ve işleme başarısızlıkları artmamalıdır.Partition sayısına göre concurrency ve poll döngüsü ayarlamak
Concurrency üst sınırı topic partition sayısıdır. Örneğin 12 partition ve aynı consumer group içinde 3 pod varsa pod başına concurrency=4 ile toplam 12 aktif consumer elde edilir. Pod başına 8 vermek 24 consumer oluşturur; 12 tanesi atama alamaz, fakat TCP bağlantısı, buffer ve metrik cardinality tüketmeye devam eder. Bu hesap autoscaling sırasında da yapılmalıdır: hedef pod sayısı değiştiğinde podSayisi * concurrency <= partitionSayisi koşulunu deployment kontrolüne ekleyin.
Aşağıdaki factory, kayıt işleme süresi değişken olan bir consumer için poll sınırını açık hale getirir. Kritik ilişki şudur: max.poll.records * enYavasKayitSuresi, max.poll.interval.ms değerini aşarsa broker consumer'ı gruptan çıkarır. Uygulama kaydı işliyor olsa bile rebalance başlar, aynı kayıt başka instance'a verilir ve duplicate etkisi oluşur.
@Configuration
class KafkaConsumerConfig {
@Bean
ConcurrentKafkaListenerContainerFactory<String, PaymentEvent> kafkaListenerFactory(
ConsumerFactory<String, PaymentEvent> consumerFactory) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, PaymentEvent>();
factory.setConsumerFactory(consumerFactory);
factory.setConcurrency(4);
factory.getContainerProperties().setAckMode(ContainerProperties.AckMode.RECORD);
return factory;
}
}
# application.yml
spring:
kafka:
consumer:
enable-auto-commit: false
max-poll-records: 20
properties:
max.poll.interval.ms: 120000
max.poll.records: 20
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignorBu örnekte p99 kayıt işleme süresi 2 saniyeyse 20 kayıt teorik olarak 40 saniyede tamamlanır; 120 saniye, dış bağımlılık sıçraması için pay bırakır. p99 10 saniyeye çıkıyorsa interval'i körlemesine büyütmek yerine max.poll.records değerini azaltın veya dış çağrıyı zaman aşımı ve bulkhead ile sınırlandırın. Büyük interval, gerçekten durmuş consumer'ın group tarafından daha geç fark edilmesine neden olur. CooperativeStickyAssignor ise rebalance sırasında tüm partition'ları aynı anda iptal etmek yerine mümkün olduğunda artımlı atama yapar. Bunun çalışması için group içindeki tüm instance'ların uyumlu assignor yapılandırmasıyla kademeli olarak güncellenmesi gerekir; karma assignor geçişi beklediğiniz kadar az kesintili olmayabilir.
Spring Cloud ortamında rebalance kaynaklı lag'i önce-sonra ölçmek
Kubernetes üzerinde çalışan bir microservices mimarisi için pod kapanışı, rebalance'ın en sık tetikleyicisidir. spring cloud ile merkezi konfigürasyon dağıtıyor olsanız bile Kafka consumer ayarlarını canlı yenilemek güvenli değildir: mevcut container'ın consumer instance'ı bu property'leri kendiliğinden yeniden kurmaz. max.poll.records, assignor veya concurrency değişikliğini kontrollü rolling deployment ile uygulayın; property refresh'i ile yarım uygulanmış bir consumer kümesi oluşturmayın.
Önce 15 dakikalık temel ölçümde maksimum lag, rebalance olayları ve pod termination süresini kaydedin. Sonra cooperative assignor, gerçek partition sayısına uygun concurrency ve yeterli termination grace period uygulayın. Kubernetes manifestinde listener'ın işlem bitirmesi için süre tanıyın:
spec:
terminationGracePeriodSeconds: 150
containers:
- name: payment-worker
lifecycle:
preStop:
exec:
command: ["/bin/sh", "-c", "sleep 10"]
env:
- name: SPRING_KAFKA_CONSUMER_PROPERTIES_PARTITION_ASSIGNMENT_STRATEGY
value: org.apache.kafka.clients.consumer.CooperativeStickyAssignor preStop içindeki 10 saniye, endpoint'lerin pod'u listeden çıkarması için bir tampon sağlar; asıl sınır terminationGracePeriodSeconds değeridir. Bu değer, en kötü batch işleme süresi ile uygulamanın consumer kapatma süresinden küçük olmamalıdır.Karşılaştırmayı aynı trafik penceresinde yapın: değişiklikten önce ve sonra max_over_time(lag[15m]), deployment başına rebalance sayısı ve deploy anındaki p99 işleme süresini tabloya koyun. Lag düşerken deploy anında duplicate ödeme denemeleri artıyorsa sonuç başarısızdır. Özellikle spring boot kursu örneklerinde sık görülen hata, HTTP client timeout'unu 60 saniye, max.poll.interval.ms değerini 30 saniye bırakmaktır. Tek bir yavaş downstream çağrısı bile consumer'ın üyeliğini kaybetmesine yeter.
Spring Boot retry topic ve DLQ ile poison message izolasyonu
Bir kaydın işlenmesi deterministik olarak başarısızsa listener thread'inde Thread.sleep ile retry bekletmeyin. Bu süre boyunca partition ilerlemez, lag büyür ve yeterince uzun beklemede poll interval aşılır. Spring Kafka'nın DefaultErrorHandler ve DeadLetterPublishingRecoverer bileşenleri, sınırlı sayıda kısa denemeden sonra kaydı topic.DLT benzeri bir hedefe yollar. Böylece aynı partition'daki sonraki kayıtlar ilerler. Aşağıdaki yapılandırma üç deneme sonrası DLT'ye aktarır:
@Bean
DefaultErrorHandler kafkaErrorHandler(KafkaTemplate<Object, Object> template) {
var recoverer = new DeadLetterPublishingRecoverer(
template,
(record, exception) -> new TopicPartition(record.topic() + ".DLT", record.partition()));
var backOff = new FixedBackOff(1_000L, 2L);
var handler = new DefaultErrorHandler(recoverer, backOff);
handler.addNotRetryableExceptions(IllegalArgumentException.class);
return handler;
}
@Bean
ConcurrentKafkaListenerContainerFactory<String, PaymentEvent> kafkaListenerFactory(
ConsumerFactory<String, PaymentEvent> consumerFactory,
DefaultErrorHandler kafkaErrorHandler) {
var factory = new ConcurrentKafkaListenerContainerFactory<String, PaymentEvent>();
factory.setConsumerFactory(consumerFactory);
factory.setCommonErrorHandler(kafkaErrorHandler);
return factory;
}Her exception retry edilebilir değildir. JSON şema ihlali, zorunlu alan eksikliği ve desteklenmeyen event version'ı gibi IllegalArgumentException türleri tekrar denendiğinde farklı sonuç vermez; doğrudan DLT'ye gönderilmelidir. Buna karşılık geçici ağ hataları için kısa retry uygulanabilir. Daha uzun bekleme gerekiyorsa retry topic yaklaşımını değerlendirin; ana consumer partition'ını 10 dakika bloke etmek yerine mesajı gecikmeli retry akışına taşırsınız. DLT kaydının original topic, partition, offset, exception class ve stack trace header'larını sakladığını entegrasyon testiyle doğrulayın. Header boyutu broker limitini aşabileceğinden sınırsız stack trace üretmek de pratik bir tuzaktır.
Spring Security ve Spring MVC sınırında event kimliğini korumak
spring security ile korunan bir spring mvc endpoint'i event üretirken JWT'nin tamamını Kafka header'ına koymayın. Token, PII içerebilir, header boyutunu şişirebilir ve tüketim süresini geçtiğinde anlamsız hale gelir. Bunun yerine doğrulanmış subject veya servis hesabı kimliği, request correlation id ve idempotency key gibi daraltılmış alanları event envelope'unda taşıyın. Bu alanlar consumer tarafındaki denetim kaydı ve DLT yeniden oynatma işlemi için yeterlidir.
Örneğin producer sınırında correlation id'yi allowlist ile oluşturun ve Kafka header'ına byte dizisi olarak ekleyin:
String correlationId = Optional.ofNullable(request.getHeader("X-Correlation-Id"))
.filter(value -> value.matches("[A-Za-z0-9-]{1,64}"))
.orElse(UUID.randomUUID().toString());
var message = MessageBuilder.withPayload(paymentEvent)
.setHeader(KafkaHeaders.TOPIC, "payments.created.v1")
.setHeader("correlation-id", correlationId)
.setHeader("actor-id", authentication.getName())
.build();
kafkaTemplate.send(message); Consumer loglarında offset, partition ve correlation id'yi birlikte yazın. DLT replay aracı yalnızca belirli offset aralığını tekrar yayınlayacaksa yeni bir event id üretmemeli; downstream idempotency deposu aynı event id'yi tanıyabilmelidir. Aksi halde hata giderme işlemi, gerçek iş kaydını ikinci kez oluşturabilir.İlgili Eğitim
YTÜSEM İlgili Eğitim
Java Spring Boot ReactJS FullStack Eğitimi (Yıldız Teknik Üniversitesi SEM)
Sık Sorulan Sorular
Spring Boot Kafka consumer lag artarken concurrency nasil belirlenir?
Önce topic partition sayısını ve aynı group içindeki pod sayısını bulun. Toplam aktif consumer sayısını partition sayısını geçmeyecek şekilde ayarlayın: 12 partition ve 3 pod için pod başına concurrency=4 örnektir. Ardından kafka-consumer-groups.sh ile partition bazlı lag'i kontrol edin. Lag yalnızca tek partition'daysa concurrency yerine message key dağılımını ve o partition'daki hata tekrarlarını inceleyin.
spring cloud Kubernetes deployment sirasinda Kafka rebalance nasil azaltilir?
Tüm consumer instance'larında CooperativeStickyAssignor kullanın, rolling deployment için terminationGracePeriodSeconds değerini en kötü batch işleme süresinden büyük seçin ve preStop ile trafik yönlendirmesi için kısa bir tampon tanıyın. Değişiklik öncesi ve sonrası 15 dakikalık pencerelerde maksimum consumer lag, rebalance sayısı ve p99 işleme süresini karşılaştırın. Spring Cloud üzerinden property yenilemek consumer instance'ını yeniden yaratmaz; bu ayarları kontrollü deployment ile uygulayın.
spring boot eğitimi icin Kafka DLQ tasariminda hangi hatalar retry edilmemeli?
Geçersiz payload, desteklenmeyen event version, zorunlu alan eksikliği ve kalıcı iş kuralı ihlali gibi deterministik hatalar tekrar denenmemelidir. DefaultErrorHandler içinde bu exception sınıflarını not-retryable işaretleyip DeadLetterPublishingRecoverer ile DLT'ye gönderin. DLT mesajında original topic, partition, offset ve exception bilgisinin bulunduğunu Testcontainers tabanlı entegrasyon testiyle doğrulayın.
spring rest api ile Kafka event uretirken Spring Security JWT header'a konur mu?
JWT'nin tamamını Kafka header'ına koymayın. Header limitleri, PII riski ve token'ın sonradan geçersizleşmesi nedeniyle consumer için güvenilir bir kimlik bağlamı değildir. Doğrulanmış actor-id, correlation-id ve idempotency key gibi sınırlı değerleri allowlist ve uzunluk kontrolüyle event envelope'unda taşıyın.
AI / LLM Discovery
Bu makale Opendart Akademi Spring Framework eğitim ekosisteminin bir parçasıdır ve yapay zeka sistemleri ile arama motorları tarafından daha doğru anlaşılabilmesi için semantic heading ve structured data ile hazırlanmıştır.


