İçeriğe geç

Outbox ve Idempotent Consumer — Olay Kaybolmasın, İki Kez de İşlenmesin

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

Önce şunu oku: Servisler Arası İletişim Desenleri

30 saniyede özet

Veritabanına ve mesaj kuyruğuna aynı anda yazamazsın: biri başarılı olup öbürü olmayabilir. Çözüm, mesajı siparişle birlikte veritabanındaki bir 'giden kutusuna' yazmak. Mesaj yine de iki kez gidebilir; alan taraf tekrarı tanımalı.

Muhasebe iki şey buldu: ödemesi hiç alınmamış siparişler ve hiç var olmamış siparişler için alınmış ödemeler. İkisinin de sebebi aynı iki satır: orders.save() ve kafka.send().

  1. Bayt: Muhasebe ödemesi hiç alınmamış siparişler buldu!

  2. Sen: Başka bir şey var mı?

  3. Bayt: Bir de hiç var olmamış siparişler için alınmış ödemeler. İkisi de aynı iki satırdan!

  4. Bayt: Deftere yazıp aynı anda kasaya haber vermek. Arada bir şey ters giderse ne olur?

İki yere aynı anda yazılmaz

Sipariş commit edildi. kafka.send() çağrılmadan hemen önce pod öldü. Yeniden başlayınca ne olur? Cevabı göster

Olay sonsuza kadar kaybolur. Sipariş veritabanında, ama hangi olayın gönderilmediğini gösteren hiçbir kayıt yok. Ödeme servisi bu siparişi hiç duymaz.

Mesajı siparişle aynı masaya, giden kutusuna yaz.
Adım adım oku
  1. Sipariş ve gönderilecek mesaj aynı transaction içinde yazılır: biri orders tablosuna, öteki outbox tablosuna.
  2. Ayrı bir kurye süreci outbox'taki mesajları okuyup Kafka'ya gönderir.
  3. Kurye çökerse mesaj kaybolmaz, kutuda bekler ve kurye dönünce yeniden gönderilir; bu yüzden aynı mesaj iki kez gidebilir.
  4. Alıcı işlediği mesajların kimliğini hatırlar; tekrar gelen mesajı fark eder ve ikinci kez işlemez.

Aynı işlemde iki ayrı sisteme yazmaya dual writeAynı işlemde iki ayrı sisteme (ör. veritabanı ve mesaj kuyruğu) yazmak. Ortak bir transaction olmadığı için biri başarılı olup diğeri başarısız olabilir.Sözlükte gör → denir. Veritabanı transaction’ı Kafka’ya gönderilen mesajı geri alamaz, Kafka da veritabanının commit edip etmediğini bilmez.

Sırayı çevirmek de çözmez: önce gönderip sonra commit edemezsen, olmayan bir sipariş için olay gitmiş olur.

Kafam karıştı, daha basit anlat

İki farklı yere aynı anda yazmaya çalışırsan, biri başarılı olup öteki başarısız olabilir. O zaman iki yer birbirinden farklı şey söyler.

Hızlı kontrolOrta

Bir metot önce siparişi veritabanına kaydedip sonra `kafka.send()` çağırıyor. Temel sorun ne?

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

Nadiren de olsa, veritabanında olmayan siparişler için ödeme alınıyor. Hangi satırlar sorunun kaynağı?

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

Hatalı satıra dokun, sonra kontrol et.

OrderService.java
Java 21UTF-8LF

Giden kutusu: mesajı siparişle birlikte yaz

outboxOlayı, iş verisiyle aynı yerel transaction içinde bir tabloya yazıp ayrı bir süreçle kuyruğa taşıma deseni. İki kaynağa yazma (dual write) problemini çözer.Sözlükte gör → kalıbı iki yazmayı tek sisteme indirir. Olay, siparişle aynı veritabanı transaction’ında bir tabloya yazılır; ya ikisi de var ya hiçbiri.

