प्रतिनिधि इंटरव्यू विषय

डेटा इंजीनियरिंग इंटरव्यू: Kafka का exactly-once सेमेन्टिक्स वास्तव में क्या गारंटी देता है?

डेटाकठिन
Offer.cc संपादकीय टीमप्रकाशित अपडेट किया गया

प्रश्न

एक consumer Kafka से आर्डर इवेंट्स को पढ़ता है, उन्हें ट्रांसफॉर्म करता है, और दूसरे topic में लिखता है। इंटरव्यूअर पूछता है कि आप डुप्लिकेट्स और डेटा हानि से कैसे बचते हैं, फिर पूछता है कि वही exactly-once डिज़ाइन किसी बाहरी डेटाबेस को स्वचालित रूप से क्यों कवर नहीं करता है। आप क्या उत्तर देंगे?

प्रश्न और संदर्भ

यह डेटा-इंजीनियरिंग और स्ट्रीम-प्रोसेसिंग प्रश्न डेटा इंजीनियरों, स्ट्रीमिंग-प्लेटफ़ॉर्म इंजीनियरों और बैकएंड डेटा-इन्फ्रास्ट्रक्चर भूमिकाओं के लिए लक्षित है। Consumer क्रैश हो सकता है या rebalance हो सकता है, और कोई producer रिस्पॉन्स खोने के बाद पुनः प्रयास (retry) कर सकता है; परिणाम पहले वापस Kafka में जाते हैं। उत्तर में exactly-once को सिस्टम-व्यापी जादुई गारंटी मानने के बजाय Kafka की एटॉमिक प्रोसेसिंग को बाहरी साइड इफेक्ट्स से अलग करना चाहिए।

इंटरव्यूअर क्या मूल्यांकन करता है

  • क्या आप exactly-once को आउटपुट विजिबिलिटी और इनपुट ऑफसेट की एटॉमिक प्रगति में विभाजित कर सकते हैं।
  • क्या आप transactional producer, transactional.id, read_committed, और मैन्युअल ऑफसेट कमिट का सही उपयोग कर सकते हैं।
  • क्या आप abort, restart, fencing, और rebalance रिकवरी की व्याख्या कर सकते हैं।
  • क्या आप पहचानते हैं कि किसी डेटाबेस, सर्च इंजन, या HTTP सर्विस को अपने स्वयं के ट्रांजैक्शन, idempotency key, या reconciliation की आवश्यकता होती है।

उत्तर देने से पहले स्पष्टीकरण प्रश्न

पहले पुष्टि करें कि क्या आउटपुट Kafka में ही रहता है, क्या Kafka Streams का उपयोग किया जाता है, क्या डिप्लॉयमेंट के दौरान consumer group रोल होता है, और क्या किसी बाहरी सिस्टम को Kafka के समान ट्रांजैक्शन में कमिट करना चाहिए। फिर "प्रति इनपुट एक दृश्य आउटपुट" को "एक बाहरी साइड इफेक्ट" से अलग करें; बाद वाले के लिए गंतव्य (destination) से सहयोग की आवश्यकता होती है। लेटेंसी, ट्रांजैक्शन बैच आकार, पुनः प्रयास विंडो, और स्वीकार्य अंतराल (lag) के बारे में पूछें।

30-सेकंड का उत्तर ढांचा

मैं इनपुट रिकॉर्ड्स, ट्रांसफॉर्म किए गए आउटपुट, और consumer offsets को एक Kafka ट्रांजैक्शन में रखूँगा। Producer transactional.id सक्षम करता है; consumer ऑटो-कमिट अक्षम करता है और read_committed का उपयोग करता है। एक सफल कमिट तीनों को एक साथ दृश्यमान बनाता है, जबकि एक abort आउटपुट को छुपाता है और ऑफसेट को ट्रांजैक्शन से पहले छोड़ देता है। एक स्थिर अद्वितीय ट्रांजैक्शन ID रीस्टार्ट के बाद पुराने इंस्टेंस को फेंस (fences) करता है। यह केवल Kafka रीड्स और राइट्स को कवर करता है; एक बाहरी डेटाबेस को परिणाम और ऑफसेट को एक स्टोरेज ट्रांजैक्शन में, या idempotency, एक outbox, और reconciliation की आवश्यकता होती है।

