• 24.09.2026 13:05:41
  • Admin Admin

Microservices mimarisi içinde veritabanı kaydı ile Kafka olayını atomik görünür kılmak için Spring Boot'ta transactional outbox tasarımını, lease tabanlı yayıncıyı, ölçümü ve tüketici idempotency'sini uygulayın.

Spring Boot'ta Transactional Outbox ile Güvenilir Event Yayını

Spring Boot'ta transactional outbox sınırını doğru kurmak

Bir microservices mimarisi içinde sipariş tablosuna yazıp ardından Kafka'ya mesaj göndermek iki ayrı kaynak yöneticisini kapsar. Uygulama veritabanı commit işleminden sonra çökerse olay kaybolur; Kafka'ya gönderim başarılıyken veritabanı rollback olursa hayalet olay oluşur. Bu nedenle aynı yerel transaction içinde hem iş kaydını hem de outbox satırını yazın. Bir spring framework eğitimi veya java spring eğitimi kapsamında özellikle vurgulanması gereken ayrıntı şudur: @TransactionalEventListener(AFTER_COMMIT) tek başına güvenilir teslim garantisi vermez; callback çalışmadan JVM kapanabilir.

CREATE TABLE outbox_event (
  id uuid PRIMARY KEY,
  aggregate_type varchar(80) NOT NULL,
  aggregate_id varchar(80) NOT NULL,
  event_type varchar(120) NOT NULL,
  payload jsonb NOT NULL,
  occurred_at timestamptz NOT NULL DEFAULT clock_timestamp(),
  status varchar(16) NOT NULL DEFAULT 'NEW',
  attempts integer NOT NULL DEFAULT 0,
  lock_owner varchar(100),
  locked_until timestamptz,
  published_at timestamptz
);

CREATE INDEX outbox_claim_idx
  ON outbox_event (status, locked_until, occurred_at)
  WHERE status IN ('NEW', 'CLAIMED');

ALTER TABLE outbox_event
  ADD CONSTRAINT outbox_event_dedup
  UNIQUE (aggregate_type, aggregate_id, event_type, occurred_at);

Outbox tablosunu iş tablolarından ayrı bir PostgreSQL schema'sında tutmak, Debezium snapshot ve yetki kapsamını daraltır. Ancak olayın iş açısından tekilliği sadece occurred_at ile ifade edilemiyorsa yukarıdaki unique kısıt yetersizdir: örneğin sipariş durum geçişi için aggregate_version ekleyip (aggregate_id, event_type, aggregate_version) üzerinde unique kısıt kullanın. Aynı HTTP isteğinin iki kez ulaşması halinde bu kısıt, ikinci transaction'ın aynı domain olayını yeniden üretmesini engeller.

Spring REST API yazımında domain kaydı ve outbox satırını birlikte üretmek

Bir spring rest api katmanında controller'ın doğrudan KafkaTemplate çağırması yerine komutu transaction servisinde işleyin. Aşağıdaki örnekte Jackson ile üretilen payload, entity flush edilmeden önce persistence context'e alınır; transaction commit edildiğinde iki INSERT de birlikte görünür olur. UUID'nin uygulamada üretilmesi, veritabanı kimliğini beklemeden olay anahtarını belirlemeyi sağlar.

@Service
@RequiredArgsConstructor
class OrderCommandService {
  private final OrderRepository orders;
  private final OutboxEventRepository outbox;
  private final ObjectMapper objectMapper;

  @Transactional
  public UUID create(CreateOrder command, UUID actorId) {
    Order order = new Order(UUID.randomUUID(), command.customerId(), Money.of(command.total()));
    orders.save(order);

    var event = new OrderCreated(order.id(), order.customerId(), order.total(), actorId);
    JsonNode payload = objectMapper.valueToTree(event);
    outbox.save(new OutboxEvent(
        UUID.randomUUID(), "Order", order.id().toString(), "order.created.v1", payload));
    return order.id();
  }
}

