Kapsamlı Apache Kafka Geliştirici Rehberi

8 Ağustos 2026 · netologist · 16 dakika, 3246 kelime ·

Kafka’nın temel mimarisini, producer/consumer iç işleyişini, exactly-once semantiğini, Kafka Streams’i, Kafka Connect’i, şema yönetimini, güvenliği, performans ayarlarını, yaygın tasarım desenlerini (pattern) ve kaçınılması gereken anti-pattern’leri derinlemesine ele alan bir referans rehber.


İçindekiler

  1. Temel Mimari
  2. Topic’ler, Partition’lar ve Replikasyon
  3. Producer’lar
  4. Consumer’lar ve Consumer Group’lar
  5. Teslimat Semantiği ve Exactly-Once
  6. Serileştirme ve Schema Registry
  7. Kafka Streams
  8. Kafka Connect
  9. Topic Tasarımı ve Partitioning Stratejisi
  10. Performans Ayarlama (Tuning)
  11. İzleme (Monitoring) ve Gözlemlenebilirlik
  12. Güvenlik
  13. Yaygın Tasarım Desenleri (Design Patterns)
  14. Anti-Pattern’ler ve Tuzaklar
  15. Operasyonel En İyi Uygulamalar

1. Temel Mimari

Kafka; dağıtık (distributed), partition’lara bölünmüş, replike edilmiş bir commit log servisidir ve büyük ölçekte publish-subscribe mesajlaşma sistemi gibi davranır.

Temel Bileşenler

Log Yapısı

Her partition, diskte segment dosyaları dizisi olarak saklanır. Kafka kayıtları asla yerinde (in-place) yeniden yazmaz — tamamen append-only’dir, bu da ona ağ hızına yakın sıralı disk I/O verimi kazandırır.

/var/lib/kafka/data/my-topic-0/
  ├── 00000000000000000000.log
  ├── 00000000000000000000.index
  ├── 00000000000000000000.timeindex
  ├── 00000000000000452312.log
  └── ...

KRaft’ta Broker Rolleri


2. Topic’ler, Partition’lar ve Replikasyon

Partition Sayısı

Replikasyon Faktörü

In-Sync Replicas (ISR)

min.insync.replicas=2
acks=all

→ Bir yazma işlemi, ancak en az min.insync.replicas kadar replikaya kopyalandıktan sonra onaylanır; bu da min.insync.replicas‘tan az sayıda broker aynı anda çökmediği sürece veri kaybı olmayacağını garanti eder.

Önemli Topic Ayarları

# Retention (Saklama)
retention.ms=604800000          # 7 gün (zaman bazlı)
retention.bytes=-1               # sınırsız (boyut bazlı)

# Compaction (Sıkıştırma/Birleştirme)
cleanup.policy=compact           # ya da "delete" ya da "compact,delete"
min.cleanable.dirty.ratio=0.5
segment.ms=604800000

# Dayanıklılık
min.insync.replicas=2
unclean.leader.election.enable=false   # senkron olmayan replica'ların leader olmasına ASLA izin vermeyin

# Verim (Throughput)
max.message.bytes=1048588

Log Compaction vs Delete


3. Producer’lar

Temel Konfigürasyon

bootstrap.servers=broker1:9092,broker2:9092
key.serializer=org.apache.kafka.common.serialization.StringSerializer
value.serializer=org.apache.kafka.common.serialization.StringSerializer

acks=all                       # all | 1 | 0
enable.idempotence=true        # Kafka 3.0'dan itibaren varsayılan true
retries=2147483647             # pratikte sonsuz, delivery.timeout.ms ile sınırlı
max.in.flight.requests.per.connection=5   # idempotence açıkken 5'e kadar güvenli
delivery.timeout.ms=120000
linger.ms=5                    # batch penceresi, throughput için gecikmeyi feda eder
batch.size=32768
compression.type=lz4           # gzip/snappy yerine lz4/zstd önerilir
buffer.memory=33554432

