1 पॉइंट द्वारा GN⁺ 2024-11-14 | 1 टिप्पणियां | WhatsApp पर शेयर करें
  • Kafka-compatible streaming system Bufstream 0.1.0–0.1.3 के सत्यापन में Bufstream की अपनी उपलब्धता से जुड़ी 2 समस्याएँ और सुरक्षा/सेफ्टी से जुड़ी 3 समस्याएँ मिलीं, और 0.1.3 तक पाँचों को ठीक कर दिया गया
  • टेस्ट Java Kafka Client 3.8.0 और मौजूदा Kafka/Redpanda Jepsen tests पर आधारित थे, और acks = all, enable.idempotence = true, enable.auto.commit = false, read_committed जैसी safety-first settings का इस्तेमाल किया गया
  • Bufstream की समस्याओं में consumer/producer का रुक जाना, गलत offset 0 response, transaction commit loss, और fetch API response size filtering bug के कारण acknowledged writes का loss शामिल था
  • जाँच के दौरान Kafka Java client और Kafka transaction protocol में भी Consumer.close() का अनिश्चितकाल तक block होना, अप्रत्याशित consumer offset, और aborted read·lost write·torn transaction समस्याएँ सामने आईं
  • Jepsen का मानना है कि Kafka transaction protocol client request order और transaction number की स्पष्ट गारंटी नहीं देता, इसलिए official Java client इस्तेमाल करने पर Kafka और Kafka-compatible systems की transactional safety टूट सकती है

Bufstream की संरचना और सत्यापन का दायरा

  • Kafka एक streaming system है जो replicated, sharded append-only log देता है, और Bufstream cloud environment में data governance और cost efficiency को प्राथमिकता देने वाला Kafka alternative implementation है
  • Bufstream, Kafka की तरह topic और partition देता है, और standard Kafka client के साथ काम करता है
    • producer producer.send() से record append करता है
    • consumer consumer.assign() या consumer.subscribe() से partition से bind होने के बाद consumer.poll() से record पढ़ता है
    • consumer group topic set के records की processing को बाँटकर संभालता है
  • Buf Schema Registry के साथ integrate करने पर Protocol Buffer records की जाँच करके record validation, field-level access control, और अन्य systems के साथ data format conversion को support किया जा सकता है
  • Kafka local disk और अपना replication protocol इस्तेमाल करता है, जबकि Bufstream data को सीधे object storage में लिखता है
    • object storage के replication traffic cost structure का फायदा उठाकर लागत घटाने का लक्ष्य रखता है
    • Bufstream node stateless auto-scaled VM के रूप में काम कर सकता है
  • तीन subsystems मिलकर Bufstream बनाते हैं
    • agent: Kafka API देने वाली stateless service
    • object store: record chunks को store करता है और readers को उपलब्ध कराता है
    • coordination service: फिलहाल etcd इस्तेमाल करता है, और यह तय करता है कि कौन-सा chunk commit हुआ है और records का order क्या है
  • अक्टूबर 2024 तक Bufstream केवल कुछ customers को ही deploy किया गया था, और documentation में “Apache Kafka का drop-in replacement” तथा Kafka transactions और exactly-once semantics compatibility को आगे रखा गया था, लेकिन specific safety claims बहुत अधिक नहीं थे

Client settings और transaction assumptions

  • Jepsen ने Kafka-compatible systems के पिछले tests की तरह अधिक सुरक्षित behavior पाने के लिए client settings adjust कीं
  • Producer settings

    • default acks = all का इस्तेमाल किया गया
    • Bufstream में acks = 0 storage का इंतज़ार किए बिना write acknowledge कर सकता है, जिससे committed write खो सकता है
    • acks = 1 और acks = all तब तक block करते हैं जब तक Bufstream को durable persist का भरोसा न हो जाए
    • Kafka producer के automatic retries में duplicate append रोकने के लिए default value enable.idempotence = true का इस्तेमाल किया गया
  • Consumer settings

    • ऐसे documents हैं कि auto-commit data loss की वजह बन सकता है, इसलिए आम तौर पर enable.auto.commit = false इस्तेमाल किया गया
    • जब committed offset न हो, तो default auto.offset.reset latest offset से शुरू करता है, इसलिए at-least-once delivery की गारंटी नहीं देता
    • consumer पूरा log देख सके, इसके लिए auto.offset.reset = earliest इस्तेमाल किया गया
    • Kafka transaction producer द्वारा भेजे गए record set और consumer द्वारा poll किए गए per-partition maximum offset map से मिलकर बनता है
    • transaction commit होने पर ही भेजे गए records durable होते हैं और read_committed consumer को अंततः दिखते हैं, और committed offset भी transaction में specified offset या उससे ऊपर चला जाता है
    • transaction commit न होने पर committed offset आगे नहीं बढ़ता, और write visibility consumer settings पर निर्भर कर सकती है
    • read_uncommitted consumer द्वारा aborted transaction की values पढ़ना aborted read(G1a) के रूप में classified है
    • Kafka documentation कहता है कि read_committed G1a को रोकता है और कुछ हद तक यह property guarantee करता है कि transaction की सभी writes दिखेंगी या कुछ भी नहीं दिखेगा, लेकिन Jepsen के Kafka·Redpanda·Bufstream tests में write cycle(G0 जैसा phenomenon) और G1c के कुछ रूप देखे गए

