1 puan yazan GN⁺ 2024-11-14 | 1 yorum | WhatsApp'ta paylaş
  • Kafka uyumlu bir akış sistemi olan Bufstream 0.1.0~0.1.3 doğrulamasında Bufstream’in kendisine ait 2 erişilebilirlik sorunu ve 3 güvenlik sorunu bulundu; 0.1.3 itibarıyla 5’inin de düzeltildiği görüldü
  • Testler Java Kafka Client 3.8.0 ile mevcut Kafka/Redpanda Jepsen testleri temel alınarak yapıldı; acks = all, enable.idempotence = true, enable.auto.commit = false, read_committed gibi güvenliği önceleyen ayarlar kullanıldı
  • Bufstream sorunları arasında consumer ve producer’ın durması, hatalı offset 0 yanıtı, transaction commit kaybı ve fetch API yanıt boyutu filtreleme hatasından kaynaklanan onaylanmış yazma kaybı yer aldı
  • İnceleme sırasında Kafka Java client ve Kafka transaction protokolünde de Consumer.close() çağrısının süresiz bloklanması, öngörülemez consumer offset’i, aborted read·lost write·torn transaction sorunları ortaya çıktı
  • Jepsen’e göre Kafka transaction protokolü, istemci istek sırasını ve transaction numarasını açıkça garanti etmediği için resmi Java client kullanıldığında Kafka ve Kafka uyumlu sistemlerin transaction güvenliği bozulabilir

Bufstream mimarisi ve doğrulama kapsamı

  • Kafka, çoğaltılmış ve shard’lanmış append-only log sağlayan bir akış sistemidir; Bufstream ise bulut ortamlarında veri yönetişimi ve maliyet verimliliğini önceleyen, Kafka’ya alternatif bir uygulamadır
  • Bufstream, Kafka gibi topic ve partition sağlar ve standart Kafka client’larıyla çalışır
    • producer, producer.send() ile record append eder
    • consumer, consumer.assign() veya consumer.subscribe() ile partition’a bağlandıktan sonra consumer.poll() ile record okur
    • consumer group, bir topic kümesindeki record işleme görevini paylaşır
  • Buf Schema Registry ile entegre edildiğinde Protocol Buffer record’larını denetleyerek record doğrulama, field-level access control ve diğer sistemlerle veri biçimi dönüşümü desteği sunabilir
  • Kafka’nın yerel disk ve kendi çoğaltma protokolünü kullanmasının aksine, Bufstream verileri doğrudan object storage’a yazar
    • object storage’ın çoğaltma trafiği maliyet yapısından yararlanarak maliyeti düşürmeyi hedefler
    • Bufstream node’ları stateless auto-scaled VM’ler olarak çalışabilir
  • Bufstream üç alt sistemden oluşur
    • agent: Kafka API’sini sağlayan stateless servis
    • object store: record chunk’larını depolar ve reader’lara sunar
    • coordination service: şu anda etcd kullanır; hangi chunk’ın commit edildiğini ve record sırasını belirler
  • Ekim 2024 itibarıyla Bufstream yalnızca bazı müşterilere dağıtılmıştı; belgeler “Apache Kafka için drop-in replacement” olmayı ve Kafka transactions ile exactly-once semantics uyumluluğunu öne çıkarıyordu, ancak somut güvenlik iddiaları çok fazla değildi

İstemci ayarları ve transaction öncülleri

  • Jepsen, Kafka uyumlu sistemlere yönelik önceki testlerde olduğu gibi daha güvenli davranış elde etmek için client ayarlarını değiştirdi
  • Producer ayarları

    • Varsayılan acks = all kullanıldı
    • Bufstream’de acks = 0, storage beklemeden yazmayı onaylayabildiği için commit edilmiş yazmalar kaybedilebilir
    • acks = 1 ve acks = all, Bufstream durable persist olduğundan emin olana kadar bloklanır
    • Kafka producer’ın otomatik yeniden denemelerinde duplicate append’i önlemek için varsayılan değer olan enable.idempotence = true kullanıldı
  • Consumer ayarları

    • Auto-commit’in veri kaybına yol açabileceğini belirten belgeler olduğundan genel olarak enable.auto.commit = false kullanıldı
    • Committed offset yokken varsayılan auto.offset.reset en yeni offset’ten başladığı için at-least-once delivery’yi garanti etmez
    • Consumer’ın tüm log’u gözlemleyebilmesi için auto.offset.reset = earliest kullanıldı
    • Kafka transaction’ları, producer’ın gönderdiği record kümesi ile consumer’ın poll ettiği partition başına maksimum offset map’inden oluşur
    • Yalnızca transaction commit edildiğinde gönderilen record’lar durable olur ve sonunda read_committed consumer’lara görünür; committed offset de transaction’da belirtilen offset’in en azına yükselir
    • Transaction commit edilmezse committed offset ilerlemez; yazma görünürlüğü consumer ayarına bağlı olarak değişebilir
    • read_uncommitted consumer’ın abort edilmiş transaction’daki değerleri okuması aborted read (G1a) olarak sınıflandırılır
    • Kafka belgeleri read_committed’ın G1a’yı engellediğini ve transaction’daki tüm yazmaların görünmesi ya da hiçbirinin görünmemesi özelliğini bir ölçüde garanti ettiğini söylese de Jepsen’in Kafka·Redpanda·Bufstream testlerinde write cycle (G0 benzeri olgu) ve bazı G1c biçimleri gözlemlendi

