Kapsamlı Apache Kafka Geliştirici Rehberi
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
- Temel Mimari
- Topic’ler, Partition’lar ve Replikasyon
- Producer’lar
- Consumer’lar ve Consumer Group’lar
- Teslimat Semantiği ve Exactly-Once
- Serileştirme ve Schema Registry
- Kafka Streams
- Kafka Connect
- Topic Tasarımı ve Partitioning Stratejisi
- Performans Ayarlama (Tuning)
- İzleme (Monitoring) ve Gözlemlenebilirlik
- Güvenlik
- Yaygın Tasarım Desenleri (Design Patterns)
- Anti-Pattern’ler ve Tuzaklar
- 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
- Broker: Veriyi depolayan ve istemcilere hizmet veren tek bir Kafka sunucusu. Bir cluster (küme), birden fazla broker’dan oluşur.
- Topic: Kayıtların (record) yayınlandığı, isimlendirilmiş, sadece-ekleme (append-only) log. Topic’ler partition’lara bölünür.
- Partition: Paralellik ve sıralamanın (ordering) temel birimi. Her partition, offset ile tanımlanan, sıralı ve değiştirilemez (immutable) bir kayıt dizisidir.
- Producer: Topic’lere kayıt yayınlar (publish).
- Consumer: Topic’lere abone olur (subscribe) ve kayıt akışını işler.
- Consumer Group: Bir topic’i birlikte tüketmek için işbirliği yapan consumer kümesi; her partition, bir group içinde yalnızca bir consumer’a atanır.
- ZooKeeper (eski) / KRaft (modern): Cluster meta verisi ve controller seçimi. Kafka 3.x’ten itibaren KRaft modu ZooKeeper bağımlılığını tamamen ortadan kaldırır — Kafka 4.0 itibarıyla ZooKeeper tamamen kaldırılmıştır. Yeni kurulumlarda daima KRaft kullanılmalıdır.
- Controller: Partition leader seçiminden ve metadata yayılımından sorumlu broker.
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
└── ...
.log— asıl kayıtlar.index— offset → dosyadaki fiziksel konum eşlemesi (seyrek/sparse index, ikili arama).timeindex— zaman damgası → offset eşlemesi (zaman bazlı arama için)
KRaft’ta Broker Rolleri
- Controller node’ları: metadata’yı yönetir (
__cluster_metadatatopic’i). - Broker node’ları: produce/fetch isteklerine hizmet eder.
- Küçük cluster’larda bir node her iki rolü de üstlenebilir (
process.roles=broker,controller); büyük cluster’larda ayrılır.
2. Topic’ler, Partition’lar ve Replikasyon
Partition Sayısı
- Maksimum paralelliği belirler (bir group’ta partition başına en fazla bir consumer çalışabilir).
- Daha fazla partition = daha fazla açık dosya tanıtıcısı, daha fazla replikasyon trafiği, daha uzun leader seçim süresi, çok küçük cluster’larda daha yüksek uçtan uca gecikme.
- Genel kural:
partition_sayısı = hedef_throughput / tek_partition_throughput. Muhafazakâr başlayın (6–12); mevcut bir topic’te partition sayısını yalnızca artırabilirsiniz, topic’i yeniden oluşturmadan azaltamazsınız (ayrıca artırmak, eski key’ler için sıralama garantisini bozar).
Replikasyon Faktörü
replication.factor=3, production için standart varsayılandır (min.insync.replicas=2ile 1 broker kaybına dayanır).- Her partition’ın bir leader‘ı ve N-1 follower‘ı vardır. Yalnızca leader okuma/yazma isteklerine hizmet eder (rack-aware follower fetching / KIP-392 ile read replica kullanılmadığı sürece).
In-Sync Replicas (ISR)
- ISR,
replica.lag.time.max.mssüresi içinde leader ile tamamen senkron olan replica kümesidir. min.insync.replicas(topic/broker ayarı), producer’dakiacks=allile birleştiğinde dayanıklılık (durability) garantinizi belirler:
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
- Delete: kayıtlar
retention.ms/retention.bytessonunda silinir. Geçici olaylar (transient event) için kullanılır. - Compact: Kafka, her key için en az son bilinen değeri kalıcı olarak saklar. Changelog topic’leri, state (durum) yeniden inşası ve Kafka Streams state store’ları / KTable’ların arka planı için kullanılır.
nulldeğer (“tombstone”),delete.retention.mssonunda o key’i siler.
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ğer | Anlamı | Dayanıklılık | Gecikme |
|---|---|---|---|
0 | Gönder-ve-unut (fire-and-forget) | Yok | En düşük |
1 | Sadece 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
- Key’li kayıtlar:
partition = hash(key) % partitionSayısı(murmur2 hash) — key başına sıralamayı garantiler. - Null key: sticky partitioner, batch dolana kadar kayıtları tek bir partition’a toplar, sonra geçiş yapar — throughput’u maksimize eder.
- Özel (custom) partitioner: iş mantığına özgü yönlendirme için
Partitionerarayüzünü uygulayın (örn. VIP müşterileri ayrı partition’lara yönlendirmek).
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
}
}
});
- Her zaman asenkron callback kullanın; gerçekten senkron onay gerekmedikçe kritik yol (hot path) içinde
.get()ile bloklamayın. RecordTooLargeException,TimeoutException,NotEnoughReplicasExceptiongibi hataları izleyin.
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();
}
transactional.id, her mantıksal producer örneği için (örn. partition/shard başına) sabit ve benzersiz olmalıdır ki fencing (kilitleme), yeniden başlatmalar arasında da düzgün çalışsın.- Downstream (aşağı akış) consumer’lar, commit edilmemiş/iptal edilmiş mesajları atlamak için
isolation.level=read_committedayarlamalıdır.
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
- Eager rebalancing (eski varsayılan,
RangeAssignor/RoundRobinAssignor): tüm consumer’lar işlemeyi durdurur, tüm partition’ları geri verir (revoke), ardından yeniden atanır (“stop-the-world”). - Cooperative Sticky rebalancing (
CooperativeStickyAssignor, önerilir): yalnızca gerçekten taşınması gereken partition’ları yeniden atar, duraklama süresini minimize eder — modern kurulumlarda bunu kullanın. - Static membership (
group.instance.id): geçici yeniden başlatmalarda (örn. rolling deploy, kısa ağ kesintisi) rebalance tetiklenmesini önler — büyük local state’e sahip stateful consumer’lar (Kafka Streams) için kritiktir.
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:
- İşlemeyi idempotent hale getirin (örn. key’e göre upsert), veya
- Consume-transform-produce hattında Kafka transaction’larını kullanın (bkz. §3), veya
- 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
| Semantik | Nasıl elde edilir | Notlar |
|---|---|---|
| At-most-once (en fazla bir kez) | enable.auto.commit=true ile işlemeden önce commit, ya da acks=0 | Veri kaybı mümkün |
| At-least-once (en az bir kez) | Başarılı işlemeden sonra commit; acks=all | Varsayı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 sink | En 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:
- Idempotent bir yazma işlemi (doğal/iş anahtarına göre upsert), veya
- Transactional outbox pattern (§13), veya
- Exactly-once destekli Kafka Connect sink connector’ları (örn. upsert kullanan JDBC sink, ya da
EXACTLY_ONCEteslimat garantisini uygulayan connector’lar, KIP-618).
6. Serileştirme ve Schema Registry
Schema Registry Neden Gerekli
Avro/Protobuf/JSON Schema + Confluent (veya Apicurio/Karapace) Schema Registry size şunları sağlar:
- Merkezi şema versiyonlama
- Hatalı veri cluster’a ulaşmadan önce uyumluluk (backward/forward/full) zorlaması
- Kompakt ikili kodlama (Avro/Protobuf) — JSON’a göre daha küçük mesajlar, daha hızlı (de)serileştirme
Uyumluluk Modları
| Mod | Consumer’lar okuyabilir | Producer’lar yazabilir |
|---|---|---|
BACKWARD | Yeni şema eski veriyi okur | Consumer önce yükseltilmeli |
FORWARD | Eski şema yeni veriyi okur | Producer önce yükseltilebilir |
FULL | Her iki yön de | En güvenli, en kısıtlayıcı |
NONE | Kontrol yok | Tehlikeli, 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
- Uyumlu bir geçiş (migration) planı olmadan zorunlu bir alanı asla kaldırmayın.
- Yeni alanları her zaman varsayılan değerle ekleyin.
- Subject organizasyonunuza uygun bir isimlendirme stratejisi (
TopicNameStrategy,RecordNameStrategy) kullanın —RecordNameStrategy, tek bir topic’te birden fazla event tipine izin verir.
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
- KStream: sınırsız kayıt akışı, her kayıt bağımsız bir olaydır (event).
- KTable: her key için en son değeri temsil eden changelog akışı (compact edilmiş, materialize edilmiş bir görünüm).
- GlobalKTable: her instance’ta tamamen replike edilmiş tablo — küçük referans/lookup verisi için idealdir.
Ö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ı
- State store’lar, yerel olarak RocksDB ile ve uzakta (compact edilmiş) changelog topic’leri ile desteklenir — bu sayede failover durumunda state her zaman yeniden inşa edilebilir.
- Hızlı failover için diğer instance’larda sıcak (hot) kopyalar tutmak üzere
standby.replicas > 0kullanın.
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)
- Tumbling: sabit, çakışmayan pencereler.
- Hopping: sabit boyutlu, çakışan (advance < size) pencereler.
- Sliding: her kayıt çifti için, zaman farkına göre pencere.
- Session: key başına dinamik, boşluk (gap) bazlı pencereler.
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
- Standalone: tek process, config bir dosyada — geliştirme/küçük kullanım senaryoları için uygundur.
- Distributed: birden fazla worker bir cluster oluşturur, config’ler dahili Kafka topic’lerinde saklanır (
connect-configs,connect-offsets,connect-status) — production standardıdır, ölçeklendirmeyi ve hata toleransını destekler.
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
- Uyumluluğu bozan şema evrimleri için bir versiyon eki ekleyin.
- Kaygıları ayırın: bilinçli olmadıkça (örn. CDC tabloları) ilgisiz event tiplerini tek bir topic’te karıştırmayın.
Key Seçimi
- Sıralama gerektiren varlığa göre key seçin (örn.
orderId,customerId). - Aşırı “sıcak” (hot) key’lerden kaçının (tek bir devasa müşteri bir partition’ı çarpıtabilir) — tüm varlık genelinde sıralama gerekmiyorsa, aşırı sıcak key’leri parçalamayı (salting, örn.
customerId#shard) düşünün.
Multi-Tenancy (Çoklu Kiracı)
- Key/header’da tenant ID’si bulunan paylaşılan topic’ler + okuma izolasyonu için ACL, YA DA
- Sıkı izolasyon için tenant başına ayrı topic’ler (daha yüksek operasyonel yük, ama güçlü kota ve güvenlik sınırları).
10. Performans Ayarlama (Tuning)
Producer Throughput
- İstek başına daha fazla batch işlemek için
linger.msvebatch.size‘ı artırın. compression.type=lz4veyazstdkullanın (en iyi sıkıştırma oranı, düşük CPU maliyeti).BufferExhaustedExceptiongörüyorsanızbuffer.memory‘yi artırın.
Consumer Throughput
- İstek yükünü azaltmak (fetch’leri batch’lemek) için
fetch.min.bytesvefetch.max.wait.ms‘i artırın. max.poll.records‘u dikkatli bir şekilde artırın —max.poll.interval.msile dengede tutun.- Consumer’ları partition sayısına kadar (ama fazlasına değil) ölçeklendirin — partition sayısının ötesindeki fazla consumer’lar boşta bekler.
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
- Page cache‘i etkin kullanın: heap’i çok büyük ayarlamayın (genellikle 4–6GB yeterlidir); log segmentlerini işletim sistemi cache’lesin.
- Adanmış disklar kullanın (yazma-ağırlıklı loglar için RAID 5/6’dan kaçının — RAID 10 veya replikasyonlu JBOD tercih edilir).
- Rack awareness (
broker.rack), replica’ları hata alanlarına (AZ’ler) yayar.
Boyutlandırma Kontrol Listesi
- Disk throughput’u: sıralı yazma olduğundan, düz SSD/NVMe genelde gerekli bile değildir — ama catch-up/yeniden işleme sırasındaki rastgele okumalara yardımcı olur.
- Ağ: replikasyon trafiği =
replication.factor × produce throughput, buna göre bütçeleyin.
11. İzleme (Monitoring) ve Gözlemlenebilirlik
Kritik Broker Metrikleri (JMX)
| Metrik | Neden önemli |
|---|---|
UnderReplicatedPartitions | > 0 ise veri kaybı riski var; hemen inceleyin |
OfflinePartitionsCount | > 0 ise erişilemezlik var demektir |
ActiveControllerCount | Cluster genelinde tam olarak 1 olmalı |
RequestHandlerAvgIdlePercent | Düşükse: broker CPU’ya bağımlı, thread pool doymuş |
BytesInPerSec / BytesOutPerSec | Throughput trendi |
ISR shrink/expand oranı | Sık küçülme = kararsız replica’lar / ağ sorunları |
Consumer Metrikleri
- Partition başına
records-lag-max/records-lag consumer group state(Stable, Rebalancing, Dead)commit-latency-avg
Araçlar
- Prometheus + JMX Exporter + Grafana — fiili açık kaynak standardı.
- Burrow — SLA tarzı değerlendirmeyle özel consumer-lag izleme.
- Cruise Control — otomatik partition rebalancing ve broker devreden çıkarma.
- kcat (kafkacat) — hızlı produce/consume/inceleme için CLI çok amaçlı araç.
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)
- SASL/SCRAM — kullanıcı adı/şifre bazlı, işletmesi basit.
- SASL/GSSAPI (Kerberos) — kurumsal AD/LDAP entegrasyonu.
- mTLS — kimlik olarak karşılıklı TLS sertifikaları.
- OAUTHBEARER — modern token bazlı kimlik doğrulama (OIDC entegrasyonu).
Yetkilendirme (ACL’ler)
kafka-acls.sh --bootstrap-server broker:9092 \
--add --allow-principal User:order-service \
--operation Write --operation Read \
--topic orders
- En az yetki (least privilege) prensibini uygulayın: producer’lar yalnızca kendi topic’lerinde
Writealsın, consumer’lar kendi consumer group’larıyla sınırlıRead+DescribeGroupalsın. - Topic başına ACL karmaşasından kaçınmak için takım/domain bazlı topic önekleri (prefix) için
--resource-pattern-type prefixedkullanın.
Veri Düzeyi
- Uyumluluk gerektiriyorsa, PII verisi Kafka’ya ulaşmadan önce alan düzeyinde (field-level) şifreleme/tokenizasyon uygulayın (Kafka düzeyindeki şifreleme transport/at-rest içindir, alan düzeyinde değildir).
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):
- Uygulama, iş satırını ve bir outbox satırını aynı yerel DB transaction’ında yazar.
- Bir CDC connector’ı (örn. Debezium) outbox tablosunu takip eder ve Kafka’ya yayınlar.
- 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-Pattern | Neden sorun | Çözüm |
|---|---|---|
| Kafka’yı mesaj-başına ack/nack ile geleneksel bir görev kuyruğu (task queue) gibi kullanmak | Kafka commit’leri mesaj-bazlı değil offset-bazlıdır; öncelik kuyrukları veya öğe-başına ack semantiğine uymaz | Klasik 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şmesi | Bunun yerine ACL izolasyonlu, key’li multi-tenant topic’ler kullanın |
| Devasa mesajlar (> birkaç MB) | Throughput’a, belleğe, GC baskısına zarar verir | Payload’u blob depolamada (S3) saklayın, Kafka’da bir referans/pointer yayınlayın |
Production’da unclean.leader.election.enable=true | Leader 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üvenmek | Kafka 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çar | Başarılı işlemeden sonra manuel commit |
Yavaş consumer’lar için max.poll.interval.ms‘i doğru ayarlamamak | Consumer işleme ortasında group’tan atılır, rebalance fırtınalarına yol açar | Interval’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şiklikleri | Production’da consumer’ları sessizce bozar | CI/CD’de Schema Registry uyumluluk kontrollerini zorunlu kılın |
| Consumer lag uyarılarını görmezden gelmek | Kesinti (outage) oluşana kadar sessiz birikim büyümesi | Yalnızca mutlak değeri değil, lag trendini de izleyin |
| Kafka’yı request-response RPC için kullanmak | gRPC/REST’e göre gereksiz gecikme/karmaşıklık ekler | Kafka’yı senkron çağrılar için değil, asenkron/event-driven akışlar için kullanın |
| Partition sayısını önceden planlamamak | Topic’i yeniden oluşturmadan partition sayısını sonradan azaltamazsınız | Topic oluşturmadan önce beklenen throughput/paralelliği modelleyin |
15. Operasyonel En İyi Uygulamalar
- Sadece produce trafiğini değil, replikasyon trafiğini de kapasite planlayın —
replication.factor=3, yazma ağ/disk yükünüzü üçe katlar. - Broker başına health check’lerle rolling restart’ları otomatikleştirin (bir sonraki broker’a geçmeden önce
UnderReplicatedPartitions=0olmasını bekleyin). - Gürültülü komşu (noisy-neighbor) tenant’ların cluster’ı aç bırakmasını önlemek için kota (quota) kullanın (
producer_byte_rate,consumer_byte_rate,request_percentage). - İstemci kütüphanelerini versiyonla sabitleyin ve broker yükseltmelerini önce staging’de test edin — Kafka güçlü geriye dönük uyumluluk sağlar ama ince davranış değişiklikleri yine de olabilir.
- Topic config’lerini ve ACL’leri kod olarak yedekleyin (Kafka topic’leri/ACL’leri için Terraform provider’ları mevcuttur) — production topic’lerini asla salt ad-hoc CLI ile yönetmeyin.
- Felaket kurtarmayı (disaster recovery) test edin: broker kaybını, AZ kaybını simüle edin ve
min.insync.replicas+acks=all‘ın sizi gerçekten koruduğunu doğrulayın. - Event şemalarınızı ve topic sahipliğini (ownership) dokümante edin — sahiplik belirsiz olduğunda event-driven sistemler sessizce başarısız olur; topic’leri sözleşmeli (contract) genel API’ler gibi ele alın.
- Bölgeler-arası (cross-region) replikasyon, felaket kurtarma veya multi-cluster mimariler için MirrorMaker 2 / Cluster Linking‘i değerlendirin.
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.