Kafka yük testi: producer, consumer, throughput, uçtan uca gecikme ve consumer lag
Kafka'ya dayanan bir sistemde darboğaz çoğu zaman HTTP katmanında değil, mesaj yolundadır: producer'lar broker'ın onayını beklerken yavaşlar, consumer'lar gelen mesaja yetişemez ve consumer lag büyür. Kafka yük testi bu yolu gerçekçi bir hızla zorlayıp üç soruya cevap arar: broker'lar istenen throughput'u hangi gecikmeyle kabul ediyor, consumer'lar aynı hızda okuyabiliyor mu, ve yük arttıkça geride kalan mesaj sayısı nasıl değişiyor.
Neyi ölçüyoruz?
- Producer gecikmesi: bir mesajın gönderilmesinden broker'ın onayına (ack) kadar geçen süre. Bu süre
acksayarına bağlıdır:acks=allbütün in-sync replica'ları bekler ve en yavaşıdır,acks=1yalnız leader'ı bekler,acks=0hiç beklemez ve ölçecek bir gecikme bırakmaz. Testte canlıdaki ayarı kullanın. - Uçtan uca gecikme: bir kaydın produce edilmesinden consume edilmesine kadar geçen süre; kullanıcının asıl hissettiği gecikme budur. Producer'ın ack süresi kısa kalırken bu süre consumer'lar geride kaldıkça büyür. Ölçerken dikkat: iki ayrı makinenin saati, ölçülen gecikmeden daha fazla farklı olabilir; gönderim zamanını bir makinede, alım zamanını başka bir makinede okuyan bir ölçüm milisaniyeleri değil saat farkını gösterir.
- Throughput: saniyede üretilen ve tüketilen mesaj sayısı, mesaj boyutuyla birlikte bayt/sn. 1 KB'lık mesajda saniyede 10.000 mesaj ile 100 KB'lık mesajda aynı sayı çok farklı yüklerdir.
- Consumer lag: bir consumer group'un bir partition'da okuduğu son offset ile partition'ın sonu arasındaki fark, yani henüz işlenmemiş mesaj sayısı. Sabit kalan küçük bir lag normaldir; yük altında sürekli büyüyen lag consumer'ların yetişemediğini gösterir.
- Hatalar: broker'ın reddettiği mesajlar, zaman aşımları, yetki hataları ve rebalance sırasında yaşanan kesintiler.
Yük modelini kurmak
Producer tarafında hedef genellikle bir hızdır ("saniyede 5.000 sipariş olayı"), bir kullanıcı sayısı değil. Bu yüzden sanal kullanıcı (VU) sayısını sabitlemek yerine sabit bir varış hızı (arrival rate) kullanın: broker yavaşlasa da mesaj hızı düşmez, gecikme olduğu gibi görünür. Kapalı bir modelde (sabit VU) broker yavaşladıkça producer'lar da yavaşlar ve sorun throughput'un sessizce düşmesi olarak gizlenir (yük testi türleri).
Consumer tarafında sayı önemlidir: bir consumer group'ta aynı anda en fazla partition sayısı kadar consumer çalışır. 12 partition'lı bir topic'e 20 consumer koyarsanız 8'i boşta bekler. Consumer sayısını canlıdaki instance sayısına eşitleyin; ölçeklemeyi test etmek istiyorsanız partition sayısını da hesaba katın.
Sık yapılan hatalar
- Gerçekçi olmayan mesajlar. Hep aynı küçük mesajı göndermek sıkıştırmayı ve önbellekleri olduğundan iyi gösterir. Boyutu ve içeriği canlıya benzeyen, her seferinde değişen mesajlar üretin.
- Key dağılımı. Mesajlar key'e göre partition'a gider. Az sayıda key kullanmak yükü birkaç partition'a yığar (hot partition); boş key ise partitioner'ın mesajları yaymasına izin verir. Canlıdaki key çeşitliliğine yakın durun.
- Paylaşılan topic. Canlı consumer'ların okuduğu bir topic'e test mesajı yazmayın. Test için ayrı bir topic ve ayrı bir consumer group kullanın, retention'ı ve temizliği önceden düşünün.
- Ortalama gecikme. Broker gecikmesinin kuyruğu uzundur: log segment geçişleri, replica senkronizasyonu ve GC duraklamaları seyrek ama büyük sıçramalar yapar. p95 ve p99'a bakın (p95 ve p99 rehberi).
- Kısa test. Lag'in büyüyüp büyümediğini görmek için yükün birkaç dakika sabit kalması gerekir. Birkaç dakikalık ramp-up ve en az 10–15 dakikalık sabit bir bölüm bırakın.
Spitfire ile
Önce Bağlantılar'da bir Kafka bağlantısı ekleyin: broker listesi, client id, gerekiyorsa SASL (PLAIN, SCRAM-SHA-256, SCRAM-SHA-512) ve TLS, ve acks (all varsayılan; leader ya da none). Parolalar şifreli saklanır, testte yalnız bağlantının adı geçer; Test düğmesi broker'lara ulaşıldığını denetler. Testte Kafka adımının iki işlemi vardır:
- produce: topic, key, value ve header'larla bir mesaj gönderir ve broker'ın onayını bekler. Adımın süresi broker'ın onayına kadar geçen süredir; Spitfire mesajı biriktirmek için beklemez (linger yok), eşzamanlı VU'ların mesajları yine de aynı istekte gruplanır. Bir runner'daki VU'lar bağlantı başına tek bir producer client paylaşır, gerçek bir servis gibi. Değerlerde
{{$uuid}},{{$randInt 10 5000}},{{$timestamp}}ya da CSV veri dosyasından değişkenler kullanılabilir; boş key mesajları partition'lara yayar. - consume: her VU kendi consumer client'ını tutar ve her adımda bir mesaj okur. Bir
groupIdverilirse VU'lar partition'ları gerçek consumer instance'ları gibi paylaşır; yeni bir group baştan başlar (birikmiş mesajları da okur), var olan group commit edilmiş offset'ten devam eder. Group verilmezse yalnız test sırasında gelen mesajlar okunur. Adımın süresi bir sonraki mesajı bekleme süresidir; bekleme süresi (varsayılan 10 sn) dolarsa adım zaman aşımıyla başarısız olur. Mesajın değeri kontrollere ve değişken çıkarmaya gider (JSON ise JSONPath ile), topic, partition, offset, timestamp ve key de header olarak okunabilir.
Saniyede 500 sipariş olayı üreten ve 10 consumer'la aynı topic'i okuyan bir test; üretim onayının p95'i 50 ms'nin, hata oranı %0,1'in altında kalmalı. Tam hali örnekler sayfasında indirilebilir.
{
"name": "Kafka: sipariş olayı üret ve tüket",
"scenarios": [
{ "name": "uretici",
"executor": { "type": "constant-arrival-rate", "rate": 500, "timeUnit": "1s",
"duration": "5m", "preAllocatedVUs": 20, "maxVUs": 100 },
"steps": [ { "id": "produce", "protocol": "kafka", "connection": "kafka",
"kafka": { "action": "produce", "topic": "orders", "key": "order-{{$uuid}}",
"value": "{\"orderId\":\"{{$uuid}}\",\"amount\":{{$randInt 10 5000}}}" } } ] },
{ "name": "tuketici",
"executor": { "type": "constant-vus", "vus": 10, "duration": "5m" },
"steps": [ { "id": "consume", "protocol": "kafka", "connection": "kafka",
"kafka": { "action": "consume", "topic": "orders",
"groupId": "spitfire-loadtest", "wait": "5s" },
"checks": [ { "type": "jsonPath", "path": "$.orderId", "op": "exists" } ] } ] }
],
"thresholds": [
{ "metric": "req_duration", "filter": { "step": "produce" }, "expr": "p(95)<50" },
{ "metric": "req_failed", "expr": "rate<0.001" }
]
}Koşu sırasında adım bazında istek/sn (üretilen ve tüketilen mesaj hızı), p95, p99, toplam gönderilen ve alınan veri ve hata türleri (yetki, bağlantı reddi, broker reddi, zaman aşımı) canlı izlenir. Consume adımının süresi uçtan uca gecikme değildir: consumer'lar geride kalmamışsa kayıt zaten beklemektedir ve süre kısalır; geride kalmışlarsa da yalnız sıradaki kaydı alma süresini görürsünüz. Uçtan uca gecikme ve consumer lag için Spitfire iki ölçüm daha yapar 0.17.0+.
Uçtan uca gecikme
Produce adımlarının gönderdiği her kayda bir spitfire-ts header'ı eklenir (gönderim anı ve kaynağı). Bir consume adımı, aynı koşuda aynı runner'da damgalanmış bir kaydı okuduğunda gönderimden consume'a geçen süreyi kafka_e2e_latency olarak kaydeder (ms; p50, p95, p99). Damga produce adımında Kayıtlara damga ekle seçeneğiyle kapatılabilir ("noStamp": true); aynı adla kendi koyduğunuz bir header'a dokunulmaz.
Neden yalnız aynı runner? Yukarıdaki saat farkı yüzünden Spitfire iki makinenin saatini hiç karşılaştırmaz: bir kayıt, onu produce eden runner aynı zamanda consume ettiğinde ölçülür (iki uç aynı monoton saati okur). Bunun sonucu: N runner'a yayılan bir koşuda kayıtların yaklaşık 1/N'i ölçülür, koşu ekranı da bunu belirtir. Başka araçların, başka koşuların ya da başka runner'ların ürettiği kayıtlar normal consume edilir ama ölçülmez. Produce ve consume adımlarını aynı testte tutun; gerçek servisinizin consumer'ları bu ölçümün dışındadır, onlar için aşağıdaki lag'e bakın.
Consumer lag: testin group'u ve gerçek servisinizin group'u
Sabit bir consumer group'la consume eden adım, o group'un kendi topic'teki lag'ini izler (Bu consumer group'un lag'ini izle; kapatmak için "noLag": true). Her Kafka adımı lagGroups ile başka group'ların lag'ini de izleyebilir; en faydalısı gerçek servisinizin group'udur: testin ürettiği yük altında servisinizin geride kalıp kalmadığını, kaldıysa ne kadar kaldığını görürsünüz. Lag, partition'ın end offset'i eksi group'un commit ettiği offset'tir (kayıt sayısı); hiç commit olmayan partition'da tuttuğu bütün kayıtlar sayılır. Group başına kafka_consumer_lag, partition başına group/partition etiketiyle kafka_consumer_lag_partition raporlanır.
- Salt okunur. Lag, Kafka admin API'siyle sorgulanır: offset'ler ve commit edilmiş offset'ler listelenir; Spitfire group'a hiç katılmaz, commit etmez, offset sıfırlamaz. Gerçek servisinizin group'unu izlemek onun çalışmasını değiştirmez.
- 2 saniyede bir, tek runner'dan. Koşu başına tek bir runner (yükün ilk payını alan; yerel koşularda CLI'ın kendisi) 2 saniyede bir sorgular, böylece her değer bir kez sayılır. Bu aralıktan kısa dalgalanmalar görünmeyebilir.
- Yetki. Bağlantının kullanıcısına topic üzerinde Describe (metadata, ListOffsets) ve izlenen her consumer group üzerinde Describe (OffsetFetch) yetkisi gerekir. Yetki yoksa koşu sürer, bir uyarı yazılır ve lag raporlanmaz.
- Sabit adlar. Topic ve group şablon değil sabit ad olmalıdır (
{{topic}}ile lag izlenmez). İki adımın izlediği group bir kez raporlanır; partition serileri group başına 64 partition'da durur, toplam yine hepsini kapsar.
Yukarıdaki testin produce adımı, gerçek servisin order-service group'unu da izlesin; eşikler uçtan uca gecikmenin p95'ini, izlenen group'ların en yüksek lag'ini ve koşu bittiğinde servisin group'unda kalan lag'i sınırlar. Testin kendi spitfire-loadtest group'u consume adımından zaten izlenir. Test bu haliyle spitfire validate ile doğrulandı.
{ "id": "produce", "name": "Sipariş olayı üret", "protocol": "kafka", "connection": "kafka",
"kafka": { "action": "produce", "topic": "orders", "key": "order-{{$uuid}}",
"value": "{\"orderId\":\"{{$uuid}}\",\"amount\":{{$randInt 10 5000}}}",
"lagGroups": [ "order-service" ] } }"thresholds": [
{ "metric": "req_duration", "filter": { "step": "produce" }, "expr": "p(95)<50" },
{ "metric": "req_failed", "expr": "rate<0.001" },
{ "metric": "kafka_e2e_latency", "expr": "p(95)<500" },
{ "metric": "kafka_consumer_lag", "expr": "max<10000" },
{ "metric": "kafka_consumer_lag", "filter": { "check": "order-service" }, "expr": "value<100" }
]Eşiklerde max izlenen hiçbir group'un bundan fazla geride kalmadığını, value koşu bittiğinde kalan lag'i sınar; "filter": { "check": "order-service" } eşiği tek bir group'a indirir. Koşu ekranındaki Kafka kartı uçtan uca gecikmenin yüzdeliklerini saniye saniye ve her group'un lag'ini koşu boyunca canlı gösterir; koşu kaydedildikten sonra son ve en yüksek lag partition kırılımıyla eklenir. İkisi de raporlarda, CLI özetinde ve karşılaştırmalarda (uçtan uca p95, en yüksek lag) yer alır.
Broker ve exporter metrikleriyle birlikte
Spitfire'ın lag ölçümü koşu süresince ve izlediğiniz group'larla sınırlıdır. Lag'i zaten Prometheus'a yazıyorsanız (örneğin bir Kafka exporter'ıyla) Observability'te bir Prometheus bağlantısına o sorguyu ekleyin: koşu sayfasının Backend sekmesi onu ve broker metriklerinizi yük basamaklarıyla hizalı gösterir, isterseniz koşunun kendi grafiklerinin üstüne çizilir. Koşu kırıldığında adında queue, lag ya da backlog geçen ve en az iki katına çıkıp 10'u aşan seriler kaynak sıralamasında öne çıkar (OpenTelemetry ile kök neden). Consumer'ların hangi hızda yetişemediğini bulmak için producer senaryosunda kırılma noktası testi kullanabilir, eşiği servisinizin group'unun lag'ine koyabilirsiniz.
Spitfire tek komutla Docker'a ya da Kubernetes'e kurulur; ücretsiz sürümde bütün test özellikleri ve protokoller açıktır.