İçeriğe geç

Kafka — Partition, Consumer Group ve Lag

İleri 10 dk Çok sık karşılaşılır

30 saniyede özet

Kafka'da bir konu, kasalara benzeyen bölümlere (partition) ayrılır ve kaç kasa varsa o kadar iş aynı anda yapılır. Kasadan fazla kasiyer eklemek hızlandırmaz, fazlası boş bekler.

Kafka’yı bir kuyruk sanmak çok doğal ama yanlış. Kafka, sayfaları hiç yırtılmayan ve birkaç bölüme ayrılmış bir kayıt defteridir: mesaj okununca değil, saklama süresi dolunca silinir.

  1. Bayt: Mesajlar birikiyor! Okuyucu sayısını ikiye katladım.

  2. Sen: Birikme durdu mu?

  3. Bayt: Yeni okuyucuların yarısı hiçbir şey yapmadan bekliyor...

  4. Bayt: Dört kasalı bir markete altı kasiyer koyarsan ne olur?

Partition kasa, consumer kasiyer.
Adım adım oku
  1. Dört partition'ı iki consumer okuyor; üretim tüketimden hızlı olduğu için bekleyen mesajlar (lag) birikiyor.
  2. Consumer sayısı partition sayısına çıkınca her kasaya bir kasiyer düşer ve birikme erir.
  3. Partition sayısından fazla consumer eklemek hızlandırmaz; fazladan gelenler boşta bekler.
  4. Aynı anahtarlı mesajlar hep aynı partition'a gider; bu yüzden birbirlerine göre sıraları korunur.

Partition: her şeyin merkezi

Bir topic, N adet partitionBir Kafka topic'inin, mesajların sırayla eklendiği ve ayrı ayrı tüketilebilen bölümlerinden biri.Hem paralelliğin hem sıralamanın birimi budur: paralellik partition sayısıyla sınırlı, sıralama yalnızca partition içinde garanti. İkisi aynı şeye bağlı olduğu için birbiriyle yarışır.Sözlükte gör →’dan oluşur. Her partition bağımsız, sıralı, append-only bir log’dur ve her mesajın o log içinde bir offset’i vardır.

Üç kritik sonuç:

  • Sıralama yalnızca partition içindedir. Topic genelinde sıralama diye bir şey yoktur.
  • Paralellik tavanı partition sayısıdır. Bir consumer group’ta bir partition’ı aynı anda yalnızca bir consumer okuyabilir.
  • Partition sayısını artırmak kolay, azaltmak imkânsızdır. Ve artırmak, anahtarların dağılımını bozar.
Kafam karıştı, daha basit anlat

Topic bir süpermarketse, partition’lar onun kasalarıdır. Her kasada sıra korunur. Bir kasaya aynı anda yalnızca bir kasiyer bakar, bu yüzden kasadan fazla kasiyer boşta bekler.

Kendin gör

Varsayılan ayarda 4 partition, 2 consumer var ve üretim tüketimden hızlı — lagÜretilen ile tüketilen offset arasındaki fark. Artıyor ve hata yoksa sorun hız farkındadır — ve paralelliğin tavanı partition sayısıdır.Sözlükte gör → büyüyor.

Kafka — partition, consumer group, rebalance ve lag

Tohum 17
grup çalışıyorüretim 6 / kapasite 4

Partition’lar · toplam lag 0 · zirve 0

  • P0C0offset 0/0 · lag 0
  • P1C0offset 0/0 · lag 0
  • P2C1offset 0/0 · lag 0
  • P3C1offset 0/0 · lag 0

Consumer group (2)

  • C0

    P0 P1

    0 işlendi

  • C1

    P2 P3

    0 işlendi

Üretilen
0
Tüketilen
0
Rebalance
0
Hız
Adım 0

Şu an ne oldu?

Grup yetişiyor — lag yok

4 partition, 2 consumer arasında range stratejisiyle paylaştırılmış durumda.

Aklında kalsın: Sıralama garantisi yalnızca partition içindedir. Bir varlığın olaylarının sırası önemliyse, o varlığın kimliğini mesaj anahtarı yaparsın.

