İçeriğe geç

Spring Kafka Hata Yönetimi — Bozuk Mesaj Partition'ı Durdurur mu?

Orta 9 dk Sık karşılaşılır

Önce şunu oku: Kafka — Partition, Consumer Group ve Lag , Outbox ve Idempotent Consumer — Olay Kaybolmasın, İki Kez de İşlenmesin

30 saniyede özet

Kafka'dan gelen bir mesaj işlenemezse Spring onu birkaç kez dener, sonra log'a yazıp atlar: mesaj kaybolur. İşlenemeyen mesajları ayrı bir 'sahipsiz' kuyruğa (DLT) koymak, kaybı önlemenin yoludur.

Tek kasalı bir postanede adresi okunmayan bir mektup yüzünden bütün kuyruk bekliyor. Sonunda memur mektubu sessizce çöpe atıyor.

  1. Bayt: Muhasebe bir günün cirosunda üç sipariş eksik buldu!

  2. Sen: Mesajlar Kafka'da duruyor mu?

  3. Bayt: Duruyor. Log'da da yalnızca üç tane Backoff exhausted satırı var.

  4. Bayt: Okunmayan mektubu çöpe atan bir postane gibi. Hadi kuyruğa bakalım.

Muhasebe, bir günün cirosunda üç siparişin eksik olduğunu buldu. Kafka’da mesajlar duruyordu, DLT diye bir topic yoktu. Log’da yalnızca üç “Backoff exhausted” satırı vardı.

Hiçbir şey ayarlamazsan ne olur?

Listener exception fırlatınca kaydı Spring Kafka’nın DefaultErrorHandler’ı devralır. Varsayılan ayarla kaydı aralıksız dokuz kez daha dener.

Denemeler biterse kayıt log’a yazılır ve offset ilerletilir. Kaydın kendisi hiçbir yerde saklanmaz: mesaj kaybolmuştur.

Bunu önlemek için handler’a bir kurtarıcı verilir. DeadLetterPublishingRecoverer, kaydı hata bilgisiyle birlikte bir dead letter topicİşlenemeyen mesajların, hata bilgisiyle birlikte yazıldığı ayrı topic. Mesaj kaybolmaz; incelenip düzeltildikten sonra yeniden gönderilebilir.Sözlükte gör →’e yazar.

Kafam karıştı, daha basit anlat

Hiçbir şey ayarlamazsan Spring mesajı birkaç kez dener, olmazsa log’a yazıp bir sonrakine geçer. Mesajın kendisi hiçbir yerde saklanmaz, yani kaybolur.

Hızlı kontrolBaşlangıç

Spring Kafka'da @KafkaListener bir exception fırlattı ve sen hiçbir hata yönetimi ayarlamadın. Varsayılan olarak ne olur?

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

Arada bir sipariş hiç işlenmiyor ve DLT'de de yok. Log'da yalnızca bir 'Backoff exhausted' satırı var. Hatalı satır hangisi?

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

Hatalı satıra dokun, sonra kontrol et.

KafkaConfig.java
Java 21UTF-8LF

Tekrar denerken kuyruk bekler

#3 başarısız oldu ve DefaultErrorHandler onu yeniden deniyor. Aynı partition'daki #4, #5 ve #6 ne yapar? Cevabı göster

Bekler. Handler başarısız offset’e geri döner ve aynı kaydı yeniden teslim eder. Bir partition sırayla işlendiği için arkadakiler o kayıt bitene kadar gelmez.

Okunmayan mektup: kuyruk bekler, sonra ya çöpe ya rafa.
Adım adım oku
  1. #3 işlenemiyor ve DefaultErrorHandler onu art arda yeniden deniyor.
  2. Aynı partition'daki #4, #5 ve #6 bu sırada bekler; tek kasalı postane gibi.
  3. Varsayılan ayarda denemeler bitince kayıt log'a yazılıp atlanır: mesaj Kafka'da durur, ama uygulama için kaybolmuştur.
  4. DeadLetterPublishingRecoverer ile kayıt sahipsiz rafına, yani DLT topic'ine yazılır ve kuyruk ilerler.

Bloklayan retry sırayı korur, bedeli lag’dir. @RetryableTopic ise başarısız kaydı gecikmeli retry topic’lerine alır ve partition akmaya devam eder. Bu sefer de o kayıt, arkasındakilerden sonra işlenir.

Kafam karıştı, daha basit anlat