Test design

  • Jepsen ने Bufstream 0.1.0 से 0.1.3 तक और कई release candidate builds को test किया
  • Test harness में Bufstream test harness, Jepsen testing library, और Java Kafka Client 3.8.0 का इस्तेमाल हुआ
  • Execution environment

    • Debian Bookworm nodes 3–5 को LXC containers और EC2 VMs दोनों में इस्तेमाल किया गया
    • etcd के लिए 1 node, Minio के लिए 1 node, और बाकी Bufstream agent के रूप में इस्तेमाल किए गए
    • producer, consumer, admin client को bootstrap_servers में सिर्फ एक node डालकर initialize किया गया, लेकिन smart client discovery को रोका नहीं गया
  • मुख्य safety settings

    • auto-commit false
    • acks = all
    • retries 1,000
    • idempotence enabled
    • isolation level read_committed
    • auto_offset_reset = earliest
    • server-side automatic topic creation disabled
    • fault injection में process pause(SIGSTOP), crash(SIGKILL), clock skew(clock_settime), network partition(iptables) शामिल थे
    • क्योंकि Bufstream agent, object store, coordination service में बँटा है, Jepsen ने ऐसे नए tools बनाए जिनसे केवल किसी specific subsystem को target करके faults inject किए जा सकें
    • उदाहरण के लिए, समय के साथ combinations बदलते हुए सिर्फ Bufstream nodes को crash करना या केवल etcd coordinator को pause करना

Queue workload और Abort workload

  • Queue workload Kafka डेटा मॉडल के अनुरूप safety का विश्लेषण करता है
    • हर logical process producer, consumer, admin client चलाता है
    • numeric key किसी खास topic-partition की पहचान करता है
    • key exponential frequency से चुने जाते हैं, इसलिए कुछ key बार-बार और कुछ key बहुत कम access होते हैं
  • तीन बुनियादी operation इस्तेमाल किए जाते हैं
    • crash: logical process को समाप्त करके उसे नए client से बदलता है
    • subscribe या assign: consumer द्वारा poll किए जाने वाले topic या partition के set को बदलता है
    • txn, poll, send: poll या send micro-operation का sequence चलाता है
  • non-transactional workload में हर send या poll में ठीक एक ही micro-operation शामिल होता है
  • transactional workload में कई micro-operation को Kafka transaction में wrap किया जाता है
  • विश्लेषण key के हिसाब से offset-to-value mapping बनाकर error ढूंढता है
    • अगर एक ही offset पर कई value दिखें, तो inconsistent offset
    • अगर एक ही value कई offset पर दिखे, तो duplicate error
    • accepted record बिल्कुल भी observe न हो, तो lost या unseen
    • अगर poll, abort किए गए operation द्वारा भेजी गई value लौटाए, तो aborted read
    • यह भी जांचता है कि transaction अपनी ही writes observe करता है या नहीं
  • main test के बाद failures हटाकर final reads चरण में प्रवेश किया जाता है
    • हर process सभी topic-partition को offset 0 से पढ़ता है और ज्ञात सबसे ऊंचे written offset तक poll करता है
    • अगर final reads timeout हो जाएं और accepted record अब भी observe न हो, तो उसे unseen के रूप में classify किया जाता है
  • Abort workload transaction abort के बाद poll offset के behavior को track करने के लिए जोड़ा गया
    • topic को single partition, process, producer, consumer तक सीमित रखा जाता है
    • transaction किसी record को poll करने के बाद जानबूझकर abort करता है, और उसके बाद poll offset को advance, rewind, rewind-further, other में classify किया जाता है

