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

डेटा इंजीनियरिंग इंटरव्यू: हॉट Kafka पार्टीशन का निदान और समाधान कैसे करें?

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

प्रश्न

एक Kafka टॉपिक में 24 पार्टीशन हैं और पीक पर प्रति सेकंड 120,000 रिकॉर्ड प्राप्त होते हैं। एक टेनेंट 45% ट्रैफ़िक उत्पन्न करता है, और प्रोड्यूसर tenant_id द्वारा पार्टीशन करता है, इसलिए एक पार्टीशन लगातार पीछे छूटता जाता है जबकि अधिकांश अन्य कंज्यूमर्स निष्क्रिय रहते हैं। एक कंज्यूमर वर्तमान में प्रति सेकंड 8,000 रिकॉर्ड प्रोसेस कर सकता है, केवल एक ऑर्डर के भीतर ही ऑर्डरिंग (क्रम) की आवश्यकता है, और आप घटना (incident) के दौरान पार्टीशन नहीं जोड़ सकते। आप समस्या का निदान, शमन और स्थायी समाधान कैसे करेंगे? ऑफ़सेट, रिबैलेंस, डुप्लिकेट प्रभाव और माइग्रेशन सत्यापन को कवर करें।

प्रॉम्प्ट और यह कब लागू होता है

एक Kafka टॉपिक में 24 पार्टीशन हैं और पीक पर प्रति सेकंड 120,000 रिकॉर्ड प्राप्त होते हैं। एक टेनेंट 45% ट्रैफ़िक उत्पन्न करता है, और प्रोड्यूसर tenant_id द्वारा पार्टीशन करता है, इसलिए एक पार्टीशन में लगातार लैग जमा होता रहता है जबकि अधिकांश अन्य कंज्यूमर्स निष्क्रिय रहते हैं। एक कंज्यूमर वर्तमान हैंडलर के साथ प्रति सेकंड 8,000 रिकॉर्ड संभाल सकता है। व्यवसाय को केवल एक order_id के भीतर ऑर्डरिंग की आवश्यकता है, न कि टेनेंट के हर इवेंट में, और घटना के दौरान पार्टीशन नहीं जोड़े जा सकते हैं।

बताएं कि आप यह कैसे साबित करेंगे कि लक्षण कंज्यूमर, ब्रोकर या डाउनस्ट्रीम विफलता के बजाय की (key) के तिरछेपन (skew) के कारण हुआ था। फिर वर्णन करें कि आप आज बैकलॉग वृद्धि को कैसे कम करेंगे, पार्टीशन की और माइग्रेशन को फिर से कैसे डिज़ाइन करेंगे, और एसिंक्रोनस प्रोसेसिंग जोड़ने के बाद सुरक्षित रूप से ऑफ़सेट कैसे कमिट करेंगे। उन मेट्रिक्स और विफलता परीक्षणों के साथ समाप्त करें जो सुधार को प्रदर्शित करते हैं।

यह डेटा इंजीनियरिंग और स्ट्रीमिंग-प्लेटफ़ॉर्म समस्या निवारण (troubleshooting) का प्रश्न है। एक वर्तमान सार्वजनिक Kafka साक्षात्कार गाइड अनिवार्य रूप से समान संयोजन प्रस्तुत करता है: एक हॉट पार्टीशन, उसका कंज्यूमर पीछे छूट रहा है, अन्यत्र निष्क्रिय कंज्यूमर्स, एक कीड (keyed) प्रोड्यूसर, और तत्काल कोई पार्टीशन-संख्या में बदलाव नहीं। यह उम्मीदवारों से प्रोड्यूसर पार्टीशनिंग, कंज्यूमर समानांतरता (parallelism) और ऑर्डरिंग को जोड़ने के लिए कहता है। Apache Kafka का दस्तावेज़ीकरण कहता है कि डिफ़ॉल्ट प्रोड्यूसर मौजूद की (key) को हैश करके एक पार्टीशन चुनता है। सिमेंटिक पार्टीशनिंग चुने गए पार्टीशन के भीतर लोकैलिटी और ऑर्डर को सुरक्षित रखती है, लेकिन यह एक की (key) के ट्रैफ़िक को भी वहीं केंद्रित कर देती है। Huawei Cloud का Kafka मार्गदर्शन भी इसी तरह कहता है कि एक पार्टीशन को एक समय में केवल एक कंज्यूमर द्वारा ही प्रोसेस किया जा सकता है और अस्थायी रूप से पार्टीशन जोड़ना मौजूदा पार्टीशन के बैकलॉग को साफ़ करने का त्वरित तरीका नहीं है।

इस प्रॉम्प्ट में थ्रूपुट, ट्रैफ़िक शेयर, कंज्यूमर क्षमता और ऑर्डरिंग का दायरा साक्षात्कार की धारणाएँ हैं, किसी विशेष कंपनी के उत्पादन आंकड़े नहीं हैं।

साक्षात्कारकर्ता क्या मूल्यांकन करता है

पहला संकेत यह है कि क्या आप समग्र मेट्रिक्स से पार्टीशन-स्तरीय साक्ष्य पर जा सकते हैं। टॉपिक-व्यापी लैग, औसत कंज्यूमर CPU और कंज्यूमर की संख्या सभी एक हॉट पार्टीशन को छिपा सकते हैं। एक मजबूत उत्तर प्रत्येक पार्टीशन की प्रोड्यूस दर, कंज्यूम दर, लैग स्लोप, लीडर ब्रोकर, की (key) वितरण और डाउनस्ट्रीम प्रोसेसिंग समय को एक ही टाइमलाइन पर संरेखित करता है। यह प्रोड्यूसर स्क्यू को धीमे हैंडलर, बार-बार रिबैलेंस या ब्रोकर संसाधन दबाव से अलग करता है।