acks Semantiği

DeğerAnlamıDayanıklılıkGecikme
0Gönder-ve-unut (fire-and-forget)YokEn düşük
1Sadece leader onayıLeader, replikasyondan önce çökerse veri kaybıOrta
all (-1)Tüm in-sync replica’ların onayıEn güçlü (min.insync.replicas ile)En yüksek

Idempotent (Etkisiz Eleman) Producer

enable.idempotence=true ayarı her producer’a bir Producer ID (PID) ve her mesaja partition başına bir sequence number (sıra numarası) atar. Broker, tekrar denemelerdeki (retry) kopyaları eler ve uygulama kodunda değişiklik yapmadan partition başına, producer oturumu başına exactly-once teslimat sağlar. Bu, transactional producer’ların temelini oluşturur.

Partitioning Stratejisi

// Varsayılan partitioner (Kafka >= 2.4): key null olduğunda sticky partitioning
// eski round-robin'e göre batching'i iyileştirir.

ProducerRecord<String, String> record =
    new ProducerRecord<>("orders", customerId, payload); // key, hash üzerinden partition'ı belirler

Backpressure ve Hata Yönetimi

producer.send(record, (metadata, exception) -> {
    if (exception != null) {
        if (exception instanceof RetriableException) {
            // producer'ın dahili retry mekanizmasına bırakın; görünürlük için loglayın
        } else {
            // tekrar denenemez: DLQ'ya gönder, uyarı ver veya hızlıca başarısız ol
        }
    }
});

Transactional Producer (Exactly-Once)

transactional.id=order-service-1
enable.idempotence=true
producer.initTransactions();
try {
    producer.beginTransaction();
    producer.send(record1);
    producer.send(record2);
    producer.sendOffsetsToTransaction(offsets, groupMetadata); // consume-transform-produce için
    producer.commitTransaction();
} catch (ProducerFencedException | OutOfOrderSequenceException | AuthorizationException e) {
    producer.close(); // ölümcül (fatal), yeniden başlatma gerekir
} catch (KafkaException e) {
    producer.abortTransaction();
}

4. Consumer’lar ve Consumer Group’lar

Temel Konfigürasyon

bootstrap.servers=broker1:9092,broker2:9092
group.id=order-processing-service
key.deserializer=org.apache.kafka.common.serialization.StringDeserializer
value.deserializer=org.apache.kafka.common.serialization.StringDeserializer

enable.auto.commit=false        # "at-least-once" üzerinde kontrol için manuel commit tercih edin
auto.offset.reset=earliest      # earliest | latest | none
max.poll.records=500
max.poll.interval.ms=300000     # ölü sayılmadan önce iki poll arasında izin verilen süre
session.timeout.ms=45000
heartbeat.interval.ms=15000
fetch.min.bytes=1
fetch.max.wait.ms=500
isolation.level=read_committed  # transactional topic okunuyorsa
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor

Consumer Group Rebalancing

group.instance.id=consumer-instance-7
session.timeout.ms=45000

Manuel Offset Commit Desenleri

while (true) {
    ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
    for (ConsumerRecord<String, String> record : records) {
        process(record);
    }
    consumer.commitSync(); // ya da daha yüksek throughput için callback'li commitAsync()
}

İşlemden-sonra-commit yaklaşımı at-least-once (en az bir kez) teslimat sağlar (çoğu sistem için varsayılan ve önerilen). Exactly-once işleme etkisi elde etmek için şunlardan birini yapın:

  1. İşlemeyi idempotent hale getirin (örn. key’e göre upsert), veya
  2. Consume-transform-produce hattında Kafka transaction’larını kullanın (bkz. §3), veya
  3. Offset’i ve işleme sonucunu, harici bir sistemde atomik olarak saklayın (örn. hem iş satırını hem offset’i yazan bir DB transaction’ı).