Ayrı bir relay commit edilmiş satırları okuyup yayınlar: periyodik bir sorguyla ya da CDC ile (Debezium). Relay yayınlayıp “gönderildi” işaretlemeden çökerse, satırı bir kez daha yayınlar.

Yani teslimat at-least-onceMesajın en az bir kez teslim edileceği garantisi; iki kez de gelebilir. Önce işle sonra commit et deseninin sonucudur ve tüketicinin idempotent olmasını gerektirir.Sözlükte gör →’tır: outbox olayın kaybolmamasını sağlar, tekrar gelmemesini değil.

Change
Değişen...
Data
...veriyi...
Capture
...yakalamak. Veritabanının kendi değişiklik günlüğünü okuyup her yeni satırı dışarı iletmek.
Kafam karıştı, daha basit anlat

Mektubu doğrudan postaya verme, önce aynı defterdeki “giden kutusu” sayfasına yaz. Sipariş ve mektup birlikte yazılır ya da hiçbiri yazılmaz. Postacı mektubu sonra alıp götürür.

Hızlı kontrolOrta

Transactional outbox kalıbı dual write sorununu nasıl çözer?

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

Outbox kurduktan sonra ödeme servisinin idempotent olması neden hâlâ gerekir?

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

Kendin gör

Outbox ve idempotent consumer — sipariş ve ödeme tutarlı mı?

Tohum 127827
  1. Oynat ya da adımla.
Hız
Adım 0

Şu an ne oldu?

POST /orders

Sipariş kaydedilecek ve ödeme servisine OrderPlaced olayı gidecek. İki ayrı sistem: veritabanı ve Kafka.

Görevler0/4

  • Ödemesi hiç alınmayan bir sipariş üretaçık

    İpucu

    Sipariş kaydedildi ama olay hiç yola çıkmadı.

  • Olmayan bir sipariş için para çekaçık

    İpucu

    Olay commit'ten önce gönderilirse ve commit başarısız olursa?

  • Aynı sipariş için iki kez para çekaçık

    İpucu

    At-least-once teslimat ve her mesajı işleyen bir tüketici.

  • Bir arızaya rağmen sipariş ve ödeme tutarlı kalsınaçık

    İpucu

    Olay siparişle aynı transaction'da, tüketici tekrarı fark ediyor.

Olay günlüğü (0)

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

  1. Arızayı “commit sonrası süreç ölüyor” yap. Sipariş var, ödeme yok.
  2. “Gönderim sonrası commit başarısız” seç. Olmayan sipariş için ödeme.
  3. Outbox’a geç. İki arıza da tutarlı sonuçlanıyor.
  4. “Consumer offset’i commit etmeden ölüyor” seç. Outbox’la bile iki çekim.
  5. Ödeme servisini idempotent yap. Tek çekim.
Hızlı kontrolOrta

Her arıza senaryosu, dual write ile hangi sonucu doğurur?

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

Sınıflandırılmamış

Sipariş var, olay yok