Bufstream में मिली 5 समस्याएँ

  • Stuck consumers (#1)

    • 0.1.0 से 0.1.3-rc.8 तक final read चरण अक्सर अटक जाता था
    • consumer.poll() तुरंत खाली नतीजा लौटाता था, लेकिन log में स्वीकार किए गए हजारों records बचे रहते थे
    • यह स्थिति कई दर्जन सेकंड से लेकर 1 घंटे से अधिक समय तक बनी रहती थी
    • एक टेस्ट में पहले 120 सेकंड के दौरान 691 स्वीकार किए गए records भेजे गए, और final reads शुरू होने के समय 40 records किसी भी poller को दिखाई नहीं दिए
    • इसके बाद 1 घंटे से अधिक समय तक consumer.poll() ने कोई नतीजा नहीं लौटाया, जिससे टेस्ट timeout हो गया
    • कारण यह था कि restart हुआ Bufstream node last stable offset और high watermark की stale cached value लौटा सकता था
    • कुछ client libraries ने मान लिया कि आगे कोई record नहीं है और stall हो गईं; Bufstream ने startup के समय cache refresh करने के लिए 0.1.3-rc.6 में patch लगाया
  • Stuck producers & consumers (#2)

    • 0.1.3-rc.6 में भी coordinator, storage और Bufstream node पर pause, crash, partition के बाद unseen write की समस्या लगातार देखी गई
    • कुछ मामलों में coordinator pause के बाद सभी Bufstream nodes चल रहे होने के बावजूद client InitProducerId का इंतज़ार करते-करते timeout वाली स्थिति में पहुँच गया
    • अन्य मामलों में listOffsets node ... being disconnected या timed out waiting for a node assignment के साथ fail हुआ, और poll पूरा तो हुआ लेकिन कोई नतीजा नहीं लौटाया
    • Bufstream node को kill करके restart करने पर समस्या हल हो गई
    • कारण etcd lease से जुड़ा था
    • Bufstream agent active agents को track करने के लिए etcd leases का उपयोग करता है
    • छोटे pause या partition की वजह से etcd ने agent lease से बंधे key को delete कर दिया, लेकिन delete update agent तक नहीं पहुँच सकता था
    • agent इस बात से अनजान रह गया कि उसने अपनी lease खो दी है
    • Bufstream टीम ने अतिरिक्त polling logic जोड़ा, और 0.1.3-rc.8 में unseen write मोटे तौर पर हल हो गया
  • Spurious zero offsets (#3)

    • 0.1.0 से 0.1.3-rc.2 तक sent value को offset 0 assign होने के बाद वह असल में किसी higher offset पर दिखाई दे सकता था
    • ऐसा तब भी हुआ जब offset 0 बहुत पहले ही assign हो चुका था
    • सिर्फ sender ने offset 0 देखा, जबकि poller ने higher offset देखा
    • single Bufstream node और etcd process pause के साथ 2 मिनट के टेस्ट में 6 writes को offset 0 मिला और बाद में वे higher offset पर दिखाई दिए
    • कारण यह था कि Bufstream के error response में ज़रूरी field गायब था
    • Bufstream ने etcd को log commit request भेजी और etcd ने उसे process किया, लेकिन pause या partition की वजह से response का इंतज़ार करते समय Bufstream timeout कर सकता था
    • Bufstream ने client को error code भेजा, लेकिन sent record offset को error signal -1 पर set नहीं किया
    • Java Kafka client ने इसे offset 0 की successful response के रूप में interpret किया
    • Bufstream test suite में इस्तेमाल किए गए Franz-go ने इस message को error के रूप में interpret किया, इसलिए यह समस्या टेस्ट में सामने नहीं आई
    • Bufstream ने इसे 0.1.3-rc.6 में ठीक किया, और उसके बाद Jepsen इसे दोबारा observe नहीं कर सका
  • Lost transaction writes (#4)

    • 0.1.2 में committed transaction के कुछ records गायब हो जाते थे और फिर दोबारा observe नहीं होते थे; write loss अक्सर होता था
    • एक टेस्ट में 100 सेकंड और 6,761 write transactions के दौरान committed transaction द्वारा लिखे गए 240 records खो गए
    • उदाहरण में key 5 की value 141 offset 274 पर सफलतापूर्वक लिखी गई बताकर लौटाई गई, लेकिन सभी consumer.poll() ने उस offset को skip कर दिया
    • कारण 0.1.2 में जोड़े गए concurrency safety mechanism का bug था
    • यह mechanism Kafka transaction protocol में idempotence की कमी को कम करने के लिए producer epoch के भीतर हर transaction को unique number देता है
    • transaction number tracking logic के bug की वजह से कई epochs में कई transactions commit होने पर कुछ commits गलत तरीके से ignore हो गए
    • जो transaction committed दिखा वह असल में abort हो सकता था, या इसका उल्टा भी हो सकता था
    • Jepsen ने transaction timeout को 1 सेकंड तक कम रखा था, इसी वजह से यह bug मिला
    • Bufstream ने 0.1.2 release के कुछ घंटों के भीतर समस्या समझ ली, customers के upgrade को रोका, और customers ने 0.1.2 पर upgrade नहीं किया
    • fix 0.1.3-rc2 में शामिल किया गया
  • Server-side filtering के कारण lost writes (#5)

    • 0.1.3-rc.8 में Bufstream process या coordinator pause, या दोनों के बीच partition जैसी छोटी failures के बाद छोटी write loss window अक्सर दिखाई दी
    • data loss transaction के इस्तेमाल से स्वतंत्र रूप से हुआ
    • एक 5 मिनट के टेस्ट में 16,770 records में से 22 acknowledge हुए, लेकिन कोई भी consumer उन्हें poll नहीं कर पाया
    • कुछ records थोड़ी देर तक poller को दिखाई दिए, फिर बाद में poll से गायब भी हो गए
    • कारण लोकप्रिय Kafka web GUI के bug को workaround करने के लिए 0.1.3-rc.8 में जोड़ा गया fetch API response size limit logic था
    • filtering logic के bug ने lagging consumer से records छिपा दिए, और यह write loss जैसा दिखा
    • Bufstream ने इसे 0.1.3-rc.12 में ठीक किया

Kafka Java client और Kafka protocol की समस्याएँ

  • KIP-588: भ्रम पैदा करने वाला ProducerFencedException

    • टेस्टिंग के दौरान ProducerFencedException: There is a newer producer with the same transactionalId which fences the current one. error अक्सर आया
    • उन tests में भी यह error दिखा जहाँ सभी producers को unique transactional ID दिया गया था, इसलिए वजह समझने में समय लगा
    • KIP-588 में लिखा है कि transaction timeout पर भी ProducerFencedException throw हो सकता है
    • Kafka Java client ज़्यादातर timeouts के लिए dedicated TimeoutException इस्तेमाल करता है, लेकिन इस मामले में ProducerFencedException throw करता है
    • असल में कोई conflicting producer न होने पर भी error message कहता है कि दूसरा producer instance मौजूद है
    • KIP-588 दो साल से खुला है, और Jepsen ने Kafka team को error message बदलने की सलाह दी है
  • KAFKA-17734: Consumer.close() अनिश्चित समय तक block कर सकता है

    • Bufstream और Kafka दोनों की tests में Java client bug की वजह से हर कुछ घंटों में test रुक जाता था
    • Consumer.close() default रूप से network IO पर block करता है
    • close() का timeout parameter अनिश्चित block को रोकना चाहिए था, लेकिन यह काम नहीं कर रहा था
    • अलग thread से consumer.wakeup() call करके IO में stuck consumer को interrupt करने का तरीका भी कारगर नहीं रहा
    • Jepsen का मानना है कि long-running programs को network error होने पर भी client, connection, thread, memory जैसे resources को reasonable समय में release कर पाना चाहिए, और उसने KAFKA-17734 दर्ज किया
  • KAFKA-17582: transaction fail होने के बाद consumer offset अप्रत्याशित

    • Kafka की official documentation transaction commit fail होने पर consumer offset कैसा होना चाहिए, इस बारे में बहुत कम बताती है
    • Confluent की Kafka design documentation कहती है कि transaction abort होने पर consumer position पिछले value पर लौट आती है, लेकिन असल Java client हमेशा ऐसा व्यवहार नहीं करता
    • Abort workload के results में, healthy cluster पर भी abort के बाद behavior predict करना मुश्किल था
    • ज़्यादातर transaction pairs आगे के offset पर advance हो गए
    • कुछ पुराने offset पर rewind हुए
    • सभी rewinds rebalance event से जुड़े थे, और सभी advances में rebalance नहीं था
    • Kafka की ओर से जवाब के मुताबिक यह behavior intentional है
    • consumer लगातार advance करता रहता है
    • rebalance होने पर committed offset के आधार पर किसी भी point पर rewind हो सकता है
    • users को transaction abort पर consumer position को manually rewind करना होगा
    • Jepsen ने KAFKA-17582 खोला और इस behavior को document करने तथा transaction abort पर default rewind बदलने पर विचार करने का सुझाव दिया
    • Queue workload को भी consumer को explicitly rewind करने के लिए modify किया गया
  • KAFKA-17754: write loss, aborted read, torn transaction

    • Bufstream 0.1.0~0.1.3 में सिर्फ Bufstream process pause, coordinator pause, crash और network partition से aborted read, lost write और atomicity violation देखे गए
    • विश्लेषण Kafka transaction protocol की बुनियादी कमी तक पहुँचा
    • उदाहरण में client ने unique transactional ID jt1234 के साथ transaction चलाया और EndTxn में committed = false भेजकर abort किया, लेकिन 15 poll() calls ने aborted transaction के writes देखे
    • उसी transaction के अन्य writes किसी भी poller ने नहीं देखे
    • packet capture और Bufstream log को साथ देखने पर कारण delayed commit message था
    • कुछ transactions पहले भेजा गया commit EndTxn एक node पर देर से process हुआ
    • client पहले ही अगले transactions चला रहा था
    • delayed commit current transaction पर apply हो गया, जिससे transaction का शुरुआती हिस्सा ही commit हुआ और बाकी अलग transaction की तरह handle होकर abort हो गया
    • Kafka protocol इस तरह design किया गया है कि client कई TCP connections और कई nodes पर requests भेज सके, लेकिन उसी client के requests का order तय करने के लिए कोई sequence number नहीं है
    • transaction number की अवधारणा भी नहीं है, इसलिए server commit या abort message मिलने पर नहीं जान सकता कि client किस transaction को खत्म करना चाहता था
    • नतीजतन ये स्थितियाँ संभव हो जाती हैं
      • जो transaction commit हुआ दिखता है, वह असल में abort हो जाता है
      • abort हुआ transaction असल में commit हो जाता है
      • transaction के कुछ writes ही बचे रहते हैं और कुछ loss हो जाते हैं, जिससे torn transaction होता है
    • official Java Kafka client timeouts को retryable मानता है और कई EndTxn messages अपने-आप भेज सकता है, इसलिए user हर transaction के लिए commit या abort सिर्फ एक बार call करे तब भी समस्या हो सकती है
    • Jepsen ने Kafka में भी process pause से aborted read और torn transaction observe किए और KAFKA-17754 खोला
    • Kafka engineers का मानना है कि KIP-890 इस समस्या को ठीक कर सकता है
    • KIP-890 transaction protocol को इस तरह बदलता है कि हर transaction पर producer epoch बढ़ता है
    • server पुराने epoch messages को reject करता है, इसलिए पिछले transaction का commit message बाद के transaction में leak होने से रोका जा सकता है
    • Bufstream ने 0.1.3 में etcd revision को logical clock की तरह इस्तेमाल करके frequency घटाने वाला mechanism जोड़ा, लेकिन यह client और Bufstream के बीच reorder को नहीं रोकता
    • Jepsen ने 0.1.3 में भी aborted read, lost write और torn transaction लगातार observe किए, और माना कि client-side fix की जरूरत है

पूरे नतीजों का सारांश

  • Bufstream की अपनी 5ों समस्याएँ ठीक कर दी गईं
    • #1: lagging highest stable offset की वजह से consumer stuck हो गया; failure की जरूरत नहीं थी; 0.1.3-rc.6 में ठीक किया गया
    • #2: etcd lease expiry की वजह से producer/consumer stuck हो गए; pause की जरूरत थी; 0.1.3-rc.8 में ठीक किया गया
    • #3: spurious zero offsets; pause की जरूरत थी; 0.1.3-rc.6 में ठीक किया गया
    • #4: lost transaction writes; failure की जरूरत नहीं थी; 0.1.3-rc.2 में ठीक किया गया
    • #5: server-side filtering की वजह से lost writes; pause की जरूरत थी; 0.1.3-rc.12 में ठीक किया गया
  • Kafka से जुड़ी समस्याएँ अब भी बनी हुई हैं
    • KIP-588: transaction timeout पर गलत error message, अनसुलझा
    • KAFKA-17734: ConsumerClient.close() अनिश्चित समय तक block कर सकता है, अनसुलझा
    • KAFKA-17582: transaction fail होने के बाद consumer offset unpredictable, अनसुलझा
    • KAFKA-17754: write loss, aborted read, torn transaction, अनसुलझा
  • Jepsen चेतावनी देता है कि experimental safety verification bug के मौजूद होने को साबित कर सकता है, लेकिन उसके न होने को साबित नहीं कर सकता
  • खास तौर पर KAFKA-17754 की वजह से Bufstream में write loss के दूसरे मामले हैं या नहीं, यह तय करना मुश्किल माना गया

Bufstream users और operations के लिए सुझाव

  • official Java Kafka client से Bufstream transaction इस्तेमाल करने वाले users को ध्यान रखना चाहिए कि फिलहाल transaction सुरक्षित नहीं हो सकते
    • abort किया गया transaction वास्तव में commit हो सकता है
    • commit किया गया transaction वास्तव में abort हो सकता है
    • transaction बीच से फटकर केवल कुछ effects ही बचा सकता है
  • Bufstream मानता है कि Franz-go client इस समस्या के प्रति कम vulnerable है, लेकिन Jepsen ने इस काम जैसी techniques से Franz-go को test नहीं किया
  • दूसरे clients vulnerable हो भी सकते हैं और नहीं भी
  • Bufstream 0.1.3 से पहले के users को ये समस्याएँ हो सकती हैं
    • producer.send() वास्तविक offset की जगह गलती से 0 offset लौटाता है
    • client stuck हो जाने वाली metastable availability issue
  • Jepsen 0.1.3 upgrade की सलाह देता है
  • Bufstream की overall architecture sound लगती है, ऐसा आकलन किया गया
    • etcd जैसी coordination service से immutable data chunk का क्रम तय करने का तरीका OLTP और streaming system में precedent वाला, अपेक्षाकृत सरल approach है
  • operations के लिहाज से दो सुधार सुझाए गए
    • startup के समय storage की shared file request fail होने पर cluster crash हो सकता है, इसलिए retry जोड़ने की सलाह दी गई; Bufstream ने retry layer जोड़ दी
    • dependency unavailable होने पर agent तुरंत मरने के बजाय चलता रहे, backpressure और system status दे, और ज्यादा smooth तरीके से recover करे—ऐसी दिशा सुझाई गई
  • 0.1.3 के हिसाब से Bufstream ने etcd के लिए अतिरिक्त retry logic जोड़ा है, लेकिन online स्थिति बनाए रखने के लिए अब भी constant supervision की जरूरत है
  • users को test करना चाहिए कि process supervisor मौजूद है और लंबे outage के दौरान भी हार माने बिना चलता रहता है

Kafka transaction documentation और protocol में बदलाव की जरूरत

  • Kafka की official documentation transaction के बारे में बहुत कम कहती है, इसलिए users को अस्पष्ट और आपस में टकराते कई sources को मिलाकर समझना पड़ता है
  • Jepsen ने Kafka team को transaction semantics स्पष्ट रूप से व्यवस्थित करने वाला central document बनाने की सलाह दी और KAFKA-17671 का उल्लेख किया
  • उस document में कम-से-कम ये बातें साफ लिखी जानी चाहिए
    • consumer कब monotonically increasing offset observe करता है
    • consumer कब acknowledged record को skip कर सकता है
    • क्या rebalance transaction के बीच में असर डाल सकता है
    • producer write offset कब monotonic तरीके से बढ़ता है
    • G0, G1a, G1b, G1c, fractured read, अपने transaction write read कब valid हैं
    • abort किए गए transaction के बाद poll() return value और offset का क्या मतलब है
    • transaction error, abort के दौरान error, rewind के दौरान error को कैसे handle करना चाहिए
  • Confluent documentation बार-बार कहती है कि Kafka defaults at-least-once delivery देते हैं, लेकिन Jepsen ने बताया कि यह सही नहीं लगता
    • auto.offset.reset = latest unprocessed record को “committed” जैसा बना सकता है
    • Confluent offset management documentation भी default auto-commit में crash के समय message progress खोने के risk की बात करती है
    • transaction abort होने पर consumer rewind होता है—यह documentation भी वास्तविकता से अलग है
  • Jepsen मानता है कि Kafka transaction protocol को मूल रूप से बदलना होगा
    • protocol implicitly ordered reliable delivery मानता है, लेकिन process pause, network unreliability, non-zero latency और कई TCP sockets के बीच unordered delivery मौजूद हैं
    • Kafka protocol message को कई nodes और TCP sockets में distribute करता है, और client message को automatically retry करता है
    • उसी client के message order को restore करने के लिए sequence number और transaction target की पुष्टि के लिए transaction number नहीं है
  • KIP-890 हर transaction commit पर epoch बढ़ाकर अधिक strict order guarantee करने की कोशिश करता है
  • client library भी message acknowledge न होने पर producer को re-initialize करके epoch बढ़ाने के तरीके से मदद कर सकती है
  • Java Kafka Client 3.8.0 इस समस्या के प्रति vulnerable है
  • Jepsen मानता है कि Franz-go timeout पर re-initialize करके समस्या को कम या रोक सकता है, लेकिन उसने दूसरी client libraries की जांच नहीं की

आगे का काम

  • कई users transaction को सीधे handle करने के बजाय Kafka Streams API की “exactly-once semantics” पर निर्भर करते हैं, इसलिए आगे चलकर Streams application की correctness की जांच की जा सकती है
  • Jepsen को KAFKA-17754 की जांच के दौरान Kafka में भी unseen write मिला, लेकिन समय की कमी के कारण उसका विश्लेषण नहीं कर सका
    • unseen write, hanging transaction, stuck consumer और data loss का संकेत हो सकता है
    • यह भी सवाल बना हुआ है कि क्या देर से आया Produce message किसी future transaction में शामिल होकर transaction guarantee का उल्लंघन कर सकता है
    • Jepsen को यह भी संदेह है कि Kafka Java Client request timeout होने पर sequence number को reuse करता है, जिससे write acknowledge तो हो सकता है लेकिन चुपचाप discard भी किया जा सकता है
  • rebalance event होने पर consumer position आगे-पीछे हो सकती है, लेकिन उसके नियम स्पष्ट नहीं हैं
  • अगर Kafka अपने intended behavior को document करे, तो Jepsen उसे verify करना चाहेगा
  • Jepsen बताता है कि यह एक random process है, इसलिए दुर्लभ anomaly ढूंढना मुश्किल है
    • जो समस्या सिर्फ एक बार होती है, उसकी debugging और reproduction बहुत कठिन होती है
  • Bufstream distributed system के पूरे setup को deterministic hypervisor और simulated network में चलाने वाले Antithesis का भी इस्तेमाल करता है
    • Jepsen की workload generation और history checking को Antithesis के deterministic, replayable environment के साथ जोड़ने से testing reproducibility बढ़ सकती है

1 टिप्पणियां

 
GN⁺ 2024-11-14
Hacker News की राय
  • KAFKA-17754 जैसी issue की जांच करते हुए अगर Kafka में अदृश्य writes भी मिली हैं, तो लगता है Jepsen के लिए Kafka को फिर से गहराई से खंगालने का समय आ गया है
    आख़िरी जांच 2013 में हुई थी(https://aphyr.com/posts/293-call-me-maybe-kafka, Kafka 0.8 beta), और अब लगता है Kafka खुद ही कई समस्याओं को अभी खोजने की अवस्था में है
    “write की पुष्टि हो जाए लेकिन वह चुपचाप फेंक दी जाए” जैसी बात काफ़ी डरावनी है

    • Kafka analysis ज़रूर करना चाहूँगा :-)
  • डिफ़ॉल्ट enable.auto.commit=true में Kafka consumer का एप्लिकेशन ने वास्तव में प्रोसेस किया या नहीं, इससे अलग offset commit कर देना बहुत चौंकाने वाला है
    मैंने auto-commit को कभी इस तरह नहीं समझा था, और अगर यही default है तो यह तर्कसंगत नहीं लगता
    documentation पूरी तरह स्पष्ट नहीं है, लेकिन समग्र रूप से पढ़ने पर ऐसा लगा कि offset तभी commit होते हैं जब processing पूरी हो चुकी हो
    auto-commit interval को adjust करना कम-से-कम एक बार processing (at-least-once) में उम्मीद के मुताबिक message loss नहीं, बल्कि duplicate processing window कम करने में मदद करता है—मैं यही समझता था

    • थोड़ा चौंकाने वाला है, और मैं सहमत हूँ कि docs इस हिस्से को ठीक से समझाती नहीं हैं
      अगर आप explicitly commit नहीं करते, तो Kafka के पास यह जानने का कोई तरीका नहीं है कि message प्रोसेस हुआ या नहीं
      Kafka मान लेता है कि जो message उसने दे दिया, वह तुरंत प्रोसेस हो गया
      Auto-commit कुछ वैसा है जैसे किसी को ice cream cone थमाकर तुरंत मुड़ जाना और मान लेना कि उसने खा लिया। कोई उसे लेते ही गिरा सकता है और एक कौर भी न खा पाए
    • मुख्य बात यह है कि message का Kafka client तक सफलतापूर्वक पहुँच जाना, यह नहीं बताता कि application ने उसे प्रोसेस भी कर लिया
      अगर आपको वह guarantee चाहिए, तो आपको explicitly acknowledge करना होगा
      उदाहरण के लिए, अगर आप सिर्फ message को database में लिख रहे हैं, तो client handler callback में आते ही उसे acknowledged मान लिया जाता है
      लेकिन व्यवहार में ज़्यादा संभावना है कि आप DB insert सफल होने के बाद acknowledgment चाहते हों
      DB अगर network, Kubernetes, firewall configuration वगैरह की वजह से unreachable हो जाए, और उसी दौरान कोई engineer restart की कोशिश करे जिससे client down हो जाए, तो unprocessed messages रह जाना बहुत आसान है
    • मेरी समझ में यह feature high-performance scenarios के लिए है
      कोई दूसरा system failure detect कर सकता है, और यह feature upper bound आगे बढ़ाकर reprocessing कम कर सकता है
      लेकिन अगर timing मेल खा जाए और failure हो जाए, तो आपको मानकर चलना होगा कि restart के बाद पहले से प्रोसेस की हुई कुछ चीज़ें फिर मिल सकती हैं
      समस्या तब होती है जब auto-commit से पहले ऐसी processing मौजूद ही न हो
      पढ़ने पर लगता है कि इरादा processing के काफ़ी बाद commit करने का है, लेकिन auto-commit होते हुए भी auto-commit समय से कुछ milliseconds पहले तक की entries ही commit होनी चाहिए—यह बात कुछ विरोधाभासी लगती है
    • इस feature के अस्तित्व को कुछ हद तक उचित ठहराया जा सकता है। यह synchronous single-threaded consumer के लिए डिज़ाइन किया गया था, और मोटे तौर पर poll कॉल करने के बाद messages को durably process करने वाले loop को ध्यान में रखता है
      भ्रम की जगह यह है कि auto-commit check timeout के बाद asynchronously नहीं होता, बल्कि अगली poll call के समय होता है
      इसलिए अगर आप अगली poll call से पहले messages को durably process करने के बजाय सिर्फ store कर रहे हैं—जैसे async processing, delay, queue आदि का उपयोग कर रहे हैं—तभी writes गिर सकती हैं
      यह Java client library के documented behavior(https://kafka.apache.org/32/javadoc/org/apache/kafka/clients...) के अनुसार है; मौजूदा implementation वास्तव में ऐसा करती है या नहीं, यह अलग बात है
      Kafka protocol high-level और low-level के बीच फँसा हुआ है, इसलिए दोनों में से कोई भी काम बहुत अच्छी तरह नहीं करता
      auto-commit एक high-level feature है जो simple applications बनाना आसान करता है, लेकिन अगर इसे अपेक्षित तरीके से न इस्तेमाल किया जाए तो इसका fail होना स्वाभाविक है
      आजकल end users को Kafka client सीधे इस्तेमाल करने के बजाय ऐसे high-level implementations इस्तेमाल करने चाहिए जो details को सही ढंग से संभालें। data use case के लिए stream processing engine, और application use case के लिए durable execution engine जैसी चीज़ें
  • product page(https://buf.build/product/bufstream) को देखकर यह जिज्ञासा होती है कि “यह सिर्फ AWS या GCP VPC के अंदर चलता है और बाहर संपर्क नहीं करता” और “compression से पहले प्रति GiB $0.002” जैसी usage-based pricing दोनों बातें एक साथ कैसे संभव हैं
    लगता नहीं कि पूरा business बस honor system पर चलाया जा रहा होगा

    • intro में लिखा है कि “अक्टूबर 2024 तक Bufstream सिर्फ चुनिंदा ग्राहकों को ही दिया गया था”, इसलिए honor system भी संभव हो सकता है
      हाँ, misuse का जोखिम है, लेकिन कुछ खास ग्राहकों को आकर्षित करने के लिए यह एक सार्थक समझौता हो सकता है
    • कोई program या तो open source होता है या नहीं होता
      अगर source public नहीं है, तो “यह बाहर संपर्क नहीं करता” जैसे दावे पर कभी भरोसा नहीं करना चाहिए
  • “Kafka transaction protocol बुनियादी रूप से टूटा हुआ है और इसे revise किया जाना चाहिए” — यह सुनना काफ़ी पीड़ादायक है
    फिर भी, हमेशा की तरह जांच और लेखन शानदार है

  • सोच रहा हूँ क्या Kyle ने NATS JetStream की समीक्षा की है। वे इसके बारे में क्या सोचेंगे, जानने की उत्सुकता है

    • अभी तक समीक्षा नहीं की, लेकिन यह मांग करने वाले आप पहले व्यक्ति नहीं हैं
      कुछ लोगों ने सुझाव दिया है कि यह... क्या कहें... दिलचस्प हो सकता है :-)
  • bufstream GitHub project नहीं मिल रहा, यह कहाँ है?

    • ओह, माफ़ कीजिए। अब तक यह ठीक हो गया होगा
    • website पर https://github.com/bufbuild/buf मिला
    • लगता है bufstream खुद open source नहीं है, लेकिन https://github.com/bufbuild/bufstream-demo शायद आपकी ज़रूरत के काफ़ी क़रीब हो
      हालाँकि अजीब बात है कि इसमें भी license नहीं है
  • संबंधित blog posts और documentation पढ़ने पर लगता है कि Kafka का “exactly once delivery” उस read-process-write workflow की property के रूप में परिभाषित है जहाँ worker topic 1 से पढ़ता है और topic 2 में लिखता है, और दोनों topics एक ही logical Kafka system के भीतर होते हैं
    अगर ऐसा है, तो क्या इसे transaction कहना ज़्यादा सही नहीं होगा?

    • Kafka वास्तव में इसे transaction ही कहता है
      बस “exactly once” को देखने के दो तरीके हैं
      एक अर्थ database transaction जैसा है, जहाँ effects duplicate भी न हों और गायब भी न हों
      दूसरा topic-partition के पार message relationships की dataflow graph property के ज़्यादा क़रीब है, और ACID की consistency के कुछ अधिक निकट है
      जैसे serializable transaction systems किसी खास domain-level consistency को guarantee कर सकते हैं, वैसे ही transaction का उपयोग करके उस dataflow property तक पहुँचा जा सकता है
      उदाहरण के लिए, serializability यह guarantee करती है कि हर transaction को अलग-अलग देखने पर जो invariants सुरक्षित रहते हैं, वे concurrent execution history में भी सुरक्षित रहें
      Kafka शायद इसी तरह “exactly once semantics” तक पहुँचने की कोशिश कर रहा है
  • इसे https://www.warpstream.com/ से भ्रमित नहीं करना चाहिए

    • सही। WarpStream transactions भी support नहीं करता
  • Errata: “Transactions may observe none, part, or all” की जगह “Consumers may observe none, part, or all” होना चाहिए, ऐसा लगता है

    • दोनों सही हैं, लेकिन स्पष्टता के लिए मैंने transactions लिखा
      transaction के बाहर consumers की semantics और भी धुंधली है
      इस workload में सारी reads transaction context में होती हैं, और transaction offset commit path से गुजरती हैं
  • यह software किस काम आता है? Instrumentation? Black box?

    • Jepsen ऐसा tool है जो अगर आपको पता न हो कि यह आपके बनाए database को test कर रहा है, तो आपको रुला सकता है
      बेशक, खुशी के आँसू। क्योंकि Jepsen की दिलचस्पी मिलना अपने-आप में एक उपलब्धि है
    • यह Kafka clone है। Kafka मूल रूप से एक durable queue है