Zehirli Mesajları (Poison Pill) Yönetmek

try {
    process(record);
} catch (DeserializationException | UnrecoverableBusinessException e) {
    deadLetterProducer.send(new ProducerRecord<>("orders-dlq", record.key(), record.value()));
} 

Tek bir hatalı biçimlendirilmiş kaydın tüm partition’ı sonsuza kadar bloklamasına asla izin vermeyin — orijinal header/metadata’yı koruyarak daha sonra tekrar oynatılabilecek (replay) bir Dead Letter Queue (DLQ)‘ya yönlendirin.

Consumer Lag (Gecikme)

Lag = (partition’daki en son offset) − (group tarafından son commit edilen offset). kafka-consumer-groups.sh --describe --group X komutuyla ya da JMX/Burrow/Prometheus exporter’ları ile izleyin. Sürekli artan lag, consumer’ın yetişemediğini gösterir — ölçeği genişletin (partition sayısına kadar consumer ekleyin), işlemeyi optimize edin veya partition sayısını artırın.


5. Teslimat Semantiği ve Exactly-Once

SemantikNasıl elde edilirNotlar
At-most-once (en fazla bir kez)enable.auto.commit=true ile işlemeden önce commit, ya da acks=0Veri kaybı mümkün
At-least-once (en az bir kez)Başarılı işlemeden sonra commit; acks=allVarsayılan öneri; hata durumunda kopyalar (duplicate) mümkün
Exactly-once (tam olarak bir kez)Idempotent producer + transaction’lar + read_committed consumer’lar, YA DA idempotent downstream sinkEn yüksek karmaşıklık, sadece kopyaların gerçekten kabul edilemez olduğu durumlarda kullanın (örn. finansal işlemler)

Kafka’daki Exactly-Once Semantics (EOS), Kafka ekosistemi içinde garanti edilir (producer → topic → consumer, Kafka Streams dahil). Pipeline’ınız harici bir sisteme (DB, S3, REST API) yazıyorsa şunlardan birine ihtiyacınız var:


6. Serileştirme ve Schema Registry

Schema Registry Neden Gerekli

Avro/Protobuf/JSON Schema + Confluent (veya Apicurio/Karapace) Schema Registry size şunları sağlar:

Uyumluluk Modları

ModConsumer’lar okuyabilirProducer’lar yazabilir
BACKWARDYeni şema eski veriyi okurConsumer önce yükseltilmeli
FORWARDEski şema yeni veriyi okurProducer önce yükseltilebilir
FULLHer iki yön deEn güvenli, en kısıtlayıcı
NONEKontrol yokTehlikeli, production’da kaçının
key.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
value.serializer=io.confluent.kafka.serializers.KafkaAvroSerializer
schema.registry.url=http://schema-registry:8081
// Örnek Avro şeması — yeni alanlar için her zaman varsayılan (default) değer verin (geriye dönük uyumluluk için)
{
  "type": "record",
  "name": "OrderCreated",
  "fields": [
    {"name": "orderId", "type": "string"},
    {"name": "amount", "type": "double"},
    {"name": "currency", "type": "string", "default": "USD"}
  ]
}

Genel Kurallar


7. Kafka Streams

Kafka Streams, doğrudan Kafka üzerine stream-processing (akış işleme) uygulamaları geliştirmek için bir istemci kütüphanesidir — ayrı bir cluster gerektirmez.

Temel Soyutlamalar

Örnek Topoloji

StreamsBuilder builder = new StreamsBuilder();

KStream<String, Order> orders = builder.stream("orders",
        Consumed.with(Serdes.String(), orderSerde));

KTable<String, Long> orderCountsByCustomer = orders
        .groupBy((key, order) -> order.getCustomerId(), Grouped.with(Serdes.String(), orderSerde))
        .count(Materialized.as("order-counts-store"));

