Java backend geliştirme servislerinde veritabanı ve Kafka arasındaki çift yazma hatasını transactional outbox ile çözün. Spring Boot eğitimi için şema, kilitleme, CDC ve ölçüm odaklı bir uygulama rehberi.
Spring Boot'ta Transactional Outbox ile Kafka Teslimat Tasarımı
Spring Framework ile çift yazma hatasını doğru modellemek
Bir sipariş kaydedip ardından Kafka'ya olay gönderen kodda iki ayrı dayanıklılık sınırı vardır: PostgreSQL işlemi ve Kafka broker onayı. İlk işlem commit olduktan sonra JVM kapanırsa olay kaybolur; Kafka gönderimi başarılı olup DB rollback olursa hayali bir sipariş olayı oluşur. Spring Framework içindeki @Transactional bu iki kaynağı atomik hale getirmez. Transactional outbox yaklaşımında iş verisi ile olay kaydı aynı yerel SQL transaction'ında yazılır; yayınlama daha sonra, tekrar denenebilir bir akışta yapılır.
CREATE TABLE outbox_event (
id UUID PRIMARY KEY,
aggregate_type VARCHAR(80) NOT NULL,
aggregate_id UUID NOT NULL,
event_type VARCHAR(160) NOT NULL,
aggregate_version BIGINT NOT NULL,
payload JSONB NOT NULL,
headers JSONB NOT NULL DEFAULT '{}'::jsonb,
created_at TIMESTAMPTZ NOT NULL DEFAULT clock_timestamp(),
published_at TIMESTAMPTZ,
lease_until TIMESTAMPTZ,
attempt_count INTEGER NOT NULL DEFAULT 0,
last_error TEXT
);
CREATE UNIQUE INDEX uq_outbox_aggregate_version
ON outbox_event(aggregate_type, aggregate_id, aggregate_version);
CREATE INDEX ix_outbox_claim
ON outbox_event(created_at)
WHERE published_at IS NULL;Buradaki aggregate_version ayrıntısı önemlidir. Aynı aggregate için iki olayın sırasını tüketicide denetlemeye imkan verir. Sadece UUID üretmek deduplikasyon sağlar, sıralama sağlamaz. PostgreSQL'de kısmi indeks, yayımlanmış milyonlarca satırı taramadan bekleyen kayıtları bulur. JSONB payload esneklik sağlar; ancak sorgulanacak alanları JSONB'den okumak yerine ayrı kolon veya üretilmiş kolon olarak modelleyin. Aksi halde olay kuyruğu büyüdüğünde GIN indeksinin yazma maliyeti, yayınlayıcının basit claim sorgusuna gereksiz yük getirir.
Spring Data JPA ve Hibernate ORM ile atomik outbox yazımı
Hibernate ORM flush sırası, servis metodunda iki repository çağrısı bulunmasının atomiklik garantisi olmadığı anlamına gelmez; ikisi de aynı Spring transaction'ına katılıyorsa tek commit'te kalıcı olur. Fakat domain olayını sadece @TransactionalEventListener(phase = AFTER_COMMIT) ile Kafka'ya göndermek outbox değildir: commit sonrası süreç çökerse olay tekrar üretilemez. Olayın kendisini tabloya yazın, yayın sinyalini değil.
@Service
@RequiredArgsConstructor
class OrderApplicationService {
private final OrderRepository orders;
private final OutboxEventRepository outbox;
@Transactional
public UUID place(PlaceOrder command) {
Order order = Order.place(command.customerId(), command.lines());
orders.save(order);
var event = new OrderPlaced(order.getId(), order.getVersion(), order.total());
outbox.save(OutboxEvent.pending(
UUID.randomUUID(),
"Order",
order.getId(),
"order.placed.v1",
order.getVersion(),
event));
return order.getId();
}
}Spring Data JPA tarafında save() çağrısının SQL'i hemen çalıştıracağı varsayımına dayanmayın. Hibernate flush çoğunlukla commit öncesinde olur. Aynı transaction içinde üretilen order ID'si için UUID veya uygulama tarafında atanmış kimlik kullanmak, IDENTITY stratejisinin erken insert zorlamasını engeller. Ayrıca order.getVersion() değerini JPA'nın artıracağı sürümden önce mi sonra mı aldığınızı test edin. Yeni aggregate için sürümü uygulama seviyesinde 1 olarak başlatmak, update olaylarında ise managed entity flush edilmeden önce beklenen sürümü açıkça hesaplamak daha güvenlidir.
Bu tasarım java programlama eğitimi kapsamında transaction sınırlarını öğretmek için iyi bir örnektir; üretimde kritik nokta repository soyutlaması değil, event satırının business row ile aynı bağlantı ve aynı commit kaydında olmasıdır. Bir java kursu projesinde H2 ile geçen test, PostgreSQL kilit ve JSONB davranışını doğrulamaz. Testcontainers PostgreSQL ile rollback testi yazın: order insert'inden sonra istisna atın ve hem orders hem outbox_event sayısının sıfır olduğunu doğrulayın.
SKIP LOCKED, lease ve idempotent Kafka tüketicisi
Birden fazla uygulama pod'u outbox taradığında FOR UPDATE SKIP LOCKED bekleyen satırları paylaşır. Ancak Kafka'ya gönderim sırasında SQL transaction'ını açık tutmak kötü bir takastır: broker gecikmesi PostgreSQL bağlantılarını ve row lock'larını tüketir. Bunun yerine kısa bir claim transaction'ında satıra lease atayın, transaction'ı kapatın, sonra gönderin. Süreç gönderimden sonra fakat published işaretinden önce ölürse aynı olay yeniden gider. Bu yüzden teslimat semantiği pratikte at-least-once'dur ve tüketici idempotent olmalıdır.
WITH claimed AS (
SELECT id
FROM outbox_event
WHERE published_at IS NULL
AND (lease_until IS NULL OR lease_until < clock_timestamp())
ORDER BY created_at
FOR UPDATE SKIP LOCKED
LIMIT :batchSize
)
UPDATE outbox_event e
SET lease_until = clock_timestamp() + interval '30 seconds',
attempt_count = e.attempt_count + 1
FROM claimed
WHERE e.id = claimed.id
RETURNING e.id, e.event_type, e.aggregate_id, e.payload, e.headers;Yayınlayıcı, dönen her kayıt için Kafka key olarak aggregate_id kullanmalıdır. Aynı key aynı partition'a gittiği için partition içi sıra korunur. Başarıdan sonra koşullu güncelleme kullanın: UPDATE outbox_event SET published_at = clock_timestamp(), lease_until = NULL WHERE id = :id AND published_at IS NULL. Hata durumunda lease'i hemen temizlemek yerine üstel gecikme için next_attempt_at kolonu eklemek, bozuk bir kaydın her 100 ms taramada tekrar alınmasını önler.
Tüketicide yalnızca Kafka offset commit'ine güvenmek yeterli değildir; tüketici DB yazdıktan sonra offset commit etmeden çökerse mesaj tekrar gelir. Hedef veritabanında processed_event(event_id UUID PRIMARY KEY, processed_at timestamptz) tablosuna önce insert deneyin. PostgreSQL'de INSERT ... ON CONFLICT DO NOTHING sonucu 0 satırsa handler'ı atlayın. Event ID, payload'dan değil Kafka header'daki değişmez event-id değerinden alınmalıdır; payload şeması v2'ye evrildiğinde deduplikasyon anahtarının değişmesi bu sayede engellenir.
CDC, Debezium ve java microservices için teslimat seçimi
Polling publisher uygulama kodunu basit tutar, fakat her pod'un periyodik sorgu çalıştırması gerekir. PostgreSQL logical decoding tabanlı Debezium Outbox Event Router ise WAL'deki outbox insert'lerini Kafka'ya aktarır. Bu yaklaşımda uygulama publish izni taşımaz; connector offset'i Kafka Connect tarafından saklanır. Debezium kullanırken outbox tablosuna doğrudan DELETE uygulamak yerine, connector'ın ilgili LSN'i okuduğundan emin olmadan silmeyin. Güvenli saklama penceresi için örneğin yayımlanmış ve 7 günden eski satırları küçük batch'lerle temizleyin: DELETE FROM outbox_event WHERE id IN (SELECT id FROM outbox_event WHERE published_at < now() - interval '7 days' LIMIT 10000).
Kafka producer tarafında broker dayanıklılığı için şu ayarları doğrulayın: acks=all, enable.idempotence=true, max.in.flight.requests.per.connection<=5. İdempotent producer, ağ zaman aşımında aynı producer session'ındaki tekrar gönderimin broker'da çift kayda dönüşmesini azaltır; fakat uygulamanın crash sonrası yeniden yayınladığı outbox olayını ortadan kaldırmaz. Bu nedenle producer idempotence ile consumer deduplikasyonu birbirinin alternatifi değildir.
Java microservices ekiplerinde event şeması için Avro, Protobuf veya JSON Schema registry seçin ve event type'a sürüm ekleyin: order.placed.v1. V1 tüketicileri kaldırılmadan alan silmeyin; yeni alanı opsiyonel ekleyin. Bu disiplin java backend geliştirme işinde REST DTO'larının event sözleşmesi olarak yeniden kullanılmamasını da sağlar: REST DTO'su istemci görünümüdür, domain event ise geçmişte yaşanmış değişmez bir olgudur.
Spring Boot eğitimi için outbox gecikmesini ölçmek ve kapasitelemek
Outbox gecikmesini 'sistem hızlı' diyerek değil, clock_timestamp() - created_at ile ölçün. Micrometer üzerinde outbox.publish.age histogramını ve outbox.pending.count gauge'ını yayınlayın. Prometheus'ta p95 yaşın 30 saniyeyi aşması, HTTP istek süresi normal görünse bile asenkron teslimat SLO'sunun ihlal edildiğini gösterir. PostgreSQL tarafında pg_stat_statements ile claim sorgusunun çağrı başına ortalama ve p95 yürütme süresini, EXPLAIN (ANALYZE, BUFFERS) ile de kısmi indeks kullandığını doğrulayın.
management:
observations:
key-values:
service: order-service
spring:
kafka:
producer:
properties:
acks: all
enable.idempotence: true
retries: 2147483647
# PromQL: son 15 dakikada p95 outbox yaşı
histogram_quantile(0.95,
sum by (le) (rate(outbox_publish_age_seconds_bucket[15m])))Önce tek kayıt claim eden bir publisher ile baz ölçüm alın: p95 publish age, claim SQL p95 ve Kafka send callback p95 değerlerini kaydedin. Ardından batch'i 1'den 100'e çıkarıp tekrar ölçün. Eğer SQL çağrı maliyeti düşerken publish age yükselirse, tek thread'in batch içindeki yavaş Kafka gönderiminde takıldığını görürsünüz; batch'i büyütmek yerine sınırlı paralel gönderim ve ayrı claim worker sayısı gerekir. JDK Flight Recorder'da jdk.SocketRead, jdk.JavaMonitorEnter ve thread park süreleriyle publisher'ın broker I/O mu yoksa connection pool beklemesi mi yaşadığını ayırın. Ölçmeden batch büyütmek, outbox tablosunda lease süresinin dolmasına ve çift gönderim oranının artmasına yol açabilir.
Spring AI, Spring MCP ve model context protocol olayları
Spring AI veya Spring MCP kullanan bir serviste model çağrısının sonucu da denetlenebilir bir domain olayı olabilir: örneğin belge sınıflandırma tamamlandıktan sonra document.classified.v1 outbox'a yazılır. Fakat prompt'un tamamını veya modelin ham cevabını varsayılan olarak event payload'a koymayın. Bu alanlar PII, erişim belirteci veya müşteri verisi içerebilir. Payload'a belge ID, model sağlayıcısı, prompt template sürümü, normalize edilmiş sonuç ve isteğe bağlı hash koyun; hassas ham içeriği şifreli, erişim kontrollü ayrı depoda tutun.
Model Context Protocol ile bir aracın yaptığı yan etkili işlemde outbox, tool çağrısı ile business state arasındaki denetim zincirini güçlendirir. Örneğin bir modelin 'refund' aracı çağrısı önce uygulamanın yetkilendirilmiş command handler'ına gelir, refund ve refund.requested.v1 aynı transaction'da yazılır. Tool sonucu doğrudan Kafka olayı olarak kabul edilirse, model tekrar denemeleri ve ağ hataları finansal işlemi çoğaltabilir. Spring MCP entegrasyonunda tool call ID'yi outbox event ID veya ayrı bir idempotency key ile ilişkilendirin.
Bu konu java eğitimi, spring boot eğitimi ve java fullstack eğitimi programlarında yalnızca mesajlaşma bölümü olarak ele alınmamalıdır. Frontend'in sipariş durumunu sorgulaması için ayrı bir read model veya polling endpoint tasarlayın; istemcinin Kafka teslimatını HTTP yanıtı içinde beklemesi, outbox'ın asenkron hata izolasyonunu bozar. Uygulanabilir kontrol listesi şudur: business write ve outbox insert aynı transaction'da mı, consumer event ID ile idempotent mi, lease süresi gözlemleniyor mu, şema geriye uyumlu mu, ham AI verisi event içinde taşınmıyor mu?
İlgili Eğitim
YTÜSEM İlgili Eğitim
Java Spring Boot ReactJS FullStack Eğitimi (Yıldız Teknik Üniversitesi SEM)
Sık Sorulan Sorular
Spring Data JPA ile transactional outbox yazarken @TransactionalEventListener yeterli mi?
Hayır. AFTER_COMMIT listener, commit ile Kafka send arasında JVM kapanırsa tekrar oynatılabilir bir kayıt bırakmaz. Listener kullanılacaksa yalnızca outbox publisher'ı uyandırmak için kullanılmalı; event payload'u business transaction içinde outbox_event tablosuna Spring Data JPA ile yazılmalıdır.
Java microservices uygulamasında outbox publisher Kafka'ya exactly-once teslimat sağlar mı?
Genel durumda hayır. Publisher Kafka'ya başarıyla yazdıktan sonra published_at güncellemesinden önce çökerse kayıt yeniden gönderilir. Kafka producer idempotence broker tekrarlarını azaltır, fakat crash sonrası uygulama tekrarını çözmez. Tüketicide event-id için UNIQUE indeks ve INSERT ON CONFLICT DO NOTHING ile idempotency uygulanmalıdır.
Hibernate ORM outbox tablosunda SKIP LOCKED kullanmak için JPA repository yeterli mi?
Basit CRUD için yeterlidir, ancak claim işlemi için native SQL veya JdbcTemplate tercih edin. FOR UPDATE SKIP LOCKED, UPDATE ... RETURNING ve clock_timestamp() tek SQL ifadesinde atomik claim sağlar. JPA entity'lerini kilitleyip Kafka gönderimi boyunca transaction açık tutmak connection pool ve row lock beklemesini büyütür.
Spring AI ve model context protocol olayları outbox payload'unda nasıl saklanmalı?
Prompt ve ham model yanıtı yerine referans ID, template sürümü, sağlayıcı adı, normalize çıktı ve bütünlük hash'i saklayın. Tool call ID'yi idempotency key olarak koruyun. Hassas içeriği event broker'a yaymak yerine şifreli bir veri deposunda tutup olaydan o kayda referans verin.
AI / LLM Discovery
Bu makale Opendart Akademi Java 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.


