Ana sayfa → Teknik
Queue Pathologies
Sabah 09:14. Emir kuyruğunda 40 bin bekleyen mesaj var ve sayı artıyor. Servisler ayakta, CPU rahat, veritabanı iyi. Sebep tek bir şey: kuyruğun başındaki bir tane bozuk mesaj.
- Poison message sonsuza kadar retry edilir. DLQ yoksa kuyruk kilitlenir.
- Jitter’sız exponential backoff, sorunu ritme sokar. Herkes aynı anda tekrar dener.
- Head-of-line blocking: büyük iş küçükleri bekletir. Ayrı lane aç ya da paralelliği artır.
- Sıra bozulur (message reordering). En sağlam çözüm sıraya bağımlı olmayan mesaj tasarlamak: delta değil, mutlak değer gönder.
- Sınırsız buffer bir çözüm değil, sadece geciktirilmiş bir bellek çökmesi.
1. Poison message: kuyruğu kilitleyen tek satır
Bir emir mesajı geldi ve ayrıştırılamıyor. Belki bir alan eksik, belki
volume alanında sayı yerine boş string var, belki karakter kodlaması
bozuk. Consumer hata alıyor, ack göndermiyor, mesaj kuyrukta kalıyor
ve yeniden teslim ediliyor.
Sonra aynı şey oluyor. Ve tekrar. Ve tekrar. Bu, saniyede yüzlerce kez dönen bir döngü ve arkadaki 40 bin emir sırada bekliyor.
Hiçbir servis çökmez. Sağlık kontrolleri yeşil. CPU normal. Grafiklerde tek anormal şey kuyruk uzunluğu — ki ona bakmıyorsan olayı ancak müşteri aradığında öğreniyorsun.
İlaç: dead letter queue (DLQ)
Kural basit: bir mesaj N kez retry edildiyse artık denenmez, kenara alınır.
void isle(Mesaj m) {
try {
emriIsle(m);
ack(m);
} catch (KaliciHata e) { // ayristirma, dogrulama, is kurali
dlqYaz(m, e); // ← hic tekrar deneme, direkt kenara
ack(m);
} catch (GeciciHata e) { // ag, zaman asimi, 503
if (m.denemeSayisi >= 4) { // 4 deneme yapildi: 100,200,400,800 ms
dlqYaz(m, e);
ack(m);
} else {
m.denemeSayisi++;
gecikmeliTekrarKuyruguna(m, bekleme(m.denemeSayisi));
ack(m); // asil kuyruktan CIKAR
}
}
}
İki ayrım kritik ve çoğu ekip birincisini atlıyor:
- Kalıcı hata hiç denenmez. Ayrıştırılamayan bir mesaj beşinci denemede de ayrıştırılamayacak. Tekrar denemek sadece zaman ve kaynak yakıyor. Sunucudan 400 aldıysan tekrar deneme; 503 aldıysan dene.
- Tekrar denenecek mesaj asıl kuyruktan çıkarılır. Ayrı bir delay queue’ya taşınır. Yerinde bırakırsan tekrar denemek ile sıradakileri bekletmek aynı şey olur. Bedeli var: kuyruktan çıkan mesaj sırasını da kaybeder — 4. bölümdeki sıra bozulmasının sebeplerinden biri tam olarak bu, ve orada cevabı var.
Kimsenin bakmadığı bir dead letter queue, mesajları çöpe atmakla aynı şey — üstelik “biz önlem aldık” rahatlığıyla. Üç şey şart:
- Sayaç:
dlq_mesaj_toplammetriği. - Alarm: DLQ’ya tek bir emir düşerse haber ver. DLQ’daki emir sayısı sıfır olmalı; bir bile anormal.
- Replay: Düzelttikten sonra mesajları asıl kuyruğa geri basacak bir düğme. Elle SQL yazarak yapılan iş, gece 3’te yapılmaz.
2. Exponential backoff — ve jitter’ın önemi
Geçici hatalarda hemen tekrar denemek işi kötüleştirir: karşı taraf zaten zorlanıyorsa, ikinci istek onu daha da zorlar. Klasik cevap her denemede beklemeyi katlamak:
100 ms → 200 ms → 400 ms → 800 ms → DLQ
Doğru ama eksik. Şunu düşün: köprü bir saniyeliğine düştü ve 800 emir aynı anda hata aldı. Hepsi 100 ms bekleyip aynı anda tekrar deniyor. Hepsi yine hata alıyor, hepsi 200 ms bekleyip yine aynı anda deniyor.
Yani rastgele bir yük yerine, senkronize dalgalar üretiyorsun. Köprü toparlanmaya çalışırken düzenli aralıklarla 800’lük tokat yiyor.
long bekleme(int deneme) {
// deneme 1..4 -> 100, 200, 400, 800 ms (semayla ayni)
long taban = Math.min(TABAN_MS * (1L << (deneme - 1)), TAVAN_MS);
// "full jitter": 0 ile taban arasinda rastgele
return ThreadLocalRandom.current().nextLong(taban + 1);
}
800 emir artık her denemede 0 ile o anki bekleme süresi arasına yayılıyor: ilk denemede 0–100 ms, dördüncüde 0–800. Köprü sabit bir yükle karşılaşıyor ve nefes alabiliyor. Tek satır ve etkisi grafiklerde anında görünüyor.
Fiyat saniyeler içinde değişir. Bir emri 800 ms bekletip tekrar denemek makul, ama 30 saniye sonra göndermek makul değil — artık müşterinin istediği fiyat orada değil.
O yüzden emir mesajlarında gecerlilik_bitis alanı var. Tekrar denemeden
önce bakılıyor: süresi geçmişse denenmiyor, doğrudan reddedilip müşteriye
bildiriliyor. Geç gelen doğru cevap, yanlış cevaptır.
3. Head-of-line blocking: büyük emir küçükleri bekletir
Tek bir Redis stream var, bütün emirler oradan akıyor. Kurumsal bir müşteri 1000 lotluk bir emir gönderiyor; bu emir 40 parçaya bölünüp köprülere dağıtılıyor ve işlenmesi 2 saniye sürüyor.
O 2 saniye boyunca arkadaki 0,01 lotluk emirler bekliyor. Grafikte gecikme ortalaması iyi görünüyor ama %99’luk dilim tavan yapıyor.
| Çözüm | Nasıl | Ne zaman |
|---|---|---|
| Ayrı lane | Büyük emirler ayrı kuyruğa, ayrı consumer’a | En basit ve en etkilisi. İlk denenecek. |
| Partitioning | Hesap kimliğine göre partition; bir hesabın büyük emri sadece kendi partition’ını bekletir | Aynı hesabın sırası korunmalıysa (detay) |
| Paralellik | Tek consumer yerine N consumer | Sıra önemli değilse; en ucuzu |
| Chunking | Büyük işi mesaj düzeyinde küçük parçalara böl | İş bölünebiliyorsa; kuyruğu hiç şişirmez |
Bizim yaptığımız birinci ve dördüncünün karışımıydı: 100 lot üstü emirler ayrı lane’e alındı ve orada parçalara bölünüp işlendi. Küçük emirlerin gecikmesi %99’luk dilimde 2100 ms’den 40 ms’ye indi — ortalama zaten iyiydi, sorun hep uçlardaydı (tail latency).
4. Message reordering: #3 önce, #1 sonra
Kısmi fill’ler (partial fill) sırayla gönderiliyor: #1 (10 lot), #2 (15 lot), #3 (5 lot). Sana geliş sırası: #3, #1, #2. Kümülatif hesap anlık olarak yanlış çıkıyor ve arayüzde müşteri tuhaf sayılar görüyor.
Neden olur?
- Mesajlar farklı partition’lara düştü, farklı consumer’lar işledi.
- #1 bir kez hata alıp tekrar denendi, bu sırada #2 ve #3 geçti.
- Ağda farklı yollar, farklı gecikmeler.
Üç çözüm, en zayıftan en sağlama
Aynı emrin bütün mesajları aynı partition’a gitsin (anahtar: order_id).
Böylece tek bir consumer sırayla işler. Basit ve etkili — ama rebalance
anında garanti yine kayboluyor.
// Beklenen: 1. Gelen: 3.
// al() ve zamanAsimi() ayni durumu (beklenenSira, bekleyenler) degistirir: ya ikisi de
// synchronized, ya da timer mesajlari isleyen thread'de calisir (tek thread'li executor).
// Aksi halde ayni mesaj iki kez islenir ya da sayac bozulur.
synchronized void al(Mesaj m) {
if (m.sira < beklenenSira) return; // gec gelen eski mesaj: ya islendi ya atlandi, sirayi KAYDIRMA
if (m.sira > beklenenSira) {
bekleyenler.put(m.sira, m); // kenarda tut
return;
}
isle(m);
beklenenSira++;
bekleyendekileriBosalt(); // 1 geldiyse kenardaki 2 ve 3 de islenir
}
// AYRI bir timer, orn. her 50 ms. Mesaj akisi durursa yalnizca burasi calisir;
// zaman asimini "yeni mesaj gelince" kontrol edersen akis kesildiginde hic tetiklenmez.
synchronized void zamanAsimi() {
if (!bekleyenler.isEmpty() && bekleyenler.enEskiYasi() > 200) { // 200 ms'den fazla bekledi
uyar("eksik mesaj: " + beklenenSira);
beklenenSira = bekleyenler.ilkAnahtar(); // atla, kilitlenme
bekleyendekileriBosalt(); // atladiktan sonra kenardakiler islenir
}
}
Çalışıyor ama bir kusuru var: eksik mesajı beklerken sen de tıkanıyorsun. O yüzden timeout şart, yoksa 3. bölümdeki hastalığı kendi elinle üretirsin.
// Kirilgan — sira sart, delta gonderiyor
{ "order_id": 991, "fill_no": 2, "eklenen_lot": 15 }
// Saglam — sira onemsiz, mutlak deger gonderiyor
{ "order_id": 991, "fill_no": 2,
"kumulatif_lot": 25, "kumulatif_tutar": 27512.50,
"kalan_lot": 5 }
İkincisinde #3 önce gelirse ne olur? Kümülatif değeri yazarsın. Sonra #1 gelir,
daha küçük bir fill_no taşıdığı için yok sayılır.
Sonuç yine doğru.
UPDATE pozisyonlar
SET kumulatif_lot = :kumulatif_lot,
son_fill_no = :fill_no
WHERE order_id = :order_id
AND COALESCE(son_fill_no, 0) < :fill_no; -- ← eski mesaj yazamaz; NULL < x false doner, o yuzden COALESCE
İlk fill için satır henüz yoksa UPDATE hiçbir şey yapmaz; gerçek kod
PostgreSQL’de INSERT ... ON CONFLICT (order_id) DO UPDATE ... WHERE ile upsert olmalı,
aynı COALESCE koşulu DO UPDATE’in WHERE’ine
gider.
Bu, fencing token ile aynı fikir: doğruluğu zamanlamaya değil sıraya bağla, ve sırayı verinin içine göm.
Aradaki fark tasarım seviyesinde: delta gönderen sistem her mesajın gelmesine ve sırasına muhtaç; mutlak değer gönderen sistem tek bir mesajla kendini toparlıyor. Yeni bir mesaj şeması tasarlarken sorulacak soru şu: “Bu mesajı tek başına alsam, doğru duruma varabilir miyim?”
5. Backpressure: 50 bin tick, yetişemeyen consumer
Piyasa hareketlendi. Saniyede 50 bin tick geliyor, streaming-engine 30 bin işleyebiliyor. Aradaki 20 bin nereye gidiyor? Dört seçeneğin var ve sonuncusu bir seçenek değil:
| Seçenek | Sonuç | Ne zaman doğru |
|---|---|---|
| Drop | Veri kaybı, ama sistem ayakta | Fiyat tick’i, metrik, canlı yayın |
| Slow down (producer’ı yavaşlat) | Gecikme artar, kayıp olmaz | Fill, bakiye — producer’ın bekleyebildiği yerde |
| Reject | Producer açık hata alır; kayıp yok, gecikme yok | Emir — bekletmek fiyatı eskitir (aşağıda) |
| Unbounded buffer | Bellek biter, süreç ölür | Hiçbir zaman |
Sonuncusu en çok yapılanı, çünkü kod yazarken en görünmezi. Sınırsız bir kuyruk oluşturmak tek satır ve o satır, sistemin ne zaman öleceğini belirsizleştiriyor. Her kuyruğa bir üst sınır koy ve sınır dolduğunda ne olacağını bilinçli seç.
Fiyat için: drop değil, conflation
Fiyat tarafında güzel bir numara var. EURUSD için 40 tick birikmişse, 40’ını da işlemenin bir anlamı yok — müşteri sadece sonuncuyu görecek. Bu yüzden kuyrukta biriktirmek yerine sembol başına son değeri tutuyorsun:
// Kuyruk degil, sembol basina tek slot
ConcurrentHashMap<String, Kotasyon> sonFiyat = new ConcurrentHashMap<>();
// Uretici: uzerine yaz, biriktirme
sonFiyat.put(k.sembol, k);
// Tuketici: yetisebildigi hizda oku; okudugunu CIKAR,
// yoksa yeni tick gelmese de her turda ayni fiyati yeniden yayinlarsin
for (String s : sonFiyat.keySet()) {
Kotasyon k = sonFiyat.remove(s);
if (k != null) yayinla(k);
}
Bellek sabit: sembol sayısı kadar. Yük ne olursa olsun büyümüyor ve müşteri her zaman en güncel fiyatı görüyor — ara tick’leri kaçırmak zaten sorun değil.
Emir için: bounded queue ve dürüst ret
Emir düşürülemez. Ama sınırsız da biriktirilemez. Doğru davranış, kuyruk dolduğunda yeni emri açıkça reddetmek:
BlockingQueue<Emir> kuyruk = new ArrayBlockingQueue<>(10_000);
if (!kuyruk.offer(emir)) { // dolu
metrik.artir("emir_reddedildi_kuyruk_dolu");
return hata("SYSTEM_BUSY, lutfen tekrar deneyin");
}
Kulağa kötü geliyor ama alternatifi daha kötü: emri kabul edip 40 saniye sonra eski fiyattan açmak. Hızlı ret, geç kabulden iyidir — müşteri en azından ne olduğunu biliyor ve kararını kendisi veriyor.
Ne izlemeli?
- Queue depth — artıyorsa consumer yetişemiyor.
- En eski mesajın yaşı — derinlikten daha iyi bir sinyal; tıkanmayı doğrudan gösterir.
- DLQ sayacı — emir tarafında sıfır olmalı, alarm bağlı olmalı.
- Tekrar deneme sayacı — zıplarsa dışarıda bir sorun var demektir.
- Gecikmenin %99’luk dilimi — ortalama değil. Head-of-line blocking önce burada görünür; ortalama uzun süre rahat kalır.
- Drop edilen tick sayısı — kayıp bilinçliyse bile ölçülmeli.
Kontrol listesi
- Kalıcı hata ile geçici hata ayrılıyor mu? (Kalıcı olan hiç denenmemeli.)
- Deneme üst sınırı var mı, sonrasında DLQ’ya gidiyor mu?
- DLQ’nun sayacı ve alarmı var mı?
- DLQ’dan replay yolu var mı? (Elle SQL sayılmaz.)
- Backoff’ta jitter var mı?
- Emirlerin son kullanma süresi kontrol ediliyor mu?
- Büyük işler küçükleri bekletiyor mu? Ayrı lane var mı?
- Mesajlar mutlak değer mi taşıyor, delta mı?
- Her kuyruğun bir üst sınırı var mı?
- Sınır dolduğunda ne oluyor — drop mu, slow down mu, reject mi? Bilinçli mi seçildi?
Sonuç
Kuyruk, sistemler arasına konan bir yastık. Ama yastık aynı zamanda bir gizleyici: arkasındaki yavaşlığı bir süre saklıyor, sonra hepsini birden önüne koyuyor.
Baştaki 40 bin mesajlık kuyruk, tek bir bozuk emirden çıkmıştı. Düzeltmesi on beş satırdı. Asıl ders, o on beş satırın ilk günden yazılması gerektiğiydi — çünkü zehirli mesaj bir ihtimal değil, bir kesinlik.