दूसरा संकेत यह है कि क्या आप इस घटना को परिमाणित (quantify) करते हैं:

text
120,000 × 45% = 54,000 records/second

यदि उस पार्टीशन को सौंपा गया कंज्यूमर प्रति सेकंड केवल 8,000 रिकॉर्ड प्रोसेस कर सकता है, तो बैकलॉग इस प्रकार बढ़ता है:

text
54,000 - 8,000 = 46,000 records/second
46,000 × 600 = 27,600,000 records in 10 minutes

यह गणना बताती है कि सामान्य कंज्यूमर्स को जोड़ने से इस पार्टीशन की सीमा क्यों नहीं बदलती है और अलर्ट थ्रेशोल्ड बदलने से घटना क्यों नहीं सुधरती है।

तीसरा संकेत यह है कि क्या आप वास्तविक ऑर्डरिंग इनवेरिएंट की पहचान करते हैं। वर्तमान की (key) ऑर्डरिंग डोमेन को पूरे टेनेंट तक विस्तारित करती है, जबकि व्यवसाय को केवल एक ऑर्डर के भीतर ही क्रम की आवश्यकता होती है। की (key) को order_id में बदलने से कार्डिनैलिटी बढ़ती है और नए ऑर्डर वितरित होते हैं, लेकिन केवल तभी जब माइग्रेशन एक ऑर्डर को दो पार्टीशन या टॉपिक में जाने से रोकता है।

चौथा संकेत ऑफ़सेट की शुद्धता (correctness) है। एक बार जब किसी एक पार्टीशन के रिकॉर्ड वर्कर पूल में चलते हैं, तो पूरा होने का क्रम ऑफ़सेट क्रम से भिन्न हो सकता है। यदि 104 के चलते रहने के दौरान ऑफ़सेट 105 पूरा हो जाता है, तो 105 तक कमिट करने से क्रैश के बाद 104 छूट सकता है। एक मजबूत डिज़ाइन उच्चतम निरंतर पूर्ण ऑफ़सेट (highest contiguous completed offset) को ट्रैक करता है और डाउनस्ट्रीम प्रभावों को इडेम्पोटेंट बनाता है क्योंकि पूर्ण लेकिन अनकमिटेड रिकॉर्ड फिर से चलाए (replay) जा सकते हैं।

अंतिम संकेत यह है कि क्या आप कठोर सीमा (hard boundary) का उल्लेख करते हैं। अधिक पार्टीशन अधिक समानांतर स्लॉट बनाते हैं, लेकिन वे एक विशाल एंटिटी को विभाजित नहीं करते हैं जो एक ही की (key) से बंधी रहती है। एक पारंपरिक कंज्यूमर समूह में अधिक कंज्यूमर्स दो कंज्यूमर्स को एक साथ एक ही पार्टीशन का स्वामित्व लेने की अनुमति नहीं देते हैं। यदि केवल एक ऑर्डर सुरक्षित एकल-पार्टिशन क्षमता से अधिक है और उसे कड़ाई से क्रमित रहना चाहिए, तो शेष उपाय सीरियल पाथ को अनुकूलित करना, इसे थ्रॉटल करना, या स्वतंत्र व्यावसायिक अनुक्रमों को फिर से परिभाषित करना हैं।

उत्तर देने से पहले स्पष्ट करने योग्य प्रश्न

  • क्या प्रति टेनेंट, प्रति ऑर्डर, या किसी छोटे इवेंट स्ट्रीम के लिए ऑर्डरिंग आवश्यक है? यदि टेनेंट-व्यापी क्रम अनिवार्य है, तो टेनेंट को विभाजित करना अमान्य है। यदि ऑर्डर-व्यापी क्रम पर्याप्त है, तो order_id अधिक सटीक पार्टीशन सीमा है।
  • क्या 45% रिकॉर्ड की संख्या, बाइट्स या प्रोसेसिंग लागत का वर्णन करता है? बड़े रिकॉर्ड या महंगे डाउनस्ट्रीम राइट्स लागत में तिरछापन पैदा कर सकते हैं, भले ही रिकॉर्ड की संख्या संतुलित दिखाई दे। रिकॉर्ड, बाइट्स और हैंडलर समय का निरीक्षण करें।
  • क्या प्रोड्यूस दर बढ़ी, या कंज्यूम क्षमता घटी? 8,000-रिकॉर्ड क्षमता के मुकाबले स्थिर 54,000 रिकॉर्ड प्रति सेकंड की (key) के तिरछेपन और प्रति-की अपर्याप्त क्षमता को इंगित करता है। यदि इनपुट अपरिवर्तित है लेकिन खपत 8,000 से घटकर 2,000 हो जाती है, तो पहले डाउनस्ट्रीम सिस्टम, कचरा संग्रहण (garbage collection), नेटवर्क, डिस्क और रिबैलेंस की जांच करें।
  • कंज्यूमर कहाँ लिखता है? एक Kafka-से-Kafka पाइपलाइन Kafka लेनदेन (transactions) के साथ आउटपुट और इनपुट ऑफ़सेट को परमाणु रूप से (atomically) कमिट कर सकती है। एक डेटाबेस, ऑब्जेक्ट स्टोर या बाहरी API को आमतौर पर एट-लीस्ट-वन्स (at-least-once) प्रोसेसिंग के साथ-साथ एक व्यावसायिक इडेम्पोटेंसी की या topic-partition-offset की आवश्यकता होती है।
  • एक ऑर्डर कितने समय तक सक्रिय रह सकता है? अल्पकालिक ऑर्डर पूरा होने तक पुराने रूट पर रह सकते हैं जबकि नए ऑर्डर नए रूट का उपयोग करते हैं। लंबे समय तक चलने वाले ऑर्डरों के लिए एक स्पष्ट बैरियर, अनुक्रम (sequence) या रूटिंग स्थिति की आवश्यकता होती है।
  • क्या घटना प्रतिक्रिया ट्रैफ़िक को थ्रॉटल या डिग्रेड कर सकती है? एक टेनेंट कोटा, विलंबित गैर-महत्वपूर्ण इवेंट, या समेकित (coalesced) स्थिति अपडेट कोड और पार्टीशन माइग्रेशन की तुलना में इनपुट को तेजी से कम कर सकते हैं।
  • क्या कंज्यूमर्स max.poll.interval.ms से अधिक समय ले रहे हैं? यदि भारी काम पोल थ्रेड को ब्लॉक करता है, तो रिबैलेंस लैग को बढ़ा देते हैं। केवल टाइमआउट बढ़ाने से पहले पोलिंग को प्रोसेसिंग से अलग करें और इन-फ़्लाइट कार्य को सीमित करें।

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