चरण-दर-चरण समाधान

1. गारंटी की सीमा निर्धारित करें

Kafka का डिज़ाइन topic-to-topic exactly-once को एक ट्रांजैक्शन में आउटपुट रिकॉर्ड्स और consumer पोजीशन को एटॉमिक रूप से अपडेट करने के रूप में परिभाषित करता है। इसका मतलब यह नहीं है कि फ़ंक्शन एक बार चलता है या कोई यादृच्छिक HTTP अनुरोध एक बार आता है। कॉन्फ़िगरेशन पर चर्चा करने से पहले इस सीमा को स्पष्ट करें।

2. Transactional producer के साथ आउटपुट और ऑफसेट को एटॉमिक रूप से लिखें

ऑटो-कमिट अक्षम करें, एक बैच प्रोसेस करें, आउटपुट भेजें, और उस बैच के ऑफसेट को ट्रांजैक्शन के भाग के रूप में सबमिट करें। मुख्य स्यूडोकोड है:

java
producer.initTransactions();
while (running) {
  ConsumerRecords<String, Order> records = consumer.poll(timeout);
  producer.beginTransaction();
  try {
    for (ConsumerRecord<String, Order> r : records) {
      producer.send(new ProducerRecord<>("orders-enriched", r.key(), enrich(r.value())));
    }
    producer.sendOffsetsToTransaction(offsets(records), groupMetadata);
    producer.commitTransaction();
  } catch (AbortableException e) {
    producer.abortTransaction();
    consumer.seekToCommitted();
  }
}

आउटपुट और ऑफसेट को एक साथ कमिट करना आउटपुट-बिफोर-ऑफसेट विफलता से दिखने वाले डुप्लिकेट्स को रोकता है और ऑफसेट-बिफोर-आउटपुट विफलता से नुकसान को रोकता है। डिप्लॉय किए गए संस्करण के अनुसार क्लाइंट के एक्सेप्शन वर्गों को संभालें; प्रत्येक एक्सेप्शन को आँख मूंदकर पुनः प्रयास करना असुरक्षित है।

3. Consumers को केवल कमिट किए गए ट्रांजैक्शन देखने योग्य बनाएं

isolation.level=read_committed सेट करें और enable.auto.commit=false बनाए रखें। read_uncommitted निरस्त (aborted) ट्रांजैक्शन के रिकॉर्ड्स को उजागर करता है, जिससे consumer उस आउटपुट को आगे बढ़ा सकता है जिसे वापस (rolled back) लिया जाना चाहिए था। read_committed विजिबिलिटी को कमिट परिणामों के साथ संरेखित करने के लिए ट्रांजैक्शन मार्कर्स का उपयोग करता है।

4. Restart, fencing, और rebalance को संभालें

प्रत्येक सक्रिय consumer इंस्टेंस को एक स्थिर, क्लस्टर-व्यापी अद्वितीय transactional.id दें। जब कोई नया इंस्टेंस समान ID के साथ रजिस्टर होता है, तो Kafka पुराने इंस्टेंस के इन-फ़्लाइट ट्रांजैक्शन को निरस्त करता है और उसे फेंस करता है, जिससे दोनों को कमिट करने से रोका जा सके। एक abort के बाद, एप्लिकेशन को consumer पोजीशन को फिर से बनाना या स्पष्ट रूप से रिवाइंड करना होगा और बैच को पुनः प्रोसेस करना होगा; इसे ऐसे लोकल कर्सर से जारी नहीं रखना चाहिए जो निरस्त ट्रांजैक्शन के अंदर आगे बढ़ा था। पार्टीशन असाइनमेंट यह सुनिश्चित करता है कि एक समय में एक ग्रुप मेंबर के पास एक पार्टीशन हो।

5. बताएं कि बाहरी सिस्टम्स को सहयोग क्यों करना चाहिए