orderCountsByCustomer.toStream().to("order-counts", Produced.with(Serdes.String(), Serdes.Long()));

KafkaStreams streams = new KafkaStreams(builder.build(), props);
streams.start();

Stream-Table Dualitesi ve Join’ler

// Stream-Stream join (pencereli/windowed, her iki taraf da zamanla sınırlı)
orders.join(payments,
    (order, payment) -> enrich(order, payment),
    JoinWindows.ofTimeDifferenceWithNoGrace(Duration.ofMinutes(5)));

// Stream-Table join (pencere yok; tablo mevcut durumu temsil eder)
orders.join(customersTable, (order, customer) -> enrich(order, customer));

State Store’lar ve Hata Toleransı

num.standby.replicas=1
state.dir=/data/kafka-streams
cache.max.bytes.buffering=10485760

Streams’te Exactly-Once

processing.guarantee=exactly_once_v2

Bu ayar, read-process-write döngüsünü otomatik olarak bir Kafka transaction’ı içine sarar — transactional kodu elle yazmadan EOS elde etmenin önerilen yoludur.

Windowing (Pencereleme)

TimeWindows.ofSizeAndGrace(Duration.ofMinutes(5), Duration.ofMinutes(1)); // grace period'u her zaman ayarlayın

8. Kafka Connect

Kafka ile harici sistemler arasında ölçeklenebilir, hataya dayanıklı entegrasyon için bir framework — yaygın entegrasyonlar için özel producer/consumer kodu yazmaya gerek kalmaz.

Modlar

Source vs Sink Connector’lar

// Debezium (CDC) source connector örneği
{
  "name": "orders-cdc",
  "config": {
    "connector.class": "io.debezium.connector.postgresql.PostgresConnector",
    "database.hostname": "postgres",
    "database.dbname": "orders_db",
    "table.include.list": "public.orders",
    "topic.prefix": "cdc",
    "plugin.name": "pgoutput"
  }
}
// JDBC sink connector örneği
{
  "name": "orders-sink",
  "config": {
    "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
    "connection.url": "jdbc:postgresql://warehouse/db",
    "topics": "orders",
    "insert.mode": "upsert",
    "pk.mode": "record_key",
    "auto.create": "true"
  }
}

Single Message Transforms (SMT’ler)

Ayrı bir stream processor gerektirmeden hafif, satır-içi dönüşümler:

"transforms": "maskPII,route",
"transforms.maskPII.type": "org.apache.kafka.connect.transforms.MaskField$Value",
"transforms.maskPII.fields": "ssn",
"transforms.route.type": "org.apache.kafka.connect.transforms.RegexRouter",
"transforms.route.regex": "(.*)",
"transforms.route.replacement": "prefixed-$1"

Hata Yönetimi (Connect için Dead Letter Queue)

errors.tolerance=all
errors.deadletterqueue.topic.name=connect-dlq
errors.deadletterqueue.context.headers.enable=true
errors.log.enable=true

9. Topic Tasarımı ve Partitioning Stratejisi

İsimlendirme Kuralları

<domain>.<entity>.<event-tipi>.v<versiyon>
örn. ecommerce.orders.created.v1
     billing.invoices.updated.v2

Key Seçimi

Multi-Tenancy (Çoklu Kiracı)


10. Performans Ayarlama (Tuning)

Producer Throughput

Consumer Throughput

Broker Ayarlama

num.network.threads=8
num.io.threads=16
num.replica.fetchers=4
socket.send.buffer.bytes=1048576
socket.receive.buffer.bytes=1048576
log.flush.interval.messages=Long.MAX_VALUE   # dayanıklılık için fsync yerine replikasyona güvenin

Boyutlandırma Kontrol Listesi


11. İzleme (Monitoring) ve Gözlemlenebilirlik

Kritik Broker Metrikleri (JMX)