"मैं कुल लैग को प्रति-पार्टिशन प्रोड्यूस दर, कंज्यूम दर और लैग स्लोप में विभाजित करूँगा, फिर उन्हें की (key) आवृत्ति, बाइट्स, हैंडलर समय, रिबैलेंस लॉग और लीडर ब्रोकर के मेट्रिक्स के साथ संरेखित करूँगा। हॉट टेनेंट प्रति सेकंड 54,000 रिकॉर्ड बनाता है जबकि एक कंज्यूमर 8,000 को संभालता है, इसलिए बैकलॉग लगभग 46,000 प्रति सेकंड बढ़ता है; सामान्य कंज्यूमर्स को जोड़ने से उस पार्टीशन में तेजी नहीं आ सकती है। आज मैं हॉट टेनेंट को थ्रॉटल या डिग्रेड करूँगा, पार्टीशन को एक समर्पित इंस्टेंस पर अलग करूँगा, और विभिन्न ऑर्डरों को समवर्ती रूप से प्रोसेस करते हुए प्रत्येक ऑर्डर को सीरियलाइज़ करके सही ऑर्डरिंग सीमा का लाभ उठाऊंगा। मैं केवल उच्चतम निरंतर पूर्ण ऑफ़सेट को कमिट करूँगा और इडेम्पोटेंट डाउनस्ट्रीम राइट्स का उपयोग करूँगा। लंबी अवधि में, नए ऑर्डर order_id द्वारा कीड किए गए एक नए टॉपिक पर चले जाएंगे, जबकि मौजूदा ऑर्डर समाप्त होने तक पुराने रूट पर बने रहेंगे। मैं विषम लोड, कंज्यूमर क्रैश और रिबैलेंस के तहत पार्टीशन-स्तरीय लैग स्लोप, एंड-टू-एंड p99, डुप्लिकेट और ऑर्डर उल्लंघनों को सत्यापित करूँगा।"

चरण-दर-चरण गहन उत्तर

चरण 1: साबित करें कि किस परत ने हॉटस्पॉट बनाया

एक पीक-टाइम विंडो का उपयोग करें और संरेखित करें:

  1. प्रति सेकंड प्रोड्यूस रिकॉर्ड, प्रति सेकंड बाइट्स, और पार्टीशन द्वारा हाई-वाटरमार्क वृद्धि;
  2. प्रति सेकंड कंज्यूम रिकॉर्ड, कमिटेड ऑफ़सेट, और पार्टीशन द्वारा लैग स्लोप;
  3. की (key) आवृत्ति, बाइट्स, और अनुमानित प्रोसेसिंग-लागत वितरण;
  4. हॉट पार्टीशन के लीडर ब्रोकर पर CPU, नेटवर्क, डिस्क प्रतीक्षा, और रिक्वेस्ट लेटेंसी;
  5. असाइन किए गए कंज्यूमर के लिए पोल अंतराल, बैच आकार, हैंडलर लेटेंसी, कचरा संग्रहण, त्रुटियां और रिबैलेंस लॉग;
  6. डाउनस्ट्रीम डेटाबेस, स्टोरेज सिस्टम या API में पार्टीशन-सहसंबद्ध लेटेंसी या थ्रॉटलिंग।

निर्णय का नियम संख्याओं से आता है। हॉट पार्टीशन को प्रति सेकंड लगभग 54,000 रिकॉर्ड प्राप्त होते हैं। अन्य 23 पार्टीशन शेष 66,000 को साझा करते हैं, यदि वह शेष यथोचित रूप से एक समान है तो औसतन लगभग 2,870 रिकॉर्ड प्रति सेकंड। 8,000-रिकॉर्ड वाले कंज्यूमर के पास एक साधारण पार्टीशन पर अतिरिक्त क्षमता होती है लेकिन वह हॉट इनपुट से मेल नहीं खा सकता है। यह एक बढ़ते पार्टीशन और अन्यत्र निष्क्रिय क्षमता की पूरी तरह से व्याख्या करता है। यदि हॉट-पार्टिशन इनपुट सामान्य है जबकि खपत डाउनस्ट्रीम लेटेंसी के साथ घटती है, तो की (key) डिज़ाइन अभी तक सिद्ध मूल कारण नहीं है।