यदि परिणाम PostgreSQL में जाता है, तो Kafka ट्रांजैक्शन में डेटाबेस कमिट अपने आप शामिल नहीं होता है। डेटाबेस-पहले क्रैश होने पर Kafka ऑफसेट अनकमिटेड रह जाता है, इसलिए पुनः प्रयास के लिए एक अद्वितीय कुंजी (unique key) या सशर्त संस्करण अपडेट की आवश्यकता होती है; ऑफसेट-पहले कमिट करने पर राइट खो सकता है। एक अधिक मजबूत डिज़ाइन परिणाम और ऑफसेट को एक डेटाबेस ट्रांजैक्शन में रखता है, या एक रिप्ले-योग्य outbox/connector रिकॉर्ड लिखता है और गंतव्य को डुप्लिकेट हटाने और मिलान (reconcile) करने देता है। गंतव्य के सहयोग के बिना, डिटेक्टेबल डुप्लिकेट्स के साथ at-least-once निष्पादन का वादा करें, एंड-टू-एंड exactly-once का नहीं।

6. फ़ॉल्ट इंजेक्शन के साथ सत्यापित करें, कॉन्फ़िगरेशन स्नैपशॉट से नहीं

sendOffsetsToTransaction से पहले और बाद में क्रैश करें, कमिट रिस्पॉन्स खोएं, rebalance ट्रिगर करें, पुराने इंस्टेंस को फेंस करें, ट्रांजैक्शन समाप्त करें, और सिंक को अनुपलब्ध बनाएं। read_committed के साथ आउटपुट का उपभोग करें और प्रति इनपुट ID दृश्यमान आउटपुट गणना, अंतिम ऑफसेट, रीस्टार्ट के बाद का अंतराल (lag), और निरस्त ट्रांजैक्शन रिकॉर्ड्स की जांच करें। बाहरी सिंक के लिए, यूनिक-की विवादों, रिप्ले, और मिलान का अलग से परीक्षण करें। Abort दर, consumer lag, प्रोसेसिंग लेटेंसी, और fencing गणना की निगरानी करें।

उच्च गुणवत्ता वाला नमूना उत्तर

मैं पहले दावे को सीमित करूँगा। यदि इनपुट और आउटपुट दोनों Kafka हैं, तो मैं Kafka Streams या इसके समकक्ष transactional consume-transform-produce लूप का उपयोग करूँगा। Consumer ऑटो-कमिट को अक्षम करता है; producer के पास एक स्थिर अद्वितीय transactional.id होता है; और प्रत्येक बैच के आउटपुट और ऑफसेट को sendOffsetsToTransaction के साथ कमिट किया जाता है। डाउनस्ट्रीम consumers read_committed का उपयोग करते हैं, इसलिए निरस्त आउटपुट अदृश्य रहता है। रीस्टार्ट पर, वही ट्रांजैक्शन ID पुराने इंस्टेंस को फेंस करती है, और एक abort के लिए अंतिम कमिटेड ऑफसेट पर रिवाइंड करने की आवश्यकता होती है। यदि गंतव्य PostgreSQL है, तो मैं Kafka की गारंटी को एंड-टू-एंड दावे में विस्तारित नहीं करूँगा: मैं परिणाम और ऑफसेट को एक डेटाबेस ट्रांजैक्शन में रखूँगा, या outbox, idempotency key, और reconciliation का उपयोग करूँगा। फिर मैं क्रैश, खोए हुए रिस्पॉन्स, rebalance, और fencing इंजेक्ट करूँगा और प्रति इनपुट ID दृश्यमान आउटपुट, ऑफसेट, और गंतव्य स्थिति को सत्यापित करूँगा।

