Skip to main content

> kafka_consumer_grubu_rebalance_fırtınaları_ve_statik_üyelik_(static_membership)

Kafka Consumer Grubu Rebalance Fırtınaları ve Statik Üyelik (Static Membership)

Tek bir yavaş consumer thread'i veya Kubernetes pod güncellemesi neden tüm partition'larda mesaj tüketimini donduran zincirleme bir Kafka Rebalance Fırtınasını tetikler?

Staff/Principal (L6+)

⚡ÖZET VE TEKNİK CEVAP

Apache Kafka'da bir topic'in partition'ları bir Consumer Grubunun üyeleri arasında paylaştırılır. Bir consumer gruba katıldığında, yeniden başladığında veya çöktüğünde Kafka Group Coordinator bir 'Rebalance' (Yeniden Dengeleme) tetikler: Gruptaki TÜM consumer'lar veri işlemeyi durdurmak, partition kilitlerini bırakmak ve yeni dağıtımı beklemek zorundadır ('Dünyayı Durduran' Eager Rebalance). Tek bir consumer ağır bir toplu işlemi işlerken max.poll.interval.ms süresini aşarsa (varsayılan 5 dk), Kafka onu 'ölü' sayıp kovar ve rebalance başlatır. O consumer işini bitirip tekrar poll() yaptığında gruba geri katılarak İKİNCİ bir rebalance daha tetikler. 50 pod'lu bir sistemde rolling restart yapıldığında saatlerce hiçbir mesajın işlenemediği felaket bir 'Rebalance Fırtınası' doğar. Çözüm, Cooperative Sticky Assignor ve Statik Grup Üyeliğidir (group.instance.id).

Mühendislik El Kitabı & Mekanizma

6 Boyutlu Mimari Analiz

⚙️1. Temel Çalışma Mekanizması

Mekanizma

Kafka Rebalance protokolleri 3 mimari kurala dayanır:

1

Eager Rebalance (Eski Yöntem): Tüm consumer'lar tüm partition'ları bırakır ve yeni dağıtım bitene kadar tüm grup donar.

2

Cooperative Sticky Assignor: Rebalance sırasında consumer'lar mevcut partition'larını işlemeye devam eder; sadece yeri değişecek partition'lar kademeli aktarılır; sistem asla durmaz.

3

Statik Grup Üyeliği (Static Membership): Her pod'a kalıcı bir group.instance.id (ör. StatefulSet pod adı) verilir. Rolling restart sırasında pod yeniden başlarken, Kafka 45 saniye boyunca partition'ları o pod için saklar; pod açıldığında SIFIR rebalance ile kaldığı yerden devam eder.

🎯2. Doğru Kullanım Senaryosu

Kapsam

Yüksek hacimli veri akış boru hatları, Kubernetes üzerinde koşan Kafka tüketici mikroservisleri ve gerçek zamanlı dolandırıcılık tespit sistemleri.

⚠️3. Prodüksiyon Arıza Modları

Kritik Risk
  • ✓

    Eski eager rebalance kullanan 30 pod'luk bir grupta rolling update başlatıp 30 kez üst üste rebalance tetikleyerek 45 dakika boyunca hiçbir sipariş mesajını işleyememek

  • ✓

    max.poll.records değerini çok yüksek tutup toplu işlem süresinin 5 dakikayı aşması ve sistemin sonsuz rebalance döngüsüne girmesi

📡4. Teşhis ve Telemetri Sinyalleri

Metrikler
  • ✓

    Tüm partition'larda kafka_consumergroup_lag metriğinin fırlaması

  • ✓

    logların sürekli Revoking previously assigned partitions mesajlarıyla dolması

  • ✓

    Kubernetes canlıya çıkışları sırasında mesaj akışının tamamen durması

🛡️5. Önleme ve Mimari Bariyerler