ब्रोकर प्लेसमेंट की भी जांच करें। एक पार्टीशन जिसका लीडर ओवरलोड ब्रोकर पर बैठता है, वह धीमी गति से उत्पादन और फ़ेचिंग का शिकार हो सकता है। लीडरशिप को स्थानांतरित करना या प्रतिकृतियों (replicas) को रिबैलेंस करना उस प्लेसमेंट बाधा को हटा सकता है, लेकिन यह इस तथ्य को नहीं बदलता है कि एक tenant_id अभी भी एक पार्टीशन पर मैप करता है।

चरण 2: बैकलॉग को खाली करने का प्रयास करने से पहले उसके स्लोप को कम करें

पहला घटना उद्देश्य है:

text
hot-partition input rate ≤ hot-partition safe processing rate

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

कंज्यूमर पक्ष पर, एक नियंत्रित रखरखाव विंडो मौजूदा समूह को रोक सकती है और इसे एक विशेष, गैर-अतिव्यापी स्पष्ट असाइनमेंट के साथ बदल सकती है: एक पर्याप्त रूप से प्रावधानित इंस्टेंस केवल हॉट पार्टीशन का मालिक है, और शेष इंस्टेंस अन्य पार्टीशन के मालिक हैं। दूसरा सामान्य कंज्यूमर समूह शुरू करना कोई मदद नहीं है; यह स्वतंत्र रूप से पूर्ण टॉपिक का उपभोग करता है और व्यावसायिक प्रभावों की नकल (duplicate) करता है। आइसोलेशन हॉट पार्टीशन को उसी प्रोसेस पर सामान्य पार्टीशन को भूखा रखने से रोकता है, हालांकि यह मूल सीरियल हैंडलर को 8,000 रिकॉर्ड प्रति सेकंड से ऊपर नहीं उठाता है।

चूंकि व्यवसाय को केवल एक ऑर्डर के भीतर क्रम की आवश्यकता होती है, हॉट कंज्यूमर order_id द्वारा डिस्पैच कर सकता है: प्रति सक्रिय ऑर्डर एक सीरियल कतार और विभिन्न ऑर्डरों में एक सीमित वर्कर पूल। पूल में इन-फ़्लाइट सीमा होनी चाहिए। जब यह भर जाए, तो पार्टीशन को रोकें या वर्कर्स के लिए जारी की गई राशि को कम करें ताकि Kafka लैग अनबाउंडेड प्रोसेस मेमोरी न बन जाए। पोल लूप उत्तरदायी रहना चाहिए; अन्यथा max.poll.interval.ms से अधिक होने पर रिबैलेंस ट्रिगर होता है और एक और पॉज़ जुड़ जाता है।

चरण 3: एक निरंतर पूर्णता वॉटरमार्क कमिट करें

इंट्रा-पार्टिशन समवर्तीता पूरा होने के क्रम को बदल देती है, लेकिन इसे कमिट क्रम को नहीं बदलना चाहिए। प्रत्येक पार्टीशन के लिए इस स्थिति को बनाए रखें:

text
nextCommitOffset = smallest unfinished offset
completed = offsets that finished but still have a gap before them

onComplete(offset):
  add offset to completed
  while completed contains nextCommitOffset:
    remove nextCommitOffset from completed
    increment nextCommitOffset
  commit nextCommitOffset

Kafka पढ़ने के लिए अगली स्थिति कमिट करता है। इसलिए, 104, 105, और 106 सभी समाप्त होने के बाद ही कमिटेड स्थिति 107 तक आगे बढ़ सकती है। यदि 104 के पुनः प्रयास करने के दौरान 105 समाप्त हो जाता है, तो कमिट बिंदु 104 पर ही रहता है। अगले कमिट से पहले एक क्रैश कुछ पूर्ण रिकॉर्ड को फिर से चलाता है, इसलिए डेटाबेस राइट्स को एक इडेम्पोटेंट अपसर्ट के लिए event_id या किसी अन्य व्यावसायिक-अद्वितीय की (key) का उपयोग करना चाहिए। यदि कोई व्यावसायिक की मौजूद नहीं है, तो topic-partition-offset स्रोत रिकॉर्ड की पहचान कर सकता है।

ऑफ़सेट कमिट और बाहरी साइड इफ़ेक्ट को स्वाभाविक रूप से परमाणु (atomic) के रूप में वर्णित न करें। एक Kafka-से-Kafka एप्लिकेशन एक Kafka लेनदेन में आउटपुट रिकॉर्ड और कंज्यूम्ड ऑफ़सेट डाल सकता है। एक बाहरी डेटाबेस के साथ, अधिक सामान्य अनुबंध एट-लीस्ट-वन्स उपभोग प्लस एक इडेम्पोटेंट राइट, या डीडुप्लिकेशन रिकॉर्ड और व्यावसायिक म्यूटेशन दोनों से युक्त एक डेटाबेस लेनदेन है।

