Jepsen’in Bufstream 0.1 doğrulama sonuçları
(jepsen.io)- 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_committedgibi güvenliği önceleyen ayarlar kullanıldı - Bufstream sorunları arasında consumer ve producer’ın durması, hatalı offset
0yanı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()veyaconsumer.subscribe()ile partition’a bağlandıktan sonraconsumer.poll()ile record okur - consumer group, bir topic kümesindeki record işleme görevini paylaşır
- producer,
- 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 = allkullanıldı - Bufstream’de
acks = 0, storage beklemeden yazmayı onaylayabildiği için commit edilmiş yazmalar kaybedilebilir acks = 1veacks = 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 = truekullanıldı
- Varsayılan
-
Consumer ayarları
- Auto-commit’in veri kaybına yol açabileceğini belirten belgeler olduğundan genel olarak
enable.auto.commit = falsekullanıldı - Committed offset yokken varsayılan
auto.offset.reseten 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 = earliestkullanı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_committedconsumer’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_uncommittedconsumer’ı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
- Auto-commit’in veri kaybına yol açabileceğini belirten belgeler olduğundan genel olarak
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_serversiç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ştirirsubscribeveyaassign: Consumer’ınpolledeceği topic veya partition kümesini değiştirirtxn,poll,send:pollveyasendmicro-operation’larından oluşan bir sequence yürütür
- Non-transactional workload’da her
sendveyapolltam 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,
InitProducerIdbeklerken timeout olan bir duruma girdi - Diğer durumlarda
listOffsets,node ... being disconnectedveyatimed out waiting for a node assignmenthatalarıyla başarısız oldu;polltamamlandı 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
0aldı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
0aldı, 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
-1olarak ayarlamadı - Java Kafka client bunu offset
0iç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
- 0.1.0’dan 0.1.3-rc.2’ye kadar, sent value offset
-
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
5için value141, offset274’e başarıyla yazılmış olarak döndü; ancak tümconsumer.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
ProducerFencedExceptionfırlatılabileceğini belirtiyor - Kafka Java client, çoğu timeout için özel
TimeoutExceptionkullanıyor; ancak bu durumdaProducerFencedExceptionfı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
- Test sırasında sık sık
-
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 ediyorclose()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
jt1234ile bir transaction çalıştırıpEndTxniçincommitted = falsegöndererek abort etti; ancak 15poll()ç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
EndTxnmessage 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ı şekilde0offset 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
Producemessage’ı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
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ü
Varsayılan
enable.auto.commit=trueayarı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
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
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
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
pollçağrısından sonra mesajları kalıcı biçimde işleyen bir döngü varsayarKafa 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şmesidirBu 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 gerekirBu, 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
Elbette kötüye kullanım riski var, ama belirli müşterileri çekmek için değerli bir ödün olabilir
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
Bazıları bunun… nasıl desem… ilginç olacağını önerdi :-)
bufstream GitHub projesini bulamadım; nerede olduğunu merak ediyorum
Ancak garip şekilde bunun da lisansı yok
İ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
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ı
Düzeltme: “Transactions may observe none, part, or all” yerine “Consumers may observe none, part, or all” olmalı gibi geliyor
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?
Tabii sevinç gözyaşları. Çünkü Jepsen’in ilgisini çekmek başlı başına bir başarıdır