Görevler0/3

  • Hiç partition alamayan bir consumer oluşturaçık

    İpucu

    Consumer sayısını partition sayısının üstüne çıkar. Bir partition’ı aynı grupta en fazla bir consumer okuyabilir.

  • Bir rebalance başlat ve bitmesini izleaçık

    İpucu

    Çalışırken "Consumer ekle" veya "Consumer düşür"e bas. Rebalance süresince tüketim durur, üretim durmaz.

  • Adım başına 6+ mesajda lag’i sıfıra indiraçık

    İpucu

    Toplam kapasite (consumer × hız) üretimi geçmeli. Ama partition’dan fazla consumer işe yaramaz, bir consumer’a düşen partition’lar da dengeli olmalı.

Olay günlüğü (0)

Henüz olay yok. Oynat veya adımla.

Sırayla dene:

  • Consumer ekle’ye iki kez bas. Önce rebalanceConsumer group'a üye girip çıktığında partition'ların yeniden dağıtılması. Klasik biçiminde tüm grup kısa süre durur — kapasite eklemenin görünmeyen bedeli budur.Sözlükte gör → yüzünden lag artıyor, sonra kapasite yettiği için eriyor. Kapasite eklemenin kısa vadeli bedeli budur.
  • Consumer sayısını 6 yap, partition 4 kalsın. İki consumer boşta etiketiyle duruyor. Hiçbir katkıları yok, sadece grup üyesi olarak rebalance’a katılıyorlar.
  • Trafik patlaması’na bas. Lag birden tırmanıyor; kapasite üretimin üstündeyse eriyor, altındaysa hiç kapanmıyor.
  • Consumer düşür. Düşen consumer’ın partition’ları sahipsiz kalıyor, rebalance başlıyor, o süre boyunca kimse iş yapmıyor.
  • Stratejiyi round-robin yap. 4 partition / 3 consumer dağılımı değişiyor — range bir consumer’a iki partition verirken round-robin daha düzgün dağıtıyor.
Hızlı kontrolOrta

Partition tam olarak nedir?

Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

Birikme basit bir hesaptır

toplam kapasite = consumer sayısı × consumer başına hız

Kapasite üretimden küçükse lag doğrusal olarak büyür. Bu bir hata değil, matematiktir.

Kafam karıştı, daha basit anlat

Müşteriler kasiyerlerin çalışabildiğinden hızlı geliyorsa sıra büyür. Bu bir arıza değil, basit bir hesap: ya kasiyer ekle ya da her birini hızlandır.

Hızlı kontrolOrta

8 partition'lı bir topic'i okuyan consumer group'a 12 consumer eklersen ne olur?

Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

Rebalance — kapasite eklemenin bedeli. Gruba bir consumer katıldığında veya bir consumer düştüğünde Kafka partition’ları yeniden dağıtır. Klasik (eager) protokolde bu stop-the-world’dür: tüm consumer’lar partition’larını bırakır, yeni atama yapılır, sonra herkes devam eder.

Bu sırada üreticiler durmaz. Lag tırmanır.

consumer eklendi
│
├─ tüm consumer'lar durdu ← lag tırmanmaya başlar
├─ koordinatör yeni atamayı hesapladı
└─ consumer'lar yeni partition'larıyla devam ediyor

Kafka 2.4 ile gelen cooperative-sticky protokolü bunu düzeltir: yalnızca gerçekten el değiştiren partition’lar durur, geri kalan consumer’lar çalışmaya devam eder.

KafkaConsumerConfig.java
props.put(
ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
CooperativeStickyAssignor.class.getName() // eager rebalance'ın duraklamasını kaldırır
);
Hızlı kontrolİleri

Consumer lag artıyor ama consumer'lar hata vermiyor. En olası sebep nedir?

Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

Anahtar ve sıralama

Mesajın anahtarı varsa partition hash(key) % partitionCount ile seçilir. Aynı anahtar her zaman aynı partition’a gider, dolayısıyla o anahtarın mesajları sıralı işlenir.