चरण 4: पार्टीशन के दायरे को वास्तविक ऑर्डरिंग डोमेन से मिलाएं

24 पार्टीशन की औसत क्षमता अनिवार्य रूप से अपर्याप्त नहीं है। यदि 8,000 रिकॉर्ड प्रति सेकंड एक मापी गई सुरक्षित अधिकतम सीमा है और नियोजित उपयोग 70% पर सीमित है, तो प्रति पार्टीशन नियोजित क्षमता 5,600 है:

text
120,000 ÷ 5,600 ≈ 21.4

समान वितरण के साथ, 24 पार्टीशन मामूली हेडरूम के साथ अनुमानित पीक को कवर करते हैं। विफलता 45% ट्रैफ़िक को एक कम-कार्डिनैलिटी की (key) के पीछे रखने से आती है, न कि कुल पार्टीशन संख्या से। एक टिकाऊ की (key) को सबसे छोटे आवश्यक ऑर्डरिंग डोमेन का प्रतिनिधित्व करना चाहिए, पर्याप्त कार्डिनैलिटी होनी चाहिए, और पीक पर पूर्वानुमेय रूप से वितरित रहना चाहिए। यहाँ, order_id स्वाभाविक विकल्प है। (tenant_id, order_id) का एक स्थिर एन्कोडिंग भी संभव है यदि टेनेंट लोकैलिटी का वास्तविक परिचालन मूल्य है।

नियंत्रित साल्टिंग (salting) केवल तभी मान्य है जब मूल की (key) के भीतर रिकॉर्ड को पुन: व्यवस्थित किया जा सकता है या डाउनस्ट्रीम चरण क्रम को पुनर्स्थापित कर सकता है। किसी एक ऑर्डर को बेतरतीब ढंग से सॉल्ट करना इस प्रॉम्प्ट के इनवेरिएंट का उल्लंघन करता है क्योंकि वह ऑर्डर कई पार्टीशन से समवर्ती रूप से आ सकता है। यदि केवल एक ऑर्डर एकल-पार्टिशन क्षमता से अधिक है, तो अधिक की कार्डिनैलिटी मदद नहीं करती है; उस ऑर्डर के सीरियल पाथ को अनुकूलित या थ्रॉटल करें, या व्यावसायिक प्रोटोकॉल को स्पष्ट रूप से स्वतंत्र अनुक्रमों में फिर से डिज़ाइन करें।

चरण 5: संस्करणित रूटिंग के साथ माइग्रेट करें

मौजूदा टॉपिक में पार्टीशन जोड़ना और कीज़ को तुरंत बदलना दो जोखिम पैदा करता है। पार्टीशन की संख्या बदलने पर डिफ़ॉल्ट हैश मैपिंग मौजूदा कीज़ को स्थानांतरित कर सकती है। कटओवर से पहले और बाद में एक ऑर्डर के इवेंट अलग-अलग पार्टीशन में भी पहुंच सकते हैं, जबकि Kafka क्रॉस-पार्टिशन ऑर्डर की गारंटी नहीं देता है।

एक सुरक्षित डिज़ाइन order_id द्वारा पार्टीशन किया गया एक नया टॉपिक बनाता है और प्रोड्यूसर रूटिंग को वर्शन करता है:

  • कटओवर के बाद बनाए गए ऑर्डर नए टॉपिक का उपयोग करते हैं;
  • जो ऑर्डर पहले से मौजूद थे, वे समाप्त होने तक पुराने टॉपिक और लीगेसी की (key) पर बने रहेंगे;
  • प्रत्येक प्रोड्यूसर अपने स्थानीय क्लॉक की तुलना कटओवर समय से करने के बजाय समान ऑर्डर-रूटिंग स्थिति का उपयोग करता है;
  • कंज्यूमर्स दोनों रूट पढ़ते हैं, लेकिन एक ऑर्डर किसी भी समय केवल एक सक्रिय रूट से संबंधित होता है;
  • पुराने ऑर्डरों के खाली होने और अवधारण (retention) आवश्यकताओं के पूरा होने के बाद पुराने टॉपिक को रिटायर कर दिया जाता है।

यदि ऑर्डर स्वाभाविक रूप से समाप्त नहीं होते हैं, तो प्रति-ऑर्डर माइग्रेशन बैरियर बनाएं: उस ऑर्डर के लिए नए इवेंट को रोकें, जब तक कि पुराना रूट रिकॉर्ड किए गए अंतिम अनुक्रम या ऑफ़सेट तक न पहुंच जाए, रूटिंग संस्करण बदलें, और फिर से शुरू करें। एक शून्य-पॉज़ विकल्प मोनोटोनिक अनुक्रम संख्याएं ले जा सकता है और डाउनस्ट्रीम दोनों रूट को मर्ज कर सकता है, लेकिन यह बफरिंग, टाइमआउट और गैप रिकवरी का परिचय देता है। यह केवल तभी उचित है जब व्यावसायिक आवश्यकता उस जटिलता का भुगतान करती है।

चरण 6: विषम ट्रैफ़िक और विफलताओं के साथ सत्यापन करें