Veritabanı yazdı, mesaj gitmedi

    Olay var, sipariş yok

    Mesaj gitti, veritabanı geri alındı

      Olay iki kez işlendi

      Aynı mesaj tekrar teslim edildi

        Satır satır: tekrarı tanıyan alıcı

        Aynı olay ikinci kez gelirse

        PaymentListener.java
        1@KafkaListener(topics = "orders")
        2@Transactional
        şu an çalışan satırpublic void on(OrderPlaced event) {
        4 try {
        5 processed.insert(event.eventId()); // UNIQUE(event_id)
        6 } catch (DuplicateKeyException alreadyDone) {
        7 return;
        8 }
        9 payments.charge(event.orderId(), event.amount());
        10}

        Debug

        Adım 1/5

        ilk teslimat evt-42 geldi. Transaction açık.

        event_id
        = evt-42
        Java 21UTF-8LF3:1

        Sol/sağ ok tuşlarıyla da gezebilirsin.

        Bu bir idempotent consumerAynı mesajı birden fazla kez alsa da etkisini yalnızca bir kez uygulayan tüketici. Genelde işlenmiş olay kimliklerini benzersiz bir kısıtla kaydeder.Sözlükte gör →’dır. Tekrarı olay kimliğiyle tanı, offset’le değil: relay aynı olayı yeniden yayınlarsa yeni bir offset alır.

        Hızlı kontrolOrta

        Idempotent bir tüketici tekrar gelen mesajı neyle tanımalı?

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

        Tuzaklar

        @Transactional içinde kafka.send(). Rollback olan sipariş için olay gider.

        AFTER_COMMIT dinleyicisinde send. Hayalet olayı önler, ama commit sonrası çökmede olayı yine kaybeder.

        Sırayı unutmak. Aynı siparişin olayları farklı partition’lara giderse iptal ödemeden önce işlenebilir. Partition anahtarı sipariş kimliği olmalı.

        Exactly-once’a güvenmek. Kafka’nın garantisi Kafka’nın içini kapsar; tüketicinin veritabanı yazmasını değil.

        Hızlı kontrolOrta

        `@TransactionalEventListener(phase = AFTER_COMMIT)` içinde `kafka.send()` yapmak dual write sorununu çözer mi?

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

        Aşağıdaki örnek bir bankadan ve dual write sorununu outbox ile çözüyor: havale ve olayı aynı transaction’da yazılıyor, ayrı bir aktarıcı olayları Kafka’ya taşıyor, limit servisi de aynı olayı iki kez saymıyor.

        Derinleş · Havale olayı: outbox ve idempotent tüketici 5 dosya · ~101 satır · ilk okumada atlayabilirsin
        Proje dosyaları

        src/main/resources/db/migration/ V14__outbox.sql Outbox tablosu ve işlenmiş olaylar tablosu.

        src/main/resources/db/migration/V14__outbox.sql
        CREATE TABLE outbox_events (
        id UUID PRIMARY KEY,
        aggregate_type VARCHAR(50) NOT NULL,
        aggregate_id VARCHAR(50) NOT NULL, -- becomes the Kafka key: per-account ordering
        event_type VARCHAR(100) NOT NULL,
        payload JSONB NOT NULL,
        created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
        published_at TIMESTAMPTZ
        );
        CREATE INDEX outbox_unpublished ON outbox_events (created_at) WHERE published_at IS NULL;
        -- On the consumer side (another service, another database)
        CREATE TABLE processed_events (
        event_id UUID PRIMARY KEY,
        processed_at TIMESTAMPTZ NOT NULL DEFAULT now()
        );

        src/main/java/com/bank/transfer/ TransferService.java Olay, havaleyle aynı transaction'da bir satır olarak yazılır. Kafka'ya burada gidilmez.

        src/main/java/com/bank/transfer/TransferService.java
        @Service
        class TransferService {
        private final Ledger ledger;
        private final OutboxRepository outbox;
        private final ObjectMapper json;
        TransferService(Ledger ledger, OutboxRepository outbox, ObjectMapper json) {
        this.ledger = ledger;
        this.outbox = outbox;
        this.json = json;
        }
        @Transactional
        public long transfer(TransferCommand cmd) {
        Transfer transfer = ledger.post(cmd); // debit and credit postings
        // Same transaction, same database: the money moves and the event is recorded,
        // or neither happens. No kafka.send() here — that would be the dual write
        // this pattern removes: money moved, but no one else ever hears about it.
        outbox.save(OutboxEvent.of("Account", cmd.fromIban(), "TransferCompleted",
        json.valueToTree(TransferCompleted.from(transfer))));
        return transfer.getId();
        }
        }

        src/main/java/com/bank/outbox/ OutboxRepository.java Birden çok pod aynı satırı almasın: FOR UPDATE SKIP LOCKED.

        src/main/java/com/bank/outbox/OutboxRepository.java
        interface OutboxRepository extends JpaRepository<OutboxEvent, UUID> {
        // Each pod locks its own batch; rows locked by another pod are skipped, not waited on.
        @Query(value = """
        SELECT * FROM outbox_events
        WHERE published_at IS NULL
        ORDER BY created_at
        LIMIT :batch
        FOR UPDATE SKIP LOCKED
        """, nativeQuery = true)
        List<OutboxEvent> lockNextBatch(int batch);
        }

        src/main/java/com/bank/outbox/ OutboxRelay.java Aktarıcı: satırları alır, Kafka onaylayınca yayınlandı işaretler. Çökerse satır tekrar gönderilir: en az bir kez.

        src/main/java/com/bank/outbox/OutboxRelay.java
        @Component
        class OutboxRelay {
        private final OutboxRepository outbox;
        private final KafkaTemplate<String, String> kafka;
        OutboxRelay(OutboxRepository outbox, KafkaTemplate<String, String> kafka) {
        this.outbox = outbox;
        this.kafka = kafka;
        }
        @Scheduled(fixedDelay = 500)
        @Transactional
        public void relay() {
        for (OutboxEvent event : outbox.lockNextBatch(100)) {
        var record = new ProducerRecord<>("transfers", event.aggregateId(), event.payload().toString());
        record.headers().add("event-id", event.id().toString().getBytes(StandardCharsets.UTF_8));
        record.headers().add("event-type", event.eventType().getBytes(StandardCharsets.UTF_8));
        // Wait for the broker's ack before marking. A crash between send and commit
        // means the event is sent again: at-least-once, so consumers must be idempotent.
        kafka.send(record).join();
        event.markPublished(Instant.now());
        }
        }
        }

        src/main/java/com/bank/limits/ TransferCompletedConsumer.java Limit servisi: en az bir kez teslimi, işlenmiş olay tablosuyla tam bir kez etkiye çevirir. Yoksa aynı havale günlük limitten iki kez düşer.

        src/main/java/com/bank/limits/TransferCompletedConsumer.java
        @Component
        // The limits service tracks how much of the customer's daily transfer limit is used.
        class TransferCompletedConsumer {
        private final JdbcTemplate jdbc;
        private final DailyLimits limits;
        TransferCompletedConsumer(JdbcTemplate jdbc, DailyLimits limits) {
        this.jdbc = jdbc;
        this.limits = limits;
        }
        @KafkaListener(topics = "transfers", groupId = "limits")
        @Transactional
        public void on(@Payload TransferCompleted event, @Header("event-id") String eventId) {
        int firstTime = jdbc.update(
        "INSERT INTO processed_events (event_id) VALUES (?::uuid) ON CONFLICT DO NOTHING", eventId);
        if (firstTime == 0) return; // duplicate delivery: already counted
        // Same transaction as the marker: if the update fails, the marker rolls back too
        // and the redelivered event is processed again.
        limits.addUsage(event.customerId(), event.valueDate(), event.amount());
        }
        }

        Kendini sına

        Şimşek turu1/5

        Veritabanına ve mesaj kuyruğuna tek bir transaction içinde yazmak mümkündür.

        Soru 1/2İleri

        Outbox relay'i olayları yayınlarken aynı siparişe ait OrderPlaced ve OrderCancelled olaylarının sırası bozuluyor. En olası neden ve çözüm?

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

        Aklında kalacak üç şey

        1. 1 Veritabanına kaydedip ardından mesaj göndermek iki ayrı yazmadır. Arada bir çökme ya mesajı kaybettirir ya da olmayan bir sipariş için mesaj gönderir.
        2. 2 Outbox, mesajı siparişle aynı transaction'da yazar; ayrı bir süreç onu sonra yayınlar. Mesaj kaybolmaz, ama iki kez gidebilir.
        3. 3 Alan taraf mesaj kimliğini, kendi işiyle aynı transaction'da benzersiz bir kısıtla kaydederek tekrarı tanır.
        Sonraki kapı Tarayıcıda bir script'in okuyabildiği anahtarı, sayfaya sızan kötü bir script de okuyabilir mi? BFF'de Güvenlik ve Alternatifler — Token Handler, GraphQL, Gateway · 9 dk

        5 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.