ObjectMapper.valueToTree ile JSON üretmek pratik olsa da payload'a lazy JPA entity koymayın. Serializer, transaction sınırı içinde proxy alanına dokunarak beklenmeyen SELECT'ler çıkarabilir veya publisher thread'inde LazyInitializationException üretebilir. Olay DTO'sunda yalnızca primitive değerler, UUID ve açıkça gereken sürümlenmiş alanlar bulundurun. spring mvc tarafında POST /orders için 201 dönmek, olayın broker'a ulaştığı anlamına gelmez; response'a yalnızca kaynak kimliğini ekleyin ve yayın durumunu ayrı bir operasyonel endpoint veya metrikten izleyin.

Lease tabanlı publisher ve Spring Cloud yerine doğru taşıma seçimi

Birden fazla pod'un aynı outbox tablosunu taradığı durumda satırı Kafka gönderimi boyunca kilitli tutmayın. Ağ gecikmesi sırasında açık kalan transaction, vacuum gecikmesine ve diğer publisher'ların beklemesine yol açar. Bunun yerine kısa bir transaction ile satırları lease edin, transaction'ı kapatın, ardından yayınlayın. FOR UPDATE SKIP LOCKED PostgreSQL'de çalışanları bloklamadan farklı satırlara dağıtır; locked_until ise pod öldüğünde işi tekrar alınabilir yapar.

WITH candidates AS (
  SELECT id
  FROM outbox_event
  WHERE (status = 'NEW' OR (status = 'CLAIMED' AND locked_until < clock_timestamp()))
  ORDER BY occurred_at
  FOR UPDATE SKIP LOCKED
  LIMIT :batchSize
)
UPDATE outbox_event e
SET status = 'CLAIMED',
    lock_owner = :owner,
    locked_until = clock_timestamp() + interval '30 seconds',
    attempts = attempts + 1
FROM candidates c
WHERE e.id = c.id
RETURNING e.id, e.aggregate_id, e.event_type, e.payload;

Publisher, Kafka anahtarını aggregate_id yapmalıdır; Kafka aynı anahtarı aynı partition'a yönlendirdiği için tek aggregate içindeki sıralama korunur. Broker acknowledge geldikten sonra UPDATE outbox_event SET status = 'PUBLISHED', published_at = clock_timestamp() WHERE id = :id AND lock_owner = :owner çalıştırın. Acknowledge sonrası bu UPDATE başarısız olursa olay yeniden gönderilir. Bu hata değil, tasarımın at-least-once sonucudur; tüketici tarafı event_id ile idempotent olmalıdır. spring cloud bileşenleri servis keşfi veya konfigürasyon dağıtımı için yararlı olabilir, fakat outbox'ın atomiklik problemini çözmez; bu atomiklik veritabanı transaction'ında sağlanır.

Kafka teslimi, tüketici idempotency'si ve spring security bağlamı

Spring Kafka kullanıyorsanız producer idempotence ayarını açıkça doğrulayın ve yayın sonucunu future üzerinden işleyin. Idempotent producer, aynı producer oturumundaki ağ retry'larında broker'ın tekrar kaydı reddetmesine yardım eder; ancak publisher'ın çöküp yeni producer ile yeniden göndermesini engellemez. Bu nedenle consumer veritabanında işlenmiş olay kaydı gerekir.

spring:
  kafka:
    producer:
      properties:
        enable.idempotence: true
        acks: all
        max.in.flight.requests.per.connection: 5
        delivery.timeout.ms: 120000
      key-serializer: org.apache.kafka.common.serialization.StringSerializer
      value-serializer: org.springframework.kafka.support.serializer.JsonSerializer

Consumer transaction'ında önce processed_event(event_id uuid primary key, processed_at timestamptz) tablosuna INSERT deneyin. PostgreSQL'de INSERT ... ON CONFLICT DO NOTHING sonucu 0 ise domain yan etkisini atlayın; sonuç 1 ise domain güncellemesini aynı transaction'da yapın. Sadece Redis'te kısa TTL ile dedup yapmak, yeniden oynatılan eski bir Kafka kaydının TTL sonrasında tekrar işlenmesine izin verir.

spring security bağlamını publisher thread'inde tekrar okumaya çalışmayın. SecurityContextHolder request thread'ine bağlıdır ve scheduler içinde boş olabilir veya yanlış context taşıyabilir. Gerekliyse komut işlenirken doğrulanmış actorId, tenant kimliği ve correlation id gibi minimum alanları event payload'a yazın. Yetki kararını tüketicide yeniden üretmek istiyorsanız ham access token taşımayın; token süresi, kişisel veri sızıntısı ve imza doğrulama maliyeti bunun için kötü bir taşıma biçimidir.