समग्र थ्रूपुट पर्याप्त प्रमाण नहीं है। कम से कम परीक्षण करें:

  • एक वितरण जिसमें एक टेनेंट 45% ट्रैफ़िक उत्पन्न करता है और ऑर्डर वितरण वास्तविक पीक जैसा दिखता है;
  • पार्टीशन द्वारा इनपुट, खपत, लैग स्लोप और अधिकतम लैग;
  • एंड-टू-एंड p50, p95, और p99 प्लस अनुमानित ड्रेन समय;
  • इन-फ़्लाइट वर्कर संख्या, सबसे पुराने कार्य की आयु, पुनः प्रयास और डेड लेटर्स;
  • डुप्लिकेट प्रभाव, प्रति-ऑर्डर क्रम उल्लंघन, और इडेम्पोटेंसी टकराव;
  • जब पूर्ण ऑफ़सेट में गैप हो तब कंज्यूमर क्रैश;
  • क्या लंबी प्रोसेसिंग रिबैलेंस को ट्रिगर करती है और रिकवरी में कितना समय लगता है;
  • क्या कोई ऑर्डर पुराने/नए टॉपिक सीमा पर केवल एक रूट पर दिखाई देता है।

पास होने की स्थिति में निरंतर व्यवहार शामिल है: स्थिर पीक पर हॉट-पार्टिशन लैग स्लोप अब सकारात्मक नहीं है; क्रैश कार्य को फिर से चला सकता है लेकिन व्यावसायिक प्रभाव नहीं खो सकता है; कोई भी ऑर्डर क्रम से बाहर नहीं देखा गया है; नए ऑर्डर पार्टीशनों में वितरित होते हैं; और पुराने ऑर्डर शेड्यूल पर खाली हो जाते हैं। यदि एक बड़ा ऑर्डर बार-बार हॉटस्पॉट बनाते हुए समग्र थ्रूपुट बढ़ता है, तो ऑर्डरिंग डोमेन या बिजनेस एडमिशन की समस्या अनसुलझी रहती है।

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

"मैं कंज्यूमर्स जोड़कर शुरुआत नहीं करूँगा क्योंकि प्रॉम्प्ट हमें पहले ही बताता है कि एक पार्टीशन पीछे है जबकि अन्यत्र कंज्यूमर्स निष्क्रिय हैं। मैं पहले प्रति-पार्टिशन रिकॉर्ड प्रति सेकंड, बाइट्स प्रति सेकंड, लैग स्लोप और की (key) आवृत्ति का उपयोग करके तिरछेपन को साबित करूँगा, जबकि लीडर ब्रोकर, रिबैलेंस और डाउनस्ट्रीम लेटेंसी को खारिज करूँगा।

हॉट टेनेंट प्रति सेकंड 54,000 रिकॉर्ड बनाता है। एक एकल पार्टीशन कंज्यूमर 8,000 प्रोसेस करता है, इसलिए लैग प्रति सेकंड लगभग 46,000 रिकॉर्ड या 10 मिनट में 27.6 मिलियन बढ़ता है। यदि अन्य 55% मोटे तौर पर 23 पार्टीशनों पर फैला हुआ है, तो प्रत्येक औसतन लगभग 2,870 रिकॉर्ड प्रति सेकंड है। यह हॉट कंज्यूमर के घाटे और अन्यत्र अतिरिक्त क्षमता दोनों की व्याख्या करता है। एक पारंपरिक कंज्यूमर समूह में, एक पार्टीशन एक समय में एक कंज्यूमर का होता है, इसलिए अधिक सामान्य इंस्टेंस इसे गति नहीं देते हैं।

आज मैं सबसे पहले एक प्रलेखित टेनेंट कोटा के साथ इनपुट स्लोप को कम करूँगा और देरी-सहिष्णु या समेकन-योग्य इवेंट्स को डिग्रेडेड पाथ पर ले जाऊँगा। एक नियंत्रित विंडो में, मैं हॉट पार्टीशन को एक समर्पित इंस्टेंस पर अलग करूँगा ताकि यह सामान्य पार्टीशन को भूखा न रखे। चूंकि केवल प्रति-ऑर्डर ऑर्डरिंग मायने रखती है, इसलिए मैं एक सीमित पूल में order_id द्वारा डिस्पैच करूँगा: एक ऑर्डर के भीतर सीरियल, ऑर्डरों के बीच समवर्ती। मैं उस कार्य को कमिट नहीं करूँगा जो सबसे तेज़ समाप्त होता है। मैं उच्चतम निरंतर पूर्ण ऑफ़सेट को ट्रैक करूँगा और किसी भी गैप पर रुकूँगा। एक क्रैश पूर्ण लेकिन अनकमिटेड रिकॉर्ड को फिर से चला सकता है, इसलिए सिंक इडेम्पोटेंसी के लिए event_id या एक व्यावसायिक-अद्वितीय की (key) का उपयोग करता है।

दीर्घकालिक रूप से, पार्टीशन जोड़ना पूरा समाधान नहीं है। 70% नियोजित उपयोग पर, प्रत्येक 8,000-रिकॉर्ड पार्टीशन प्रति सेकंड लगभग 5,600 रिकॉर्ड का योगदान देता है, इसलिए 24 समान रूप से लोड किए गए पार्टीशन मान ली गई 120,000-रिकॉर्ड पीक को कवर कर सकते हैं। समस्या यह है कि tenant_id 45% को एक पार्टीशन पर पिन करता है। मैं order_id द्वारा कीड एक नया टॉपिक बनाऊंगा। कटओवर के बाद नए ऑर्डर इसका उपयोग करते हैं, जबकि सक्रिय पुराने ऑर्डर पूरा होने तक पुराने रूट पर बने रहते हैं। साझा रूटिंग स्थिति यह सुनिश्चित करती है कि एक ऑर्डर कभी भी दोनों टॉपिक में न फैले।

