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

स्ट्रीमिंग एग्रीगेशन में आप लेट और आउट-ऑफ-ऑर्डर इवेंट्स को कैसे संभालते हैं?

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

प्रश्न

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

1. प्रश्न

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

2. बाधाएं और स्पष्टीकरण

  • पुष्टि करें कि विंडो इवेंट टाइम, राइट टाइम, या प्रोसेसिंग टाइम का उपयोग करती हैं; यह समस्या इवेंट टाइम का उपयोग करती है।
  • अधिकतम अव्यवस्था की सीमा तय करें, जैसे कि 30 सेकंड, और परिभाषित करें कि इसके बाद क्या होता है।
  • तय करें कि परिणाम अपेंड-ओनली इवेंट्स हैं या अपडेट करने योग्य स्नैपशॉट; डाउनस्ट्रीम कंज्यूमर्स को संशोधनों (revisions) की पहचान करने में सक्षम होना चाहिए।
  • निष्क्रिय (idle) पार्टिशन्स, अमान्य टाइमस्टैम्प, डुप्लिकेट इवेंट्स और रीप्ले जॉब्स पर चर्चा करें ताकि संदेश-रहित इनपुट को स्ट्रीम का अंत न समझ लिया जाए।

3. मुख्य अवधारणाएं

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

4. संदर्भ प्रवाह

text
onRecord(event):
  ts = extractEventTimestamp(event)
  key = canonicalKey(event.customerId)
  updateWatermarkGenerator(ts)

  window = floorToFiveMinutes(ts)
  if ts <= currentWatermark + allowedLateness:
    state[window, key] = aggregate(state[window, key], event)
    emitUpsert(window, key, state[window, key], revision + 1)
  else:
    routeToLateData(window, key, event)

onWatermark(wm):
  finalizeWindowsBefore(wm)
  expireStateAfterRetention()

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

5. सटीकता और लेटेंसी ट्रेड-ऑफ

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

6. सत्यापन और ऑब्जर्वेबिलिटी

  • क्रमित, आउट-ऑफ-ऑर्डर, 30-सेकंड-से-कम लेट, और समय सीमा से परे के इवेंट्स उत्पन्न करें; ऑफ़लाइन बेसलाइन के साथ प्रत्येक विंडो की तुलना करें।
  • निष्क्रिय पार्टिशन्स, क्लॉक जंप, डुप्लिकेट इवेंट्स और टास्क रीस्टार्ट इंजेक्ट करें; सत्यापित करें कि वॉटरमार्क रुकते नहीं हैं या पीछे नहीं जाते हैं।
  • वर्तमान वॉटरमार्क, प्रोसेसिंग-टाइम बनाम इवेंट-टाइम लैग, विंडो-स्टेट साइज, लेट-इवेंट रेट, साइड-आउटपुट वॉल्यूम और सुधार संख्या रिकॉर्ड करें।
  • प्रत्येक अपडेट में विंडो कुंजी, संशोधन और इनपुट इवेंट ID शामिल करें; उसी लॉग को रीप्ले करें और अंतिम स्नैपशॉट की तुलना करें।

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

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

8. साक्षात्कार स्कोरिंग बिंदु

समय की तीन अवधारणाओं में अंतर करना

उम्मीदवार को इवेंट, इन्जेशन और प्रोसेसिंग टाइम को परिभाषित करना चाहिए और समझाना चाहिए कि विंडो का चयन सटीकता और लेटेंसी को कैसे बदलता है।

वॉटरमार्क उत्पन्न और मर्ज करना

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

लेट-डेटा और सुधार पथ डिज़ाइन करना

उम्मीदवार को समय सीमा के भीतर और समय सीमा के बाद के दोनों डेटा के लिए अलाउड-लेटनेस समय सीमा, साइड आउटपुट, संशोधन और डाउनस्ट्रीम इडेम्पोटेंट मर्ज को परिभाषित करना चाहिए।

रीप्ले के साथ परिणाम का सत्यापन करना

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

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

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