Mesajlar tek sıra hâlinde gelir. Öndeki mesaj yeniden denenirken arkadakiler bekler. Onu kenara ayırırsan sıra akar, ama o mesaj arkadakilerden sonra işlenir.

Hızlı kontrolOrta

DefaultErrorHandler bir kaydı yeniden denerken aynı partition'daki sonraki kayıtlara ne olur?

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

@RetryableTopic ile bloklamayan retry'a geçtin. Ne kazanır, ne kaybedersin?

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

Kendin gör

Bozuk mesaj partition'ı durdurur mu?

Tohum 186651
  1. #1

    sırada

  2. #2

    sırada

  3. #3

    sırada

  4. #4

    sırada

  5. #5

    sırada

  6. #6

    sırada

İşlenme sırası: —

Oynat ya da adımla.

Hız
Adım 0

Şu an ne oldu?

DefaultErrorHandler, varsayılan (9 yeniden deneme, sonra atla)

Altı sipariş aynı partition'da. Üçüncüsü işlenemiyor: arkadakiler bekleyecek mi, mesaj nereye gidecek?

Görevler0/3

  • Bir siparişi sessizce kaybetaçık

    İpucu

    Kalıcı hata, kurtarıcısı olmayan handler.

  • Poison pill'i boşa deneme yapmadan DLT'ye gönderaçık

    İpucu

    Exception'ı sınıflandır.

  • Geçici hatayı sırayı bozmadan atlataçık

    İpucu

    Bloklayan retry, doğru sınıflandırma.

Olay günlüğü (0)

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

  1. Varsayılanla oynat. Poison pill on kez denendi, sonra kayboldu.
  2. Handler’ı “DLT” yap. Kayıt kurtuldu, ama boşa iki deneme yapıldı.
  3. “Yeniden denenmez”i aç. İlk hatada DLT’ye gitti.
  4. Hatayı “Geçici” yap ve “yeniden denenmez”i kapat. Üçüncü denemede geçti, sıra korundu.
  5. Handler’ı “@RetryableTopic” yap. Partition aktı, sıra bozuldu.
Hızlı kontrolOrta

Her hatayı nasıl ele alınması gerektiğine göre ayır.

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

Sınıflandırılmamış

Back-off ile yeniden dene