Test tasarımı

  • Jepsen, Bufstream 0.1.0’dan 0.1.3’e kadar olan sürümleri ve birkaç release candidate build’i test etti
  • Test altyapısı Bufstream test harness, Jepsen testing library ve Java Kafka Client 3.8.0 kullandı
  • Çalıştırma ortamı

    • Hem LXC container hem de EC2 VM üzerinde 3~5 Debian Bookworm node’u kullanıldı
    • etcd için 1 node, Minio için 1 node, geri kalanlar Bufstream agent olarak kullanıldı
    • producer, consumer ve admin client’lar bootstrap_servers içine yalnızca tek bir node koyarak başlatıldı; ancak smart client discovery engellenmedi
  • Başlıca güvenlik ayarları

    • auto-commit false
    • acks = all
    • retries 1.000
    • idempotence enabled
    • isolation level read_committed
    • auto_offset_reset = earliest
    • server-side automatic topic creation disabled
    • Hata enjeksiyonu process pause (SIGSTOP), crash (SIGKILL), clock skew (clock_settime) ve network partition (iptables) içerdi
    • Bufstream agent, object store ve coordination service olarak ayrıldığı için yalnızca belirli bir alt sistemi hedefleyerek hata enjekte edebilen yeni Jepsen araçları geliştirildi
    • Örneğin yalnızca Bufstream node’unu crash etmek veya yalnızca etcd coordinator’ı pause etmek gibi kombinasyonlar zaman içinde değiştirildi

Queue workload ve Abort workload

  • Queue workload, Kafka veri modeline göre güvenliği analiz eder
    • Her mantıksal süreç bir producer, consumer ve admin client çalıştırır
    • Sayısal key, belirli bir topic-partition’ı tanımlar
    • Key’ler üstel frekansla seçilir; bu yüzden bazı key’lere sık, bazılarına seyrek erişilir
  • Üç temel operation kullanılır
    • crash: Mantıksal süreci sonlandırır ve yeni bir client ile değiştirir
    • subscribe veya assign: Consumer’ın poll edeceği topic veya partition kümesini değiştirir
    • txn, poll, send: poll veya send micro-operation’larından oluşan bir sequence yürütür
  • Non-transactional workload’da her send veya poll tam olarak yalnızca bir micro-operation içerir
  • Transactional workload’da birden çok micro-operation bir Kafka transaction’ı içine sarılır
  • Analiz, key bazında offset-to-value mapping oluşturduktan sonra hataları arar
    • Aynı offset’te birden çok value görülürse inconsistent offset
    • Aynı value birden çok offset’te görülürse duplicate error
    • Kabul edilmiş bir record hiç gözlemlenmezse lost veya unseen
    • Abort edilmiş bir operation’ın gönderdiği value’yu poll döndürürse aborted read
    • Bir transaction’ın kendi yazımlarını gözlemleyip gözlemlemediği de kontrol edilir
  • Ana testten sonra arızalar giderilir ve final reads aşamasına geçilir
    • Her process, tüm topic-partition’ları offset 0’dan başlayarak okur ve bilinen en yüksek written offset’e kadar poll eder
    • final reads timeout olursa ve kabul edilmiş record hâlâ gözlemlenmemişse unseen olarak sınıflandırılır
  • Abort workload, transaction abort sonrasında poll offset davranışını izlemek için eklendi
    • Topic tek bir partition, process, producer ve consumer ile sınırlandırılır
    • Transaction bir record’u poll ettikten sonra kasıtlı olarak abort eder; ardından poll offset’i advance, rewind, rewind-further, other olarak sınıflandırılır