रोलआउट से पहले, मैं उसी 45% टेनेंट स्क्यू को फिर से चलाऊँगा और प्रति-पार्टिशन थ्रूपुट, लैग स्लोप, एंड-टू-एंड p99, डुप्लिकेट और ऑर्डर उल्लंघनों का निरीक्षण करूँगा। जब ऑफ़सेट क्रम से बाहर पूरे होते हैं, तो मैं कंज्यूमर को क्रैश करूँगा और सत्यापित करूँगा कि पुनरारंभ केवल इडेम्पोटेंट रीप्ले का कारण बनता है, रिबैलेंस ट्रिगर करूँगा और रिकवरी को मापूंगा, और सत्यापित करूँगा कि रूट सीमा पर प्रत्येक ऑर्डर केवल एक टॉपिक पर दिखाई देता है। यदि एक ऑर्डर स्वयं एकल-पार्टिशन क्षमता से अधिक है, तो मैं कठोर सीमा का उल्लेख करूँगा: न तो अधिक कंज्यूमर्स और न ही अधिक पार्टीशन उस ऑर्डर के सीक्वेंसिंग मॉडल को अनुकूलित, थ्रॉटल या बदले बिना इसे हल करते हैं।"

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

  • केवल टॉपिक-व्यापी लैग को देखना → एक औसत एक पार्टीशन के इनपुट और कंज्यूम स्लोप को छुपाता है → प्रति पार्टीशन रिकॉर्ड, बाइट्स, लैग और की (key) वितरण का ग्राफ़ बनाएं।
  • जब भी लैग बढ़े तो कंज्यूमर्स जोड़ना → पारंपरिक समूह में एक पार्टीशन का एक समय में एक कंज्यूमर मालिक होता है → पहले पार्टीशन संख्या, असाइनमेंट और प्रति-पार्टिशन क्षमता की तुलना करें।
  • तुरंत पार्टीशन जोड़ना → मौजूदा बैकलॉग स्वचालित रूप से पुनर्वितरित नहीं होता है, और डिफ़ॉल्ट की (key) मैपिंग बदल सकती है → पहले बैकलॉग वृद्धि को रोकें, फिर एक संस्करणित टॉपिक और माइग्रेशन सीमा का उपयोग करें।
  • हॉट की (key) को बेतरतीब ढंग से सॉल्ट करना → एक ऑर्डर पार्टीशनों को पार कर सकता है और क्रम से बाहर आ सकता है → केवल तभी सॉल्ट करें जब पुन: क्रमित करना स्वीकार्य हो; इस प्रॉम्प्ट के वास्तविक ऑर्डरिंग दायरे के लिए order_id का उपयोग करें।
  • जैसे ही उसका वर्कर समाप्त हो, ऑफ़सेट को कमिट करना → क्रैश के बाद निचले अधूरे ऑफ़सेट को छोड़ा जा सकता है → केवल निरंतर पूर्णता वॉटरमार्क कमिट करें।
  • ऑफ़सेट कमिट को एग्जैक्ट्ली-वन्स (exactly-once) प्रोसेसिंग कहना → एक बाहरी डेटाबेस म्यूटेशन और Kafka ऑफ़सेट आमतौर पर एक लेनदेन नहीं होते हैं → एट-लीस्ट-वन्स, इडेम्पोटेंसी, और लेनदेन सीमाओं का उल्लेख करें।
  • मदद के लिए दूसरा कंज्यूमर समूह शुरू करना → दूसरा समूह अपनी पूरी प्रति पढ़ता है और प्रभावों की नकल करता है → विशेष स्पष्ट असाइनमेंट या नियंत्रित प्रोसेसिंग रीडिज़ाइन का उपयोग करें।
  • केवल एकसमान ट्रैफ़िक का परीक्षण करना → एक पासिंग औसत यह साबित नहीं करता है कि हॉट की (key) चली गई है → एक यथार्थवादी तिरछापन फिर से चलाएं और केवल औसत ही नहीं, अधिकतम पार्टीशन का निरीक्षण करें।
  • एकल-एंटिटी सीमा की उपेक्षा करना → एक कड़ाई से क्रमित एंटिटी को मुफ्त में पार्टीशनों में समानांतर नहीं किया जा सकता है → क्रम, थ्रॉटलिंग और सीरियल प्रोसेसिंग के बीच कठिन ट्रेड-ऑफ़ का उल्लेख करें।

अनुवर्ती प्रश्न और प्रतिक्रियाएं

अनुवर्ती 1: पार्टीशन की संख्या तुरंत 24 से बढ़ाकर 48 क्यों न की जाए?

नए पार्टीशन भविष्य के समानांतर स्लॉट बनाते हैं, लेकिन वे पुराने पार्टीशन में पहले से संग्रहीत बैकलॉग को विभाजित नहीं करते हैं, न ही वे एक tenant_id को कई पार्टीशनों में मैप करते हैं। डिफ़ॉल्ट की (key) हैशिंग के साथ, पार्टीशन संख्या बदलने से मौजूदा कीज़ भी रीमैप हो सकती हैं और कटओवर के दौरान एक ही ऑर्डर को विभिन्न पार्टीशनों पर रखा जा सकता है। पहले ऑर्डरिंग डोमेन और माइग्रेशन को ठीक करें, फिर यह तय करने के लिए विषम लोड परीक्षण का उपयोग करें कि क्या अधिक कुल पार्टीशन आवश्यक हैं।