MetrikNeden önemli
UnderReplicatedPartitions> 0 ise veri kaybı riski var; hemen inceleyin
OfflinePartitionsCount> 0 ise erişilemezlik var demektir
ActiveControllerCountCluster genelinde tam olarak 1 olmalı
RequestHandlerAvgIdlePercentDüşükse: broker CPU’ya bağımlı, thread pool doymuş
BytesInPerSec / BytesOutPerSecThroughput trendi
ISR shrink/expand oranıSık küçülme = kararsız replica’lar / ağ sorunları

Consumer Metrikleri

Araçlar


12. Güvenlik

Şifreleme

listeners=SSL://broker1:9093
ssl.keystore.location=/certs/kafka.server.keystore.jks
ssl.keystore.password=changeit
ssl.truststore.location=/certs/kafka.server.truststore.jks

Production’da tüm broker-arası ve client-broker trafiği için TLS kullanın (security.inter.broker.protocol=SSL).

Kimlik Doğrulama (Authentication)

Yetkilendirme (ACL’ler)

kafka-acls.sh --bootstrap-server broker:9092 \
  --add --allow-principal User:order-service \
  --operation Write --operation Read \
  --topic orders

Veri Düzeyi


13. Yaygın Tasarım Desenleri (Design Patterns)

Transactional Outbox Pattern

Dağıtık transaction kullanmadan “çift yazma” (dual write) problemini çözer (DB yazımı + Kafka publish atomik olmalıdır):

  1. Uygulama, iş satırını ve bir outbox satırını aynı yerel DB transaction’ında yazar.
  2. Bir CDC connector’ı (örn. Debezium) outbox tablosunu takip eder ve Kafka’ya yayınlar.
  3. Outbox satırı daha sonra temizlenir.

Bu, event’in yalnızca ve yalnızca DB transaction’ı commit edildiğinde yayınlanmasını garanti eder.

Event Sourcing

Durum değişikliklerini sıralı, değiştirilemez bir event dizisi olarak saklayın (Kafka topic’i = gerçeğin tek kaynağı/source of truth), event’leri materialize edilmiş bir görünüme (KTable / harici DB) tekrar oynatarak mevcut durumu yeniden inşa edin.

CQRS (Command Query Responsibility Segregation)

Yazma modelini (event yayınlayan komut işleyicileri) okuma modelinden (bu event’leri tüketen Kafka Streams/Connect ile inşa edilen materialize görünümler) ayırın — okuma/yazmayı bağımsız olarak ölçeklendirmeyi ve tek bir event akışından birden fazla okumaya-optimize edilmiş projeksiyon üretmeyi mümkün kılar.

Saga Pattern (Event’ler Üzerinden Koreografi)

Servisler arası uzun süren iş transaction’ları, tamamen yayınlanan/tüketilen event’ler üzerinden koordine edilir; her servis bir önceki adımın event’ine tepki verir ve hata durumunda telafi edici (compensating) bir event yayınlar — dağıtık 2PC transaction’lardan kaçınır.

Change Data Capture (CDC)

Veritabanı satır düzeyindeki değişiklikleri neredeyse gerçek zamanlı olarak Kafka topic’lerine akıtmak için Debezium/Connect kullanın; downstream sistemleri doğrudan DB erişiminden ayırır ve eski (legacy) veritabanları üzerinde event-driven mimariler kurmayı mümkün kılar.

KV Store Olarak Compact Edilmiş Topic

Referans/config verisi için dayanıklı, replike edilmiş, tekrar oynatılabilir bir key-value store olarak cleanup.policy=compact topic’lerini kullanın — consumer’lar, baştan tekrar oynatarak bellek-içi veya RocksDB üzerinde bir map inşa eder.

Dead Letter Queue (DLQ)

İşlenemeyen mesajları, header’larda hata metadata’sı bulunan ayrı bir topic’e yönlendirin; bu, ana pipeline’ın devam etmesini ve inceleme sonrası manuel/otomatik tekrar oynatmayı (replay) mümkün kılar.