Bir süre sonra kendiliğinden düzelebilir.

    Hemen DLT

    Aynı veriyle hiçbir zaman geçmez.

      Asla düzelmeyecek mesaj

      poison pillKaç kez denenirse denensin işlenemeyecek bir mesaj: bozuk format, geçersiz veri. Yeniden denemek yalnızca arkasındaki mesajları bekletir.Sözlükte gör →, kaç kez denenirse denensin işlenemeyecek bir kayıttır: çözülemeyen bayt, geçersiz veri, tanınmayan olay tipi. Onu yeniden denemek yalnızca arkadaki kayıtları bekletir.

      addNotRetryableExceptions(...) ile bu hataları sınıflandır; ilk hatada kurtarıcıya giderler. Geçici hatalar ise (503, deadlock) back-off ile yeniden denenir.

      Kafam karıştı, daha basit anlat

      Bozuk bir mesajı kaç kez denersen dene işlenmez. Onu hemen ayrı bir kutuya koy ki arkadaki sağlam mesajlar beklemesin.

      Hızlı kontrolOrta

      Kayıt geçersiz bir JSON içeriyor ve doğrulama her seferinde başarısız oluyor. En doğru yaklaşım hangisi?

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

      Tuzaklar

      Çözme hatası listener’a ulaşmaz. Deserializer’ı ErrorHandlingDeserializer ile sar; yoksa bozuk kayıt poll sırasında patlar ve handler onu göremez.

      Tekrarlanan yan etki. Her deneme listener’ı baştan çalıştırır. Hatadan önce gönderilen e-posta her denemede yeniden gider; listener idempotent olmalı.

      Uzun back-off. Bloklayan retry’ın toplam beklemesi max.poll.interval.ms’yi aşarsa consumer gruptan düşer ve rebalance başlar.

      Aşağıdaki örnek bir bankanın kart harcamalarını muhasebe defterine işleyen servisinden ve dersin bütün kurallarını birlikte uyguluyor: çözme hataları handler’a ulaşıyor, kalıcı hatalar denenmeden DLT’ye gidiyor, geçici hatalar back-off ile tekrar deneniyor ve listener idempotent.

      Derinleş · Kart harcamalarını deftere işleyen servis: uçtan uca 6 dosya · ~140 satır · ilk okumada atlayabilirsin
      Proje dosyaları

      src/main/java/com/bank/ledger/kafka/ KafkaErrorConfig.java Hata yönetimi tek yerde: back-off, yeniden denenmez hatalar ve DLT'ye yazan kurtarıcı.

      src/main/java/com/bank/ledger/kafka/KafkaErrorConfig.java
      @Configuration
      class KafkaErrorConfig {
      // Boot wires a CommonErrorHandler bean into its listener container factory.
      @Bean
      DefaultErrorHandler kafkaErrorHandler(KafkaTemplate<Object, Object> dltTemplate) {
      var recoverer = new DeadLetterPublishingRecoverer(dltTemplate,
      // Same partition number on the DLT: create card-transactions-dlt with at least as many partitions.
      (record, ex) -> new TopicPartition(record.topic() + "-dlt", record.partition()));
      var backOff = new ExponentialBackOffWithMaxRetries(3); // 1s, 2s, 4s, then recover
      backOff.setInitialInterval(1_000);
      backOff.setMultiplier(2.0);
      var handler = new DefaultErrorHandler(recoverer, backOff);
      // Data that can never succeed goes straight to the DLT: no retries, no blocked partition.
      handler.addNotRetryableExceptions(InvalidTransactionException.class, MethodArgumentNotValidException.class);
      return handler;
      }
      // The DLT receives raw bytes for records that failed to deserialize,
      // and CardTransactionCaptured objects for everything else.
      @Bean
      KafkaTemplate<Object, Object> dltTemplate(ProducerFactory<Object, Object> producerFactory) {
      var serializer = new DelegatingByTypeSerializer(Map.of(
      byte[].class, new ByteArraySerializer(),
      CardTransactionCaptured.class, new JsonSerializer<>()));
      var factory = new DefaultKafkaProducerFactory<Object, Object>(
      producerFactory.getConfigurationProperties(), new StringSerializer(), serializer);
      return new KafkaTemplate<>(factory);
      }
      }

      src/main/java/com/bank/ledger/kafka/ CardTransactionListener.java Asıl iş. Önce idempotency (aynı harcama iki kez borç yazılmaz), sonra doğrulama, sonra muhasebe kaydı.

      src/main/java/com/bank/ledger/kafka/CardTransactionListener.java
      @Component
      class CardTransactionListener {
      private final ProcessedEvents processed;
      private final Ledger ledger;
      CardTransactionListener(ProcessedEvents processed, Ledger ledger) {
      this.processed = processed;
      this.ledger = ledger;
      }
      @KafkaListener(topics = "card-transactions", groupId = "ledger")
      @Transactional
      public void on(CardTransactionCaptured event) {
      // 1. Idempotency first: a retry or a redelivery must not debit the card twice.
      if (!processed.markIfFirst(event.eventId())) return;
      // 2. Validation: a permanent problem, classified as not retryable.
      // Retrying a negative amount or an unknown currency will never make it valid.
      if (event.amount().signum() <= 0 || !SupportedCurrencies.contains(event.currency())) {
      throw new InvalidTransactionException("unpostable transaction " + event.transactionId());
      }
      // 3. The side effect: a double-entry posting (card account debit, merchant
      // settlement credit). A transient failure here (DB deadlock) is retried;
      // the transaction rolls back, so the processed-marker is rolled back too.
      ledger.post(event.cardAccountId(), event.merchantId(), event.amount(), event.currency());
      }
      }

      src/main/java/com/bank/ledger/kafka/ ProcessedEvents.java Aynı olayı ikinci kez işlememek için unique kısıtlı bir tablo.

      src/main/java/com/bank/ledger/kafka/ProcessedEvents.java
      @Repository
      class ProcessedEvents {
      private final JdbcTemplate jdbc;
      ProcessedEvents(JdbcTemplate jdbc) {
      this.jdbc = jdbc;
      }
      /**
      * Returns true only for the first delivery of an event.
      * Table: processed_events(event_id UUID PRIMARY KEY, processed_at TIMESTAMP NOT NULL)
      */
      boolean markIfFirst(UUID eventId) {
      int inserted = jdbc.update("""
      INSERT INTO processed_events (event_id, processed_at)
      VALUES (?, now())
      ON CONFLICT (event_id) DO NOTHING
      """, eventId);
      return inserted == 1;
      }
      }

      src/main/java/com/bank/ledger/kafka/ CardTransactionDltListener.java DLT'yi dinler: hatayı ve kaynağı log'lar, mutabakat ekibi için alarm üretir. Otomatik yeniden oynatmaz.

      src/main/java/com/bank/ledger/kafka/CardTransactionDltListener.java
      @Component
      class CardTransactionDltListener {
      private static final Logger log = LoggerFactory.getLogger(CardTransactionDltListener.class);
      private final MeterRegistry meters;
      CardTransactionDltListener(MeterRegistry meters) {
      this.meters = meters;
      }
      // Raw bytes on purpose: a record that could not be deserialized must still be readable here.
      @KafkaListener(topics = "card-transactions-dlt", groupId = "ledger-dlt",
      properties = "value.deserializer=org.apache.kafka.common.serialization.ByteArrayDeserializer")
      public void on(ConsumerRecord<String, byte[]> record) {
      String error = header(record, KafkaHeaders.DLT_EXCEPTION_MESSAGE);
      String source = header(record, KafkaHeaders.DLT_ORIGINAL_TOPIC);
      log.error("dead letter from {} key={} error={}", source, record.key(), error);
      meters.counter("ledger.dlt.records", "source", String.valueOf(source)).increment();
      // No automatic replay: an unposted card transaction means the customer's balance
      // and the scheme's settlement file will disagree. The reconciliation team fixes
      // the data or the code, then replays on purpose.
      }
      private static String header(ConsumerRecord<?, ?> record, String name) {
      Header header = record.headers().lastHeader(name);
      return header == null ? null : new String(header.value(), StandardCharsets.UTF_8);
      }
      }

      src/main/java/com/bank/events/ CardTransactionCaptured.java Olayın kendisi: değişmez bir record.

      src/main/java/com/bank/events/CardTransactionCaptured.java
      public record CardTransactionCaptured(
      UUID eventId,
      String transactionId, // the authorization's reference, for reconciliation
      long cardAccountId,
      String merchantId,
      BigDecimal amount,
      String currency,
      Instant occurredAt) {
      }

      src/main/resources/ application.yml Çözme hataları ErrorHandlingDeserializer ile sarılıyor; aksi hâlde handler bozuk kaydı hiç göremez.

      src/main/resources/application.yml
      spring:
      kafka:
      bootstrap-servers: kafka:9092
      consumer:
      group-id: ledger
      auto-offset-reset: earliest
      key-deserializer: org.apache.kafka.common.serialization.StringDeserializer
      # The wrapper catches deserialization failures and hands them to the error handler.
      value-deserializer: org.springframework.kafka.support.serializer.ErrorHandlingDeserializer
      properties:
      spring.deserializer.value.delegate.class: org.springframework.kafka.support.serializer.JsonDeserializer
      spring.json.trusted.packages: com.bank.events
      spring.json.value.default.type: com.bank.events.CardTransactionCaptured
      # Retries block the partition; keep total back-off well below this.
      max.poll.interval.ms: 300000
      producer:
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

      Kendini sına

      Şimşek turu1/5

      Varsayılan DefaultErrorHandler, denemeler bitince kaydı log'a yazıp atlar.

      Soru 1/3İleri

      Listener önce bir e-posta gönderiyor, sonra veritabanına yazarken exception alıyor. Kayıt 3 kez deneniyor ve üçüncüsünde geçiyor. Kullanıcı ne yaşar?

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

      Aklında kalacak üç şey

      1. 1 Ayar yapılmamış DefaultErrorHandler, denemeler bitince mesajı yalnızca log'a yazar. Kaybedilmemesi gereken mesajlar bir dead letter topic'e gitmelidir.
      2. 2 Bekleyerek tekrar denemek sırayı korur ama arkadaki mesajları durdurur; @RetryableTopic akışı sürdürür ama sırayı bozar. İkisi aynı anda olmaz.
      3. 3 Hiç düzelmeyecek hatalar tekrar denenmez. Onları denemek yalnızca arkadaki mesajları bekletir.
      Sonraki kapı Yeni sürüm çıkarken yarım kalan istekler ne oluyor? Graceful Shutdown — Deploy Sırasında Kaç İstek Düşer? · 8 dk

      5 kart sonraki derste seni bekliyor

      0/5 kart bu dersten toplandı