अनुवर्ती 2: इंट्रा-पार्टिशन समवर्तीता जोड़ने के बाद आप अनबाउंडेड मेमोरी को कैसे रोकते हैं?

प्रति पार्टीशन अधिकतम इन-फ़्लाइट रिकॉर्ड संख्या और अधिकतम अनकमिटेड ऑफ़सेट विंडो सेट करें। जब कोई भी सीमा पूरी हो जाए, तो उस पार्टीशन को रोकें या भारी प्रोसेसिंग से अलग पोलिंग रखते हुए वर्कर पूल के लिए कम रिकॉर्ड जारी करें। सबसे पुराने अधूरे ऑफ़सेट की आयु को ट्रैक करें। यदि एक ऑर्डर ब्लॉक रहता है, तो प्रत्येक बाद के ऑफ़सेट को अनिश्चित काल तक मेमोरी का उपभोग करने की अनुमति देने के बजाय हस्तक्षेप के लिए इसे अलग करें, पुनः प्रयास करें या रूट करें।

अनुवर्ती 3: क्या होगा यदि डाउनस्ट्रीम डेटाबेस इडेम्पोटेंट अपसर्ट का समर्थन नहीं करता है?

एक डीडुप्लिकेशन रिकॉर्ड डालें और एक डेटाबेस लेनदेन में व्यावसायिक म्यूटेशन निष्पादित करें। एक अद्वितीय की (key) के रूप में एक व्यावसायिक event_id या topic-partition-offset का उपयोग करें; विशिष्टता टकराव का अर्थ है कि इवेंट पहले ही लागू हो चुका है। एक गैर-लेनदेन बाहरी API के लिए, इसकी इडेम्पोटेंसी की, एक आउटबॉक्स, या एक क्वेरिएबल ऑपरेशन स्थिति का उपयोग करें। यदि कोई इडेम्पोटेंसी सीमा नहीं बनाई जा सकती है, तो आप यह वादा नहीं कर सकते कि क्रैश रीप्ले का कोई डुप्लिकेट प्रभाव नहीं होगा।

अनुवर्ती 4: यदि व्यवसाय को बाद में सख्त टेनेंट-व्यापी क्रम की आवश्यकता हो तो क्या बदलता है?

तब tenant_id अविभाज्य ऑर्डरिंग डोमेन है, और क्रॉस-ऑर्डर समवर्तीता अब मान्य नहीं है। यदि एक टेनेंट एकल-पार्टिशन क्षमता से अधिक है, तो उस सीरियल पाथ को अनुकूलित करें, टेनेंट को थ्रॉटल करें, या इस बात पर फिर से बातचीत करें कि कौन से इवेंट अनुक्रम स्वतंत्र हैं। डाउनस्ट्रीम मर्जर के साथ विश्व स्तर पर अनुक्रमित शार्डेड स्ट्रीम संभव है, लेकिन यह ऑर्डरिंग प्रतीक्षा, गैप रिकवरी और उपलब्धता लागत को कंज्यूमर में ले जाता है; यह मुफ़्त स्केल नहीं है।

अनुवर्ती 5: 27.6-मिलियन-रिकॉर्ड बैकलॉग को खाली करने में कितना समय लगेगा?

इनपुट पहले प्रोसेसिंग क्षमता से नीचे गिरना चाहिए। यदि थ्रॉटलिंग हॉट-पार्टिशन इनपुट को 3,000 रिकॉर्ड प्रति सेकंड तक कम कर देती है और ऑप्टिमाइज़ेशन प्रोसेसिंग क्षमता को 12,000 तक बढ़ा देता है, तो शुद्ध ड्रेन दर 9,000 है:

text
27,600,000 ÷ 9,000 ≈ 3,067 seconds ≈ 51 minutes

यह अभी भी एक स्थिर-दर अनुमान है। एक वास्तविक रिकवरी योजना में पुनः प्रयास, डाउनस्ट्रीम थ्रॉटलिंग, रिकॉर्ड-आकार भिन्नता और सुरक्षा मार्जिन शामिल होता है, फिर देखे गए लैग स्लोप से अनुमान को लगातार संशोधित किया जाता है।

अनुवर्ती 6: आप कैसे साबित करेंगे कि कोई भी ऑर्डर पुराने और नए टॉपिक को पार नहीं करता है?

प्रत्येक ऑर्डर के लिए एक partitioning_version संग्रहीत करें। प्रत्येक प्रोड्यूसर समान संस्करणित रूटिंग रिकॉर्ड को पढ़ता या कैश करता है, और माइग्रेशन बैरियर के सफल होने के बाद ही संस्करण बदलता है। कंज्यूमर्स किसी ऑर्डर के लिए देखे गए पहले टॉपिक और संस्करण को रिकॉर्ड करते हैं, दोनों सक्रिय रूट पर दिखाई देने वाले ऑर्डर पर अलर्ट करते हैं, और उस ऑर्डर के लिए स्वचालित प्रगति को रोकते हैं। लोड परीक्षणों और कैनरी रोलआउट के दौरान, केवल समान कुल रिकॉर्ड संख्या की जाँच करने के बजाय प्रोड्यूसर लॉग, दोनों टॉपिक पर ऑफ़सेट और डाउनस्ट्रीम ऑर्डर अनुक्रमों का मिलान (reconcile) करें।

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

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