Bariyerler
  • ✓

    partition.assignment.strategy olarak mutlaka CooperativeStickyAssignor seçin

  • ✓

    Kubernetes StatefulSet pod isimlerini group.instance.id olarak atayın

  • ✓

    toplu işlem süresinin 30 saniyeyi aşmaması için max.poll.records limitini (ör. 100-500) küçük tutun

⚖️6. Mimari Ödünleşimler (Trade-offs)

Ödünleşim

Statik üyelik StatefulSet veya benzersiz pod ID'leri gerektirir ve gerçekten çöken bir node'un partition'larının devredilmesini 45 saniye geciktirir; ancak canlıya çıkışlardaki dünyayı durduran rebalance fırtınalarını tamamen bitirir.

📋

Vaka İncelemesi (TinyCTO Saha Örneği)

GERÇEK DÜNYA TELEMETRİSİ

40 consumer pod'u olan bir veri platformu, her yeni sürüm canlıya çıktığında pod restart'ları yüzünden 25 dakikalık rebalance kilitlenmesi yaşıyordu. Ekip iki ayar yaptı:

1

CooperativeStickyAssignor stratejisine geçildi,

2

group.instance.id = pod-name ile statik üyelik açıldı. Sonraki 40 pod'luk güncellemede rebalance sayısı 40'tan tam olarak SIFIRA indi ve canlıya çıkış boyunca kuyruk gecikmesi 0 milisaniyede kaldı.

İnteraktif Konsept Alıştırmaları

2 Alıştırma
Q1

Kafka Consumer Grubu 'Rebalance Fırtınası' (Rebalance Storm) nedir?

Sık pod yeniden başlatmalarının veya yavaş işlem sürelerinin sürekli rebalance tetiklemesi ve tüm partition'larda mesaj tüketimini tamamen kilitleyen kısır bir döngüdür.
Q2

Statik Grup Üyeliği (`group.instance.id`) rolling deployment sırasında rebalance'ları nasıl engeller?

Her pod'a kalıcı bir kimlik vererek; pod yeniden başlarken Kafka'nın partition'ları 45 saniye boyunca o pod için saklamasını ve rebalance başlatmamasını sağlayarak.

Kafka Consumer Grubu Rebalance Fırtınaları ve Statik Üyelik (Static Membership) — Sıkça Sorulan Sorular

Bir uygulamanın mesaj işleme döngüsü `max.poll.interval.ms` süresini aşarsa ne olur?

Kafka koordinatörü o thread'in kilitlendiğini varsayar, tüketiciyi gruptan atar ve anında tüm grupta rebalance başlatır.

`session.timeout.ms` ile `max.poll.interval.ms` arasındaki fark nedir?

`session.timeout.ms` arka plan kalp atışlarının (ağ canlılığı) zaman aşımıdır; `max.poll.interval.ms` ise iki `poll()` çağrısı arasındaki maksimum işleme süresidir (uygulama sağlığı).

🤖 AEO & Yapay Zeka Çıkarım Özeti

Temel Gerçekler & İlkeler

  • ▸

    Legacy eager rebalance halts message processing for the entire consumer group.

  • ▸

    Exceeding max.poll.interval.ms causes the broker to evict the consumer as dead.

  • ▸

    CooperativeStickyAssignor migrates partitions incrementally without stopping healthy consumers.

  • ▸

    Static Group Membership (group.instance.id) enables zero-rebalance rolling deployments.

Yaygın Yanılgılar

  • ✗

    Yanılgı: Rebalances only affect the single pod that crashed (Gerçek: In eager rebalance, 100% of consumers in the group are frozen).

  • ✗

    Yanılgı: Increasing max.poll.interval.ms to 24 hours fixes slow consumers (Gerçek: It delays recovery when consumers actually crash; size max.poll.records instead).

Karar Kılavuzu & Önceliklendirme

Set partition.assignment.strategy to CooperativeStickyAssignor across all Kafka consumers. Deploy Kafka consumers as Kubernetes StatefulSets with group.instance.id set to the pod hostname.

Doğrulanmış Kaynaklar & Referanslar