- 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
0response, 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 को बाँटकर संभालता है
- producer
- 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 = 0storage का इंतज़ार किए बिना write acknowledge कर सकता है, जिससे committed write खो सकता है acks = 1औरacks = allतब तक block करते हैं जब तक Bufstream को durable persist का भरोसा न हो जाए- Kafka producer के automatic retries में duplicate append रोकने के लिए default value
enable.idempotence = trueका इस्तेमाल किया गया
- default
-
Consumer settings
- ऐसे documents हैं कि auto-commit data loss की वजह बन सकता है, इसलिए आम तौर पर
enable.auto.commit = falseइस्तेमाल किया गया - जब committed offset न हो, तो default
auto.offset.resetlatest 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_committedconsumer को अंततः दिखते हैं, और committed offset भी transaction में specified offset या उससे ऊपर चला जाता है - transaction commit न होने पर committed offset आगे नहीं बढ़ता, और write visibility consumer settings पर निर्भर कर सकती है
read_uncommittedconsumer द्वारा aborted transaction की values पढ़ना aborted read(G1a) के रूप में classified है- Kafka documentation कहता है कि
read_committedG1a को रोकता है और कुछ हद तक यह property guarantee करता है कि transaction की सभी writes दिखेंगी या कुछ भी नहीं दिखेगा, लेकिन Jepsen के Kafka·Redpanda·Bufstream tests में write cycle(G0 जैसा phenomenon) और G1c के कुछ रूप देखे गए
- ऐसे documents हैं कि auto-commit data loss की वजह बन सकता है, इसलिए आम तौर पर
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याsendmicro-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 वाली स्थिति में पहुँच गया - अन्य मामलों में
listOffsetsnode ... 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
0assign होने के बाद वह असल में किसी 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 नहीं कर सका
- 0.1.0 से 0.1.3-rc.2 तक sent value को offset
-
Lost transaction writes (#4)
- 0.1.2 में committed transaction के कुछ records गायब हो जाते थे और फिर दोबारा observe नहीं होते थे; write loss अक्सर होता था
- एक टेस्ट में 100 सेकंड और 6,761 write transactions के दौरान committed transaction द्वारा लिखे गए 240 records खो गए
- उदाहरण में key
5की value141offset274पर सफलतापूर्वक लिखी गई बताकर लौटाई गई, लेकिन सभी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 पर भी
ProducerFencedExceptionthrow हो सकता है - Kafka Java client ज़्यादातर timeouts के लिए dedicated
TimeoutExceptionइस्तेमाल करता है, लेकिन इस मामले मेंProducerFencedExceptionthrow करता है - असल में कोई 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 किया, लेकिन 15poll()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 मानता है और कई
EndTxnmessages अपने-आप भेज सकता है, इसलिए 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 की जगह गलती से0offset लौटाता है- 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 = latestunprocessed 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 का संकेत हो सकता है
- यह भी सवाल बना हुआ है कि क्या देर से आया
Producemessage किसी 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 टिप्पणियां
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 की पुष्टि हो जाए लेकिन वह चुपचाप फेंक दी जाए” जैसी बात काफ़ी डरावनी है
डिफ़ॉल्ट
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 कम करने में मदद करता है—मैं यही समझता था
अगर आप explicitly commit नहीं करते, तो Kafka के पास यह जानने का कोई तरीका नहीं है कि message प्रोसेस हुआ या नहीं
Kafka मान लेता है कि जो message उसने दे दिया, वह तुरंत प्रोसेस हो गया
Auto-commit कुछ वैसा है जैसे किसी को ice cream cone थमाकर तुरंत मुड़ जाना और मान लेना कि उसने खा लिया। कोई उसे लेते ही गिरा सकता है और एक कौर भी न खा पाए
अगर आपको वह guarantee चाहिए, तो आपको explicitly acknowledge करना होगा
उदाहरण के लिए, अगर आप सिर्फ message को database में लिख रहे हैं, तो client handler callback में आते ही उसे acknowledged मान लिया जाता है
लेकिन व्यवहार में ज़्यादा संभावना है कि आप DB insert सफल होने के बाद acknowledgment चाहते हों
DB अगर network, Kubernetes, firewall configuration वगैरह की वजह से unreachable हो जाए, और उसी दौरान कोई engineer restart की कोशिश करे जिससे client down हो जाए, तो unprocessed messages रह जाना बहुत आसान है
कोई दूसरा system failure detect कर सकता है, और यह feature upper bound आगे बढ़ाकर reprocessing कम कर सकता है
लेकिन अगर timing मेल खा जाए और failure हो जाए, तो आपको मानकर चलना होगा कि restart के बाद पहले से प्रोसेस की हुई कुछ चीज़ें फिर मिल सकती हैं
समस्या तब होती है जब auto-commit से पहले ऐसी processing मौजूद ही न हो
पढ़ने पर लगता है कि इरादा processing के काफ़ी बाद commit करने का है, लेकिन auto-commit होते हुए भी auto-commit समय से कुछ milliseconds पहले तक की entries ही commit होनी चाहिए—यह बात कुछ विरोधाभासी लगती है
pollकॉल करने के बाद messages को durably process करने वाले loop को ध्यान में रखता हैभ्रम की जगह यह है कि auto-commit check timeout के बाद asynchronously नहीं होता, बल्कि अगली
pollcall के समय होता हैइसलिए अगर आप अगली
pollcall से पहले 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 पर चलाया जा रहा होगा
हाँ, misuse का जोखिम है, लेकिन कुछ खास ग्राहकों को आकर्षित करने के लिए यह एक सार्थक समझौता हो सकता है
अगर source public नहीं है, तो “यह बाहर संपर्क नहीं करता” जैसे दावे पर कभी भरोसा नहीं करना चाहिए
“Kafka transaction protocol बुनियादी रूप से टूटा हुआ है और इसे revise किया जाना चाहिए” — यह सुनना काफ़ी पीड़ादायक है
फिर भी, हमेशा की तरह जांच और लेखन शानदार है
सोच रहा हूँ क्या Kyle ने NATS JetStream की समीक्षा की है। वे इसके बारे में क्या सोचेंगे, जानने की उत्सुकता है
कुछ लोगों ने सुझाव दिया है कि यह... क्या कहें... दिलचस्प हो सकता है :-)
bufstream GitHub project नहीं मिल रहा, यह कहाँ है?
हालाँकि अजीब बात है कि इसमें भी license नहीं है
संबंधित blog posts और documentation पढ़ने पर लगता है कि Kafka का “exactly once delivery” उस read-process-write workflow की property के रूप में परिभाषित है जहाँ worker topic 1 से पढ़ता है और topic 2 में लिखता है, और दोनों topics एक ही logical Kafka system के भीतर होते हैं
अगर ऐसा है, तो क्या इसे 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/ से भ्रमित नहीं करना चाहिए
Errata: “Transactions may observe none, part, or all” की जगह “Consumers may observe none, part, or all” होना चाहिए, ऐसा लगता है
transaction के बाहर consumers की semantics और भी धुंधली है
इस workload में सारी reads transaction context में होती हैं, और transaction offset commit path से गुजरती हैं
यह software किस काम आता है? Instrumentation? Black box?
बेशक, खुशी के आँसू। क्योंकि Jepsen की दिलचस्पी मिलना अपने-आप में एक उपलब्धि है