Fan-Out / Birden Fazla Consumer Group

Tek bir topic, her biri kendi consumer group’una sahip birçok servis tarafından bağımsız olarak tüketilebilir — Kafka, geleneksel kuyruklardan farklı olarak mesajları tüm group’lar için saklar, bu da yeniden yayınlamaya (republish) gerek kalmadan gerçek pub-sub fan-out sağlar.


14. Anti-Pattern’ler ve Tuzaklar

Anti-PatternNeden sorunÇözüm
Kafka’yı mesaj-başına ack/nack ile geleneksel bir görev kuyruğu (task queue) gibi kullanmakKafka commit’leri mesaj-bazlı değil offset-bazlıdır; öncelik kuyrukları veya öğe-başına ack semantiğine uymazKlasik kuyruklama için RabbitMQ/SQS kullanın, ya da DLQ’larla dikkatlice modelleyin
Çok fazla küçük topic (ölçekte müşteri-başına topic)Broker metadata yükü, partition leader seçim yükü, ZK/KRaft metadata şişmesiBunun yerine ACL izolasyonlu, key’li multi-tenant topic’ler kullanın
Devasa mesajlar (> birkaç MB)Throughput’a, belleğe, GC baskısına zarar verirPayload’u blob depolamada (S3) saklayın, Kafka’da bir referans/pointer yayınlayın
Production’da unclean.leader.election.enable=trueLeader failover’da sessiz veri kaybıfalse bırakın, bunun yerine düzgün replikasyona yatırım yapın
Partition’lar arası mesaj sıralamasına güvenmekKafka yalnızca bir partition içinde sıralamayı garanti ederİlgili event’lerin aynı partition’a düşmesi için kayıtları key’leyin
Yavaş/başarısız işleme ile auto-commitİşlenemeyen kayıtlar için offset commit edilebilir, sessiz veri kaybına yol açarBaşarılı işlemeden sonra manuel commit
Yavaş consumer’lar için max.poll.interval.ms‘i doğru ayarlamamakConsumer işleme ortasında group’tan atılır, rebalance fırtınalarına yol açarInterval’i gerçekçi en kötü durum işleme süresine göre ayarlayın, ya da ağır işi asenkron olarak dışarı taşıyın
Registry/uyumluluk kontrolü olmadan şema değişiklikleriProduction’da consumer’ları sessizce bozarCI/CD’de Schema Registry uyumluluk kontrollerini zorunlu kılın
Consumer lag uyarılarını görmezden gelmekKesinti (outage) oluşana kadar sessiz birikim büyümesiYalnızca mutlak değeri değil, lag trendini de izleyin
Kafka’yı request-response RPC için kullanmakgRPC/REST’e göre gereksiz gecikme/karmaşıklık eklerKafka’yı senkron çağrılar için değil, asenkron/event-driven akışlar için kullanın
Partition sayısını önceden planlamamakTopic’i yeniden oluşturmadan partition sayısını sonradan azaltamazsınızTopic oluşturmadan önce beklenen throughput/paralelliği modelleyin

15. Operasyonel En İyi Uygulamalar


Hızlı Referans: Önerilen Production Varsayılanları

# Topic
replication.factor=3
min.insync.replicas=2
unclean.leader.election.enable=false

# Producer
acks=all
enable.idempotence=true
compression.type=lz4
linger.ms=5

# Consumer
enable.auto.commit=false
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor
isolation.level=read_committed

Bu rehber, 2026 başı itibarıyla modern Kafka (3.x/4.x, KRaft tabanlı) en iyi uygulamalarını yansıtır. Varsayılanlar ve mevcut config’ler sürümler arasında değişebildiğinden, çalıştırdığınız spesifik sürüm için her zaman resmi Apache Kafka dokümantasyonuyla çapraz kontrol yapın.