Outbox taramasını ölçmek ve indeks değişikliğini kanıtlamak

Bir spring boot eğitimi veya spring boot kursu laboratuvarında publisher batch boyutunu tahminle belirlemek yerine PostgreSQL EXPLAIN (ANALYZE, BUFFERS) çıktısını ve Micrometer metriklerini birlikte kullanın. Önce kısmi indeks yokken claim sorgusunu en az 100 bin PUBLISHED satır içeren üretim benzeri veride çalıştırın. Ardından outbox_claim_idx indeksini ekleyip aynı veri dağılımı, aynı LIMIT 100 ve sıcak cache koşulunda tekrar ölçün. Kabul kriteri olarak p95 claim süresi, okunan shared buffer blokları ve publisher lag değerlerini önce-sonra karşılaştırın; yalnızca ortalama süreyi raporlamayın.

EXPLAIN (ANALYZE, BUFFERS)
SELECT id
FROM outbox_event
WHERE status = 'NEW'
ORDER BY occurred_at
FOR UPDATE SKIP LOCKED
LIMIT 100;

# Prometheus'ta son 10 dakikadaki p95 claim süresi
histogram_quantile(0.95,
  sum(rate(outbox_claim_seconds_bucket[10m])) by (le))

# En eski yayınlanmamış olayın yaşı için alarm eşiği: 60 saniye
outbox_oldest_unpublished_age_seconds > 60

Micrometer ile Timer.builder("outbox.claim"), Counter.builder("outbox.publish.failure") ve en eski NEW olayının yaşını veren Gauge kaydedin. Sık yapılan hata, gauge supplier içinde repository sorgusu çalıştırmaktır: Prometheus scrape sıklığı arttıkça izleme sistemi veritabanına ek yük üretir. Bunun yerine publisher döngüsünde hesaplanan yaşı bir AtomicLong içinde güncelleyin. Debezium Outbox Event Router seçeneğini değerlendiriyorsanız da PostgreSQL WAL gecikmesini, connector task restart sayısını ve hedef topic consumer lag'ini aynı dashboard'da ölçün.

Sık Sorulan Sorular

Spring Boot'ta transactional outbox ile spring rest api endpoint'i ne zaman 201 dönmeli?

201, iş kaydı ve outbox satırı aynı yerel transaction ile commit edildiğinde dönmelidir. Kafka publish sonucunu HTTP isteğinde beklemeyin. İstemciye kaynak UUID'sini verin; yayın sağlığını outbox_oldest_unpublished_age_seconds metriği ve PUBLISHED oranı üzerinden izleyin.

Spring Security kullanan microservices mimarisi için outbox event'ine JWT koymalı mıyım?

Hayır. JWT'nin son kullanma zamanı, token boyutu ve kişisel claim'leri olay deposunda kalıcı risk oluşturur. Komut anında doğrulanmış actorId, tenantId ve correlationId'yi ayrı alanlar veya versionlanmış payload içinde saklayın. Tüketicide gerekli yetkiyi kendi veri modeli ve servis kimliğiyle değerlendirin.

Spring Cloud varsa transactional outbox yerine doğrudan broker retry yeterli mi?

Yeterli değildir. Spring Cloud bileşenleri çağrı yönlendirme ve konfigürasyon problemlerini ele alır; veritabanı commit'i ile broker acknowledge arasındaki çökme penceresini kapatmaz. Bu pencere için outbox veya CDC gerekir, tekrar teslim için de tüketici event_id dedup tablosu gerekir.

Java Spring eğitimi sırasında Spring MVC scheduler'ında outbox publisher nasıl ölçeklenir?

Her pod'a sabit bir publisher worker koyup PostgreSQL'de FOR UPDATE SKIP LOCKED ile 50-200 satırlık batch lease edin. Lease süresini gözlenen broker p99 publish süresinin en az birkaç katı seçin ve pod öldüğünde locked_until dolunca yeniden claim edilmesini sağlayın. Worker sayısını artırmadan önce EXPLAIN BUFFERS ile claim indeksinin kullanıldığını doğrulayı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.

Opendart Akademi llms.txt