सामान्य गलतियाँ

  • एक idempotent producer को exactly-once कहना → यह मुख्य रूप से producer के पुनः प्रयासों पर डुप्लिकेट लॉग प्रविष्टियों को रोकता है और आउटपुट एवं ऑफसेट को एटॉमिक नहीं बनाता है → ट्रांजैक्शन और ऑफसेट कमिट जोड़ें।
  • केवल read_committed सेट करना → यह विजिबिलिटी को बदलता है लेकिन ट्रांजैक्शन को कमिट नहीं करता है या क्रैश से रिकवर नहीं करता है → ऑटो-कमिट बंद करके पूर्ण ट्रांजैक्शन लाइफसाइकिल लागू करें।
  • सक्रिय इंस्टेंसेस में एक transactional.id साझा करना → इंस्टेंसेस एक-दूसरे को फेंस करते हैं और थ्रूपुट को अस्थिर करते हैं → प्रति सक्रिय इंस्टेंस एक स्थिर अद्वितीय ID असाइन करें।
  • Abort के बाद स्थानीय स्थिति से जारी रखना → वह स्थिति अनकमिटेड बैच के अंदर हो सकती है → कमिटेड ऑफसेट को पुनः लोड करें या स्पष्ट रूप से seek करें।
  • यह दावा करना कि Kafka डेटाबेस साइड इफेक्ट को एक बार निष्पादित करता है → दोनों सिस्टम्स में स्वचालित एटॉमिक कमिट का अभाव है → गंतव्य ट्रांजैक्शन, idempotent राइट, outbox, या reconciliation का उपयोग करें।

फॉलो-अप प्रश्न और उत्तर

क्या exactly-once का अर्थ है कि व्यावसायिक फ़ंक्शन केवल एक बार चलता है?

नहीं। फ़ंक्शन चल सकता है और फिर कमिट विफलता के बाद दोबारा चल सकता है; गारंटी यह है कि Kafka का दृश्यमान आउटपुट और ऑफसेट कमिट आपस में सहमत हों। फ़ंक्शन को जहाँ संभव हो बाहरी साइड इफेक्ट्स से मुक्त रखें, या उन प्रभावों को दोहराने योग्य और मिलान योग्य (reconcilable) बनाएं।

read_committed अभी भी लेटेंसी क्यों बढ़ा सकता है?

Consumer को निरस्त रिकॉर्ड्स को छोड़ना होगा और ट्रांजैक्शन पूरा होने की प्रतीक्षा करनी होगी। खुले या लंबे ट्रांजैक्शन विजिबिलिटी लेटेंसी और अंतराल (lag) को बढ़ाते हैं। बैच आकार और टाइमआउट को सीमित करें, ट्रांजैक्शन अवधि की निगरानी करें, और कम-लेटेंसी वाले पाथ्स के लिए at-least-once प्रोसेसिंग की तुलना में लागत का मूल्यांकन करें।

जब ट्रांजैक्शन के बीच में rebalance होता है तो क्या होता है?

नया मालिक अंतिम कमिटेड ऑफसेट से फिर से शुरू करता है। वर्तमान ट्रांजैक्शन को निरस्त कर दिया जाना चाहिए; पुराना इंस्टेंस रिलीज़ हो जाता है या फेंस हो जाता है; और नया इंस्टेंस अनकमिटेड बैच को पुनः प्रोसेस करता है। जब तक वर्तमान सदस्य के पास पार्टीशन का स्वामित्व न हो, ऑफसेट कमिट न करें।

आप Kafka आउटपुट को PostgreSQL पर कैसे भेजेंगे?

व्यावसायिक परिणाम और consumer पोजीशन को एक PostgreSQL ट्रांजैक्शन में लिखें, या एक अद्वितीय इवेंट ID के साथ outbox पंक्ति लिखें और एक विश्वसनीय रिले को इसे डिलीवर करने दें। यदि वे ट्रांजैक्शन साझा नहीं कर सकते हैं, तो at-least-once डिलीवरी, यूनिक कंस्ट्रेंट्स, सशर्त संस्करण अपडेट, और मिलान चुनें; इसे exactly-once के रूप में प्रचारित न करें।

कौन से मेट्रिक्स साबित करते हैं कि डिज़ाइन काम करता है?

प्रति इनपुट इवेंट ID अद्वितीय दृश्यमान read_committed आउटपुट की गणना करें, फिर कमिटेड ऑफसेट, abort और fence गणना, consumer lag, रीस्टार्ट के बाद रिप्ले वॉल्यूम, और सिंक पर डुप्लिकेट-कुंजी विवादों को सहसंबंधित करें। फ़ॉल्ट इंजेक्शन से पहले और बाद के नमूने रखें; एक कॉन्फ़िगरेशन स्नैपशॉट इनवेरिएंट को साबित नहीं कर सकता।

सार्वजनिक स्रोत

संबंधित प्रश्न