Bufstream’de bulunan 5 sorun

  • Takılı kalan consumer’lar (#1)

    • 0.1.0’dan 0.1.3-rc.8’e kadar final read aşaması sık sık takılı kaldı
    • consumer.poll() hemen boş sonuç döndürdü, ancak loglarda onaylanmış binlerce record kalmıştı
    • Bu durum onlarca saniyeden 1 saatin üzerine kadar sürdü
    • Bir testte ilk 120 saniye boyunca onaylanmış 691 record gönderildi ve final reads başladığında bunların 40’ı hiçbir poller tarafından gözlemlenmedi
    • Ardından 1 saatten uzun süre consumer.poll() sonuç döndürmediği için test timeout oldu
    • Neden, yeniden başlatılan Bufstream node’unun last stable offset ve high watermark için bayat önbelleğe alınmış değer döndürebilmesiydi
    • Bazı client library’ler daha ileride record olmadığına karar verip stall etti; Bufstream, startup sırasında cache’i refresh edecek şekilde 0.1.3-rc.6’da patch uyguladı
  • Takılı kalan producer ve consumer’lar (#2)

    • 0.1.3-rc.6’da da coordinator, storage ve Bufstream node’larına yönelik pause, crash ve partition sonrasında unseen write sorunu gözlemlenmeye devam etti
    • Bazı durumlarda coordinator pause sonrasında tüm Bufstream node’ları çalışıyor olmasına rağmen client, InitProducerId beklerken timeout olan bir duruma girdi
    • Diğer durumlarda listOffsets, node ... being disconnected veya timed out waiting for a node assignment hatalarıyla başarısız oldu; poll tamamlandı ama sonuç döndürmedi
    • Bufstream node’u kill edip restart edince sorun çözüldü
    • Neden etcd lease ile ilgiliydi
    • Bufstream agent’ı, active agent takibi için etcd leases kullanıyor
    • Kısa bir pause veya partition nedeniyle etcd, agent lease’ine bağlı key’i sildi; ancak silme update’i agent’a iletilmemiş olabiliyordu
    • Agent, kendi lease’ini kaybettiğini bilmez durumda kalıyordu
    • Bufstream ekibi ek polling logic ekledi ve 0.1.3-rc.8’de unseen write büyük ölçüde çözüldü
  • Sahte sıfır offset’ler (#3)

    • 0.1.0’dan 0.1.3-rc.2’ye kadar, sent value offset 0 aldıktan sonra daha yüksek gerçek offset’te görünebiliyordu
    • offset 0 çok daha önce atanmış olduğunda bile bu yaşandı
    • offset 0’ı yalnızca sender gözlemledi; poller ise daha yüksek offset’i gözlemledi
    • Tek bir Bufstream node’u ve etcd process pause içeren 2 dakikalık testte 6 write offset 0 aldı, sonra daha yüksek offset’lerde göründü
    • Neden, Bufstream’in error response’unda gerekli bir field’ın eksik olmasıydı
    • Bufstream etcd’ye log commit isteği gönderiyor ve etcd bunu işliyordu; ancak pause veya partition nedeniyle Bufstream response beklerken timeout olabiliyordu
    • Bufstream client’a error code gönderdi, fakat sent record offset’ini hata sinyali olan -1 olarak ayarlamadı
    • Java Kafka client bunu offset 0 için başarılı yanıt olarak yorumladı
    • Bufstream test suite’inin kullandığı Franz-go bu message’ı error olarak yorumladığı için sorun testlerde ortaya çıkmadı
    • Bufstream bunu 0.1.3-rc.6’da düzeltti ve Jepsen sonrasında tekrar gözlemleyemedi
  • Kaybolan transaction write’ları (#4)

    • 0.1.2’de commit edilmiş transaction’lardaki bazı record’ların kaybolduğu ve tekrar gözlemlenmediği write loss sık sık yaşandı
    • Bir testte 100 saniye ve 6.761 write transaction boyunca, commit edilmiş transaction’ların yazdığı 240 record kayboldu
    • Örnekte key 5 için value 141, offset 274’e başarıyla yazılmış olarak döndü; ancak tüm consumer.poll() çağrıları bu offset’i atladı
    • Neden, 0.1.2’ye eklenen concurrency safety mechanism içindeki bir bug’dı
    • Bu mechanism, Kafka transaction protocol’ünün idempotence eksikliğini hafifletmek için producer epoch içindeki her transaction’a unique number veriyor
    • transaction number tracking logic’indeki bug nedeniyle, birden çok epoch boyunca birden çok transaction commit edilirken bazı commit’ler yanlışlıkla yok sayıldı
    • Commit edilmiş gibi görünen transaction gerçekte abort edilmiş olabiliyor veya bunun tersi yaşanabiliyordu
    • Jepsen, transaction timeout değerini 1 saniye gibi düşük tuttuğu için bu bug’ı buldu
    • Bufstream, 0.1.2 release’inden birkaç saat sonra sorunu tespit etti, müşterilerin upgrade yapmasını engelledi ve müşteriler 0.1.2’ye upgrade etmedi
    • Düzeltme 0.1.3-rc2’ye dahil edildi
  • Server-side filtering nedeniyle lost writes (#5)

    • 0.1.3-rc.8’de Bufstream process veya coordinator pause’u ya da ikisi arasındaki partition gibi küçük arızalardan sonra kısa bir write loss penceresi sık sık görüldü
    • Data loss, transaction kullanılıp kullanılmamasından bağımsız olarak gerçekleşti
    • 5 dakikalık bir testte 16.770 record’dan 22’si acknowledge edildi, ancak hiçbir consumer bunları poll edemedi
    • Bazı record’lar bir süre poller’lara görünüp daha sonra poll’dan kayboldu
    • Neden, popüler bir Kafka web GUI’sindeki bug’ı aşmak için 0.1.3-rc.8’e eklenen fetch API response size sınırlama logic’iydi
    • filtering logic’indeki bug, lagging consumer’lardan record’ları gizledi ve write loss gibi göründü
    • Bufstream bunu 0.1.3-rc.12’de düzeltti

Kafka Java client ve Kafka protokolü sorunları

  • KIP-588: Yanıltıcı ProducerFencedException

    • Test sırasında sık sık ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one. hatası oluştu
    • Tüm producer’ların benzersiz transactional ID aldığı testlerde bile bu hata görüldüğü için nedenini anlamak zaman aldı
    • KIP-588, transaction timeout durumunda da ProducerFencedException fırlatılabileceğini belirtiyor
    • Kafka Java client, çoğu timeout için özel TimeoutException kullanıyor; ancak bu durumda ProducerFencedException fırlatıyor
    • Gerçekte çakışan bir producer olmamasına rağmen hata mesajı ikinci bir producer instance’ı varmış gibi söylüyor
    • KIP-588 iki yıldır açık ve Jepsen, Kafka ekibine error message’ı değiştirmesini öneriyor
  • KAFKA-17734: Consumer.close() süresiz olarak block edebilir

    • Hem Bufstream hem Kafka testlerinde, Java client bug’ı nedeniyle testler birkaç saatte bir durdu
    • Consumer.close() varsayılan olarak network IO üzerinde block ediyor
    • close() içindeki timeout parameter’ın süresiz block etmeyi engellemesi gerekiyordu, ancak çalışmadı
    • Ayrı bir thread’den consumer.wakeup() çağırarak IO’da takılı kalan consumer’ı interrupt etme yöntemi de etkili olmadı
    • Jepsen, uzun süre çalışan programların network error durumlarında bile client, connection, thread ve memory gibi resource’ları makul sürede serbest bırakabilmesi gerektiğini düşünüyor ve KAFKA-17734’ü açtı
  • KAFKA-17582: Transaction başarısızlığından sonra consumer offset’i öngörülemez

    • Kafka’nın resmi dokümantasyonu, transaction commit başarısız olduğunda consumer offset’in ne olması gerektiği konusunda neredeyse hiçbir şey söylemiyor
    • Confluent’in Kafka design documentation’ı, transaction abort edilirse consumer position’ın önceki değerine döndüğünü söylüyor; ancak gerçek Java client her zaman böyle davranmıyor
    • Abort workload sonuçlarında, healthy cluster’da bile abort sonrası davranışın öngörülmesinin zor olduğu görüldü
    • Transaction pair’lerinin çoğu daha ilerideki offset’e advance ediyor
    • Bazıları önceki offset’e rewind ediliyor
    • Tüm rewind’lar rebalance event ile ilişkiliydi; tüm advance’lerde ise rebalance yoktu
    • Kafka tarafının yanıtına göre bu davranış intentional
    • Consumer advance etmeye devam ediyor
    • Rebalance gerçekleşirse committed offset’e göre rastgele bir noktaya rewind edilebilir
    • Kullanıcıların transaction abort sırasında consumer position’ı elle rewind etmesi gerekiyor
    • Jepsen KAFKA-17582’yi açarak bu davranışın dokümante edilmesini ve transaction abort sırasında varsayılan rewind davranışının değiştirilmesinin değerlendirilmesini önerdi
    • Queue workload da consumer’ı açıkça rewind edecek şekilde düzeltildi
  • KAFKA-17754: Write loss, aborted read, torn transaction

    • Bufstream 0.1.0~0.1.3’te yalnızca Bufstream process pause, coordinator pause, crash ve network partition ile aborted read, lost write ve atomicity violation gözlemlendi
    • Analiz, Kafka transaction protocol’ündeki temel bir kusura işaret etti
    • Örnekte client, benzersiz transactional ID jt1234 ile bir transaction çalıştırıp EndTxn için committed = false göndererek abort etti; ancak 15 poll() çağrısı abort edilmiş transaction’ın write’larını gözlemledi
    • Aynı transaction’daki diğer write’lar ise hiçbir poller tarafından gözlemlenmedi
    • Packet capture ve Bufstream log’ları birlikte incelendiğinde nedenin gecikmiş bir commit message olduğu görüldü
    • Birkaç transaction önce gönderilen commit EndTxn, bir node’da geç işlendi
    • Client bu sırada sonraki transaction’lara çoktan devam etmişti
    • Gecikmiş commit mevcut transaction’a uygulandı; transaction’ın yalnızca baş kısmı commit edildi, geri kalanı ayrı bir transaction gibi işlenerek abort edildi
    • Kafka protocol, client’ın birden fazla TCP connection ve birden fazla node’a request gönderebilmesine izin verecek şekilde tasarlanmış; ancak aynı client’tan gelen request’lerin sırasını belirleyen bir sequence number yok
    • Transaction number kavramı da olmadığı için server, commit veya abort message aldığında client’ın hangi transaction’ı bitirmeye çalıştığını bilemiyor
    • Bunun sonucunda aşağıdaki durumlar mümkün hale geliyor
      • Commit edilmiş gibi görünen bir transaction gerçekte abort edilebilir
      • Abort edilmiş bir transaction gerçekte commit edilebilir
      • Transaction içindeki write’ların yalnızca bir kısmının korunup bir kısmının kaybolduğu torn transaction oluşabilir
    • Resmi Java Kafka client timeout’ları retryable kabul ediyor ve otomatik olarak birden fazla EndTxn message gönderebiliyor; bu nedenle kullanıcı her transaction için commit veya abort’u yalnızca bir kez çağırsa bile sorun çıkabiliyor
    • Jepsen, Kafka’da da process pause ile aborted read ve torn transaction gözlemledi ve KAFKA-17754’ü açtı
    • Kafka engineer’ları KIP-890’ın bu sorunu düzeltme ihtimali olduğunu düşünüyor
    • KIP-890, her transaction için producer epoch’u artıran bir yöntemle transaction protocol’ünü değiştiriyor
    • Server eski epoch message’larını reddettiği için geçmiş transaction’a ait commit message’ın sonraki transaction’a sızmasını engelleyebilir
    • Bufstream 0.1.3’te etcd revision’ı logical clock olarak kullanıp sıklığı azaltan bir mechanism ekledi; ancak client ile Bufstream arasındaki reorder’ı engelleyemiyor
    • Jepsen, 0.1.3’te de aborted read, lost write ve torn transaction gözlemlemeye devam etti ve client tarafında çözüm gerektiğini düşünüyor

Genel sonuç özeti

  • Bufstream’in kendi sorunlarının 5’i de düzeltildi
    • #1: lagging highest stable offset nedeniyle consumer stuck oluyordu; hata gerekmeden ortaya çıkıyordu, 0.1.3-rc.6’da düzeltildi
    • #2: etcd lease expiry nedeniyle producer/consumer stuck oluyordu; pause gerekiyordu, 0.1.3-rc.8’de düzeltildi
    • #3: spurious zero offsets; pause gerekiyordu, 0.1.3-rc.6’da düzeltildi
    • #4: lost transaction writes; hata gerekmeden ortaya çıkıyordu, 0.1.3-rc.2’de düzeltildi
    • #5: server-side filtering nedeniyle lost writes; pause gerekiyordu, 0.1.3-rc.12’de düzeltildi
  • Kafka ile ilgili sorunlar hâlâ duruyor
    • KIP-588: transaction timeout durumunda hatalı error message, çözülmedi
    • KAFKA-17734: ConsumerClient.close() süresiz block edebilir, çözülmedi
    • KAFKA-17582: transaction başarısızlığından sonra consumer offset’i öngörülemez, çözülmedi
    • KAFKA-17754: write loss, aborted read, torn transaction, çözülmedi
  • Jepsen, deneysel güvenlik doğrulamasının bug’ların varlığını kanıtlayabileceğini ama yokluğunu kanıtlayamayacağını hatırlatıyor
  • Özellikle KAFKA-17754 nedeniyle Bufstream’de başka write loss vakaları olup olmadığını ayırt etmenin zor olduğu görüşünde

Bufstream kullanıcıları ve operasyon önerileri

  • Resmî Java Kafka client ile Bufstream transaction kullanan kullanıcılar, mevcut transaction’ların güvenli olmayabileceğini dikkate almalı
    • abort edilmiş bir transaction gerçekte commit edilmiş olabilir
    • commit edilmiş bir transaction gerçekte abort edilmiş olabilir
    • bir transaction ikiye bölünerek yalnızca bazı etkileri korunabilir
  • Bufstream, Franz-go client’ın bu soruna daha az açık olduğunu düşünüyor; ancak Jepsen, Franz-go’yu bu çalışmadaki tekniklerle test etmedi
  • Diğer client’lar savunmasız olabilir de olmayabilir de
  • Bufstream 0.1.3 öncesini kullananlar şu sorunlarla karşılaşabilir
    • producer.send() gerçek offset yerine hatalı şekilde 0 offset döndürebilir
    • client’ın stuck olduğu metastable availability issue
  • Jepsen 0.1.3’e upgrade yapılmasını öneriyor
  • Bufstream’in genel mimarisinin sound göründüğünü değerlendiriyor
    • etcd gibi bir coordination service ile immutable data chunk sırasını belirleme yöntemi, OLTP ve streaming system’larda örnekleri olan görece basit bir yaklaşım
  • Operasyon tarafında iki iyileştirme öneriliyor
    • startup sırasında storage’ın shared file isteği başarısız olursa cluster crash edebildiği için retry eklenmesi önerildi; Bufstream de bir retry layer ekledi
    • dependency unavailable olduğunda agent’ın hemen ölmesi yerine çalışmaya devam edip backpressure ve system status sağlaması, daha yumuşak recover etmesi öneriliyor
  • 0.1.3 itibarıyla Bufstream, etcd için ek retry logic koydu; ancak online durumda kalmak için hâlâ constant supervision gerekiyor
  • Kullanıcılar bir process supervisor bulunduğunu ve uzun süreli outage sırasında da vazgeçmeden çalıştığını test etmeli

Kafka transaction dokümantasyonu ve protokol düzeltmesi gerekiyor

  • Kafka’nın resmî dokümantasyonu transaction konusunda neredeyse hiçbir şey söylemediği için kullanıcılar muğlak ve birbiriyle çelişen çeşitli source’ları birleştirmek zorunda kalıyor
  • Jepsen, Kafka ekibine transaction semantics’i net biçimde toparlayan merkezi bir doküman oluşturmasını önerdi ve KAFKA-17671’den bahsetti
  • Bu doküman en azından şunları belirtmeli
    • consumer’ın ne zaman monotonically increasing offset gözlemlediği
    • consumer’ın ne zaman acknowledge edilmiş record’ları atlayabileceği
    • rebalance’ın transaction ortasında etki edip edemeyeceği
    • producer write offset’inin ne zaman monoton arttığı
    • G0, G1a, G1b, G1c, fractured read, kendi transaction write’ını read etmenin ne zaman yasal olduğu
    • abort edilmiş transaction sonrasında poll() dönüş değerinin ve offset’in ne anlama geldiği
    • transaction error, abort sırasında error, rewind sırasında error durumlarının nasıl ele alınacağı
  • Confluent dokümanı, Kafka varsayılanlarının at-least-once delivery sağladığını tekrar tekrar söylüyor; ancak Jepsen bunun doğru görünmediğine dikkat çekiyor
    • auto.offset.reset = latest, işlenmemiş record’ları “committed” gibi gösterebilir
    • Confluent offset management dokümanı da varsayılan auto-commit’te crash sırasında message progress kaybı riskinden bahsediyor
    • transaction abort olduğunda consumer’ın rewind edildiğini söyleyen dokümantasyon da gerçekle uyuşmuyor
  • Jepsen, Kafka transaction protocol’ünün temelden düzeltilmesi gerektiğini düşünüyor
    • protocol örtük olarak ordered reliable delivery varsayıyor; ancak process pause, network unreliability, non-zero latency ve birden fazla TCP socket arasında unordered delivery var
    • Kafka protocol, message’ları birden fazla node ve TCP socket’e dağıtıyor; client ise message’ları otomatik retry ediyor
    • Aynı client message sırasını geri kuracak sequence number ve transaction hedefini doğrulayacak transaction number yok
  • KIP-890, her transaction commit’inde epoch’u artırarak daha sıkı bir sıra garantilemeye çalışıyor
  • Client library de message acknowledge edilmediğinde producer’ı re-initialize edip epoch’u artırarak yardımcı olabilir
  • Java Kafka Client 3.8.0 bu soruna açık
  • Jepsen, Franz-go’nun timeout sırasında re-initialize yaparak sorunu hafifletebileceğini veya önleyebileceğini düşünüyor; ancak diğer client library’leri incelemedi

Gelecekteki çalışmalar

  • Birçok kullanıcı transaction’ları doğrudan ele almak yerine Kafka Streams API’nin “exactly-once semantics” özelliğine güvendiğinden, gelecekte Streams uygulamalarının doğruluğu incelenebilir
  • Jepsen, KAFKA-17754’ü incelerken Kafka’da da unseen write ile karşılaştı, ancak zaman kısıtları nedeniyle analiz edemedi
    • unseen write, hanging transaction, stuck consumer ve data loss belirtisi olabilir
    • Gecikmiş bir Produce message’ın gelecekteki bir transaction’a girerek transaction guarantee’yi ihlal edip edemeyeceği de soru işareti olarak kaldı
    • Kafka Java Client’ın request timeout durumunda sequence number’ı yeniden kullanması nedeniyle write’ın acknowledge edilmiş olsa da sessizce discard edilme olasılığından da şüpheleniliyor
  • Bir rebalance event meydana geldiğinde consumer position ileri geri hareket edebilir, ancak bunun kuralları belirsiz
  • Kafka amaçlanan davranışı belgelediğinde Jepsen bunu doğrulamak istiyor
  • Jepsen, random process olduğu için nadir anomaly’leri keşfetmenin zor olduğunu açıklıyor
    • Bir kez ortaya çıkan sorunların debugging ve reproduction’ı çok zor
  • Bufstream, dağıtık sistemin tamamını deterministic hypervisor ve simulated network üzerinde çalıştıran Antithesis’i de kullanıyor
    • Jepsen’in workload generation ve history checking’ini Antithesis’in deterministic, replayable environment’ı ile birleştirmek testlerin yeniden üretilebilirliğini artırabilir

1 yorum

 
GN⁺ 2024-11-14
Hacker News yorumları
  • KAFKA-17754 gibi sorunları araştırırken Kafka’da görünmez yazmalar da bulunduysa, Jepsen’in Kafka’yı yeniden derinlemesine inceleme zamanı gelmiş gibi görünüyor
    Son inceleme 2013’teydi (https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta) ve şimdi Kafka’nın kendisinde birden fazla sorunun yeni yeni bulunmaya başlandığı bir aşamadaymışız gibi duruyor
    “Bir yazma onaylanmış olsa da sessizce atılabilir” gibi bir şey epey ürkütücü

    • Kafka analizi yapmayı kesinlikle isterim :-)
  • Varsayılan enable.auto.commit=true ayarında Kafka tüketicisinin, uygulamanın gerçekten işleyip işlemediğinden bağımsız olarak offset’i commit edebilmesi çok şaşırtıcı
    Otomatik commit’i hiç böyle anlamamıştım; varsayılan buysa bana mantıksız geliyor
    Dokümantasyon açıklaması çok net değil, ama genel olarak offset’in yalnızca işleme bittiyse commit edildiği şeklinde okunuyordu
    Otomatik commit aralığını ayarlamanın, en az bir kez işleme (at-least-once) beklentisindeki gibi mesaj kaybını değil, yinelenen işleme penceresini daraltmaya yardımcı olduğunu düşünmüştüm

    • Biraz şaşırtıcı ve dokümantasyonun bu kısmı iyi açıklamadığına katılıyorum
      Açıkça commit etmezseniz Kafka’nın mesajı işleyip işlemediğinizi bilmesinin yolu yok
      Kafka, verdiği mesajın hemen işlendiğini varsayar
      Otomatik commit, birine dondurma külahı verip hemen arkanızı dönerek onun bunu yediğini varsaymaya benzer. Bazı insanlar alır almaz düşürebilir ve tek lokma bile yiyemeyebilir
    • Esas nokta şu: Mesajın Kafka istemcisine başarıyla teslim edilmiş olması, uygulamanın onu işlediği anlamına gelmez
      Bu garantiyi istiyorsanız açıkça onay yanıtı vermeniz gerekir
      Örneğin tek yaptığınız mesajı bir veritabanına yazmaksa, mesaj istemci handler callback’ine girdiği anda onaylanmış kabul edilir
      Ama gerçekte muhtemelen DB ekleme işlemi başarılı olduktan sonra onaylanmasını istersiniz
      DB ağ, Kubernetes, güvenlik duvarı ayarları vb. nedeniyle erişilemez hale gelir ve bu sırada bir mühendis yeniden başlatmayı denerken istemci kapanırsa, işlenmemiş mesajların ortaya çıkması kolaydır
    • Bu özelliğin yüksek performans senaryoları için olduğunu düşünüyorum
      Başka bir sistem başarısızlık olup olmadığını belirleyebilir ve bu özellik üst sınır konumunu ilerleterek yeniden işlemeyi azaltabilir
      Ancak zamanlama denk gelir ve arıza yaşanırsa, yeniden başlatmadan sonra zaten işlenmiş bazılarını tekrar alabileceğinizi varsaymalısınız
      Sorun, otomatik commit’ten önce böyle bir işlemenin olmadığı durumlardır
      Okuyunca, commit’in işleme sonrasından epey sonra yapılması amaçlanmış gibi görünüyor; ama otomatik commit iken yalnızca otomatik commit anından birkaç milisaniye önceki öğelerin commit edilmesi gerektiği fikri de çelişkili görünüyor
    • Bu özelliğin varlığı bir ölçüde gerekçelendirilebilir. Senkron, tek iş parçacıklı tüketiciler için tasarlandı ve kabaca poll çağrısından sonra mesajları kalıcı biçimde işleyen bir döngü varsayar
      Kafa karıştıran nokta, otomatik commit kontrolünün zaman aşımından sonra asenkron olarak değil, bir sonraki poll çağrısı sırasında gerçekleşmesidir
      Bu yüzden, yeniden poll çağırmadan önce mesajları kalıcı biçimde işlemeyip yalnızca saklıyorsanız — örneğin asenkron işleme, gecikme, kuyruklar vb. kullanıyorsanız — ancak o zaman yazmaları düşürebilmesi gerekir
      Bu, Java istemci kütüphanesinin dokümante edilmiş davranışına (https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...) göredir; mevcut implementasyonun gerçekten böyle olup olmadığı ayrı konu
      Kafka protokolü yüksek seviye ile düşük seviye arasında sıkışmış durumda ve ikisini de pek iyi yapamıyor
      Otomatik commit, basit uygulamaları kolaylaştırmaya yardımcı olan yüksek seviyeli bir özellik; ama beklenen şekilde kullanmazsanız elbette başarısız olabilir
      Bence günümüzde son kullanıcılar Kafka istemcisini doğrudan kullanmak yerine ayrıntıları düzgün şekilde ele alan daha yüksek seviyeli implementasyonlar kullanmalı. Veri amaçlıysa akış işleme motoru, uygulama amaçlıysa sürekli çalışan bir yürütme motoru gibi
  • Ürün sayfasına (https://buf.build/product/bufstream) bakınca, “yalnızca AWS veya GCP VPC içinde çalışır ve dışarıyla iletişim kurmaz” açıklaması ile “sıkıştırma öncesi GiB başına $0.002” şeklindeki kullanıma dayalı ücretlendirme nasıl bir arada var olabiliyor merak ediyorum
    Herhalde tüm işi bir onur sistemiyle yürütmüyorlardır

    • Tanıtımda “Ekim 2024 itibarıyla Bufstream yalnızca seçilmiş müşterilere dağıtıldı” dendiğine göre, onur sistemi de mümkün olabilir diye düşünüyorum
      Elbette kötüye kullanım riski var, ama belirli müşterileri çekmek için değerli bir ödün olabilir
    • Program ya açık kaynaktır ya da değildir
      Kaynak kodu açık değilse “dışarıyla iletişim kurmaz” iddiasına asla güvenilmemeli
  • “Kafka transaction protokolü temelden bozuk ve revize edilmeli” demek kulağa acı geliyor
    Yine de her zamanki gibi inceleme ve yazı harika

  • Kyle’ın NATS JetStream’i hiç inceleyip incelemediğini merak ediyorum. Ne düşüneceğini merak ettim

    • Henüz incelemedim, ama bunu isteyen ilk kişi değilsiniz
      Bazıları bunun… nasıl desem… ilginç olacağını önerdi :-)
  • bufstream GitHub projesini bulamadım; nerede olduğunu merak ediyorum

  • İlgili blog yazılarını ve dokümanları okuyunca, Kafka’nın “exactly-once delivery” kavramı, bir çalışanın topic 1’den okuyup topic 2’ye yazdığı ve iki topic’in de aynı mantıksal Kafka sistemi içinde olduğu oku-işle-yaz işinin bir özelliği olarak tanımlanıyor gibi görünüyor
    Doğruysa buna transaction demek daha iyi olmaz mı diye düşünüyorum

    • Kafka da bunu gerçekten transaction olarak adlandırıyor
      Ancak “exactly once”a bakmanın iki yolu var
      Biri, veritabanı transaction’larında olduğu gibi etkinin yinelenmemesi veya kaybolmaması gerektiği anlamı
      Diğeri, topic-partition’lar arasındaki mesaj ilişkilerine dair bir veri akışı grafiği özelliğine daha yakın ve ACID’deki tutarlılığa biraz daha benziyor
      Serileştirilebilir transaction sistemlerinin belirli alan düzeyi tutarlılıkları garanti etmesine benzer şekilde, transaction’ları kullanarak bu veri akışı özelliğine ulaşabilirsiniz
      Örneğin serileştirilebilirlik, her transaction’a ayrı ayrı bakıldığında korunan invariant’ların eşzamanlı yürütme geçmişlerinde de korunduğunu garanti eder
      Kafka’nın “exactly-once semantics”e bu şekilde ulaşmaya çalıştığı söylenebilir
  • https://www.warpstream.com/ ile karıştırılmamalı

    • Doğru. WarpStream transaction’ları da desteklemiyor
  • Düzeltme: “Transactions may observe none, part, or all” yerine “Consumers may observe none, part, or all” olmalı gibi geliyor

    • İkisi de doğru, ama açıklık için transaction dedim
      Transaction dışındaki tüketicilerin semantiği daha bulanık
      Bu iş yükündeki tüm okumalar transaction bağlamında gerçekleşiyor ve transaction offset commit yolundan geçiyor
  • Bu yazılımın nerede kullanıldığını merak ediyorum. Enstrümantasyon mu? Kara kutu mu?

    • Jepsen, geliştirdiğiniz veritabanını test ettiğini bilmiyorsanız sizi ağlatabilecek bir araçtır
      Tabii sevinç gözyaşları. Çünkü Jepsen’in ilgisini çekmek başlı başına bir başarıdır
    • Kafka klonu. Kafka genel olarak kalıcı bir kuyruktur