Sıralama gereken durum
// Bir siparişin olayları hep aynı partition'a düşsün diye anahtar = orderId
producer.send(new ProducerRecord<>("orders", order.id(), event));

Anahtar yoksa mesajlar partition’lara dağıtılır ve sıralama garantisi kalmaz.

Topic 4 partition'lı ve anahtar orderId. Trafik arttı, partition sayısını 8'e çıkardın. Sipariş 42'nin olayları hâlâ sıralı işlenir mi? Cevabı göster

Geçiş anında hayır. hash(key) % 4 ile hash(key) % 8 farklı sonuç verebilir; sipariş 42 başka bir partition’a taşınabilir. Eski olayları eski partition’da, yenileri yenisinde kalır ve iki consumer onları aynı anda işleyebilir.

Offset commit — teslimat garantisini belirleyen yer.

StratejiNe zaman commitGarantiRisk
Otomatik (enable.auto.commit=true)ZamanlayıcıylaBelirsizİşlenmeden commit → mesaj kaybı
İşlemden sonra manuelİş bitinceat-least-onceÇökme → tekrar işleme
İşlemden önce manuelİş başlamadanat-most-onceÇökme → kayıp

Pratikte doğru cevap neredeyse her zaman at-least-once + idempotent tüketici’dir. exactly-once yalnızca Kafka’dan okuyup Kafka’ya yazan akışlarda mümkündür; araya bir veritabanı girince yine idempotency’ye dönersin.

Aşağıda üçü bir arada: olaylar hesap numarasıyla anahtarlanıyor, offset iş bitince commit ediliyor ve tekrar gelen olay bakiyeyi ikinci kez değiştirmiyor.

Derinleş · Kart harcamaları: anahtar, commit ve tekrar gelen olay 5 dosya · ~67 satır · ilk okumada atlayabilirsin
Proje dosyaları

card-ledger/src/main/resources/ application.yml Tüketici ayarları: otomatik commit kapalı, her kayıt işlenince commit. cooperative-sticky ile rebalance'ta yalnızca el değiştiren partition'lar durur.

card-ledger/src/main/resources/application.yml
spring:
kafka:
consumer:
group-id: card-balance-projection
enable-auto-commit: false # commit only after the balance is written
auto-offset-reset: earliest
max-poll-records: 100 # keep one poll well inside max.poll.interval.ms
isolation-level: read_committed
properties:
partition.assignment.strategy: org.apache.kafka.clients.consumer.CooperativeStickyAssignor
listener:
ack-mode: record # at-least-once: commit each record once it is processed

card-auth/src/main/java/com/bank/card/ CardTransactionPublisher.java Üretici: anahtar hesap numarası. Bir hesabın bütün olayları aynı partition'a, dolayısıyla aynı sıraya düşer.

card-auth/src/main/java/com/bank/card/CardTransactionPublisher.java
@Component
class CardTransactionPublisher {
private final KafkaTemplate<String, CardTransactionEvent> kafka;
CardTransactionPublisher(KafkaTemplate<String, CardTransactionEvent> kafka) {
this.kafka = kafka;
}
void publish(CardTransactionEvent event) {
// Key = account: a purchase and its later reversal land on the same partition,
// so they are consumed in the order they happened.
kafka.send("card-transactions", event.accountId(), event);
}
}

card-ledger/src/main/java/com/bank/card/ BalanceProjectionListener.java Tüketici: olayı ve 'işlendi' kaydını tek transaction'da yazar. Aynı olay tekrar gelirse bakiye değişmez.

card-ledger/src/main/java/com/bank/card/BalanceProjectionListener.java
@Component
class BalanceProjectionListener {
private final ProcessedEvents processed;
private final AvailableBalances balances;
BalanceProjectionListener(ProcessedEvents processed, AvailableBalances balances) {
this.processed = processed;
this.balances = balances;
}
// The offset is committed after this method returns. If the pod dies before that,
// the event comes again, and the processed_event row makes the second run a no-op.
@KafkaListener(topics = "card-transactions")
@Transactional
void on(CardTransactionEvent event) {
if (!processed.markIfNew(event.eventId())) {
return; // already applied: a redelivery after a crash or a rebalance
}
switch (event.type()) {
case PURCHASE -> balances.hold(event.accountId(), event.amount());
case REVERSAL -> balances.release(event.accountId(), event.amount());
}
}
}

card-ledger/src/main/resources/db/migration/ V3__processed_event.sql İşlenen olaylar tablosu: olay kimliği birincil anahtar. Tekrar gelen olayı veritabanı reddeder.

card-ledger/src/main/resources/db/migration/V3__processed_event.sql
CREATE TABLE processed_event (
event_id UUID PRIMARY KEY, -- a redelivered event cannot be inserted twice
processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
-- markIfNew runs:
-- INSERT INTO processed_event (event_id) VALUES (?) ON CONFLICT DO NOTHING
-- and treats "1 row inserted" as new, "0 rows" as already processed.

counter-example/ RandomKeyPublisher.java Karşı örnek: anahtar işlem kimliği. Aynı hesabın harcaması ve iadesi farklı partition'lara düşer; iade harcamadan önce işlenebilir.

counter-example/RandomKeyPublisher.java
// Counter-example: keyed by the transaction, not the account.
void publish(CardTransactionEvent event) {
// A purchase and its reversal have different transaction ids, so they hash to
// different partitions. The reversal can be consumed first and release a hold
// that does not exist yet; the purchase then holds money that was already refunded.
kafka.send("card-transactions", event.transactionId(), event);
}
Kafam karıştı, daha basit anlat

Aynı siparişin bütün olayları aynı kasaya gitsin istiyorsan, anahtar olarak sipariş numarasını ver. Aynı anahtar hep aynı kasaya düşer.

Hızlı kontrolİleri

Her ifadeyi doğru teslimat garantisine yerleştir.

Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

Sınıflandırılmamış

at-least-once

Mesaj kaybolmaz ama tekrar işlenebilir

    at-most-once

    Tekrar işlenmez ama kaybolabilir

      Tuzaklar: kaç partition seçmeli?

      Fazla partition da bedava değil: her partition broker’da dosya tanıtıcısı ve bellek tüketir, lider seçimi süresini uzatır ve uçtan uca gecikmeyi artırır.

      Pratik yaklaşım: hedef throughput’u tek bir partition’ın ölçülen throughput’una böl ve büyüme payı ekle. Azaltamayacağın için biraz cömert davran, ama “ne olur ne olmaz” diye yüzlerce partition açmak gerçek bir maliyettir.

      Kendini sına

      Şimşek turu1/5

      Kafka'da bir mesaj okunduğu anda silinir.

      Soru 1/3İleri

      8 partition'lı bir topic'i okuyan consumer group'ta 12 consumer var. Ne olur?

      Cevabı biliyor musun?Önce birini seç. Tekrar zamanlaması buna göre ayarlanıyor.

      Aklında kalacak üç şey

      1. 1 Aynı anda çalışabilecek okuyucu sayısının tavanı partition sayısıdır. Fazla okuyucu boşta bekler.
      2. 2 Okuyucular yeniden dağıtılırken (rebalance) grup bir süre durur, mesajlar ise birikmeye devam eder. cooperative-sticky ayarı bu duruşu büyük ölçüde azaltır.
      3. 3 Sıra garantisi bütün konuda değil, tek bir partition içindedir. Bir şeyin olayları sıralı gelsin istiyorsan onun kimliğini mesaj anahtarı yaparsın.
      Sonraki kapı Kısa kodu adresin hash'inden kesip alıyorsun. Milyarlarca linkte bir gün ne olur? URL Kısaltıcı Tasarımı — Parçaları Tek Sistemde Birleştirmek · 10 dk

      4 kart sonraki derste seni bekliyor

      0/5 kart bu dersten toplandı

      Bu dersin üstüne kurulanlar

      Bunlar bu dersi temel alıyor; hazır olduğunda devam edebilirsin.