समस्या और दायरा
मोबाइल क्लाइंट कम से कम event_id, event_time, user_id, और पहले से ही एक मुद्रा में परिवर्तित amount के साथ खरीदारी इवेंट रिपोर्ट करते हैं। नेटवर्क पुनः प्रयास (retries) डुप्लिकेट बना सकते हैं, ऑफ़लाइन डिवाइस कई घंटे बाद अपलोड कर सकते हैं, और विभिन्न पार्टीशन वैश्विक क्रम सुरक्षित नहीं रखते हैं। 20,000 इवेंट प्रति सेकंड के पीक इनपुट के साथ, event_time द्वारा प्रत्येक UTC क्लॉक घंटे के लिए राजस्व की गणना करें।
व्यवसाय को एक मिनट के भीतर वर्तमान घंटे के लिए एक अनुमान, वॉटरमार्क द्वारा विंडो के अंत को पार करने के बाद एक ऑन-टाइम परिणाम, और उस विंडो के समाप्त होने के 24 घंटे बाद तक स्वचालित सुधार की आवश्यकता है। 24 घंटे के बाद आने वाला डेटा चुपचाप गायब नहीं होना चाहिए; यह ऑडिट और ऑफ़लाइन समाधान (reconciliation) के लिए जाता है। थ्रूपुट, समयबद्धता और सुधार अवधि समस्या के इनपुट हैं, न कि स्ट्रीम-प्रोसेसिंग इंजन के बारे में प्रदर्शन के दावे।
मान लें कि निर्माता एक स्थिर event_id असाइन करता है, और उस ID के साथ वैध पुनः प्रयास समान व्यावसायिक सामग्री ले जाते हैं। रॉ इवेंट्स रिप्ले करने योग्य स्टोरेज में रहते हैं। इस दायरे में समय शब्दार्थ (time semantics), वॉटरमार्क, विंडो ट्रिगर, डिडुप्लिकेशन, परिणाम अपडेट, स्थिति क्षमता (state capacity), रिकवरी और मिलान शामिल हैं। ब्रोकर चयन और मुद्रा रूपांतरण दायरे से बाहर हैं। यह एक डेटा-इंजीनियरिंग प्रश्न है: इसका लक्ष्य एक सत्यापन योग्य डेटा अनुबंध है जो पूर्णता, विलंबता और स्थिति-लागत के ट्रेडऑफ़ को स्पष्ट करता है।
साक्षात्कारकर्ता क्या मूल्यांकन करते हैं
पहला संकेत यह है कि क्या उम्मीदवार तीन घड़ियों को अलग करता है। event_time किसी इवेंट वाले व्यावसायिक घंटे को निर्धारित करता है। processing_time बताता है कि इंजन इसे कब देखता है और एक मिनट के प्रारंभिक रिफ्रेश को चला सकता है। एक वॉटरमार्क इवेंट-टाइम प्रगति का इंजन का अनुमान है। यदि प्रोसेसिंग समय विंडो निर्धारित करता है, तो एक ही इतिहास को अलग-अलग समय पर फिर से चलाने से अलग-अलग प्रति घंटा परिणाम मिल सकते हैं।
दूसरा संकेत वॉटरमार्क को एक पूर्ण वादे के बजाय प्रगति अनुमान के रूप में मानना है। इसे तेजी से आगे बढ़ाने से समय पर आउटपुट मिलता है लेकिन अधिक रिकॉर्ड देर से आने वाले (late) के रूप में वर्गीकृत होते हैं। इसे रोके रखने से पूर्णता में सुधार होता है लेकिन विंडो बंद होने में देरी होती है और अधिक स्थिति बरकरार रहती है। एक मजबूत उत्तर एक निश्चित पांच मिनट या एक घंटे की देरी को याद करने के बजाय मापे गए आगमन-विलंब डेटा और सुधार SLO से नीति प्राप्त करता है।
तीसरा संकेत वॉटरमार्क, अनुमत विलंबता (allowed lateness) और डिडुप्लिकेशन प्रतिधारण (retention) को अलग करना है। वॉटरमार्क ऑन-टाइम निष्पादन को नियंत्रित करता है। अनुमत विलंबता नियंत्रित करती है कि विंडो स्थिति कितने समय तक सुधार योग्य रहती है। डिडुप्लिकेशन स्थिति उस अवधि तक फैली होनी चाहिए जिसमें किसी ईवेंट की दूसरी प्रति आ सकती है। मान संबंधित हो सकते हैं, लेकिन वे एक ही सेटिंग नहीं हैं।
अंत में, आउटपुट को अभिसरित (converge) होना चाहिए। देर से आने वाले इवेंट एक विंडो के लिए कई परिणाम उत्पन्न करते हैं। प्रत्येक स्नैपशॉट को जोड़ने से डाउनस्ट्रीम योग में विंडो बार-बार गिनी जाती है। परिणामों के लिए विंडो द्वारा अनुक्रमित (keyed) एक इडेम्पोटेंट अपसर्ट की आवश्यकता होती है, जिसमें एक संशोधन या अंतिमता मार्कर हो। एक चेकपॉइंट ऑपरेटर स्थिति की सुरक्षा करता है; यदि सिंक चेकपॉइंट के साथ कमिट नहीं कर सकता है, तो डिज़ाइन को अभी भी ट्रांजेक्शनल राइट्स या संस्करण-आधारित इडेम्पोटेंट प्रतिस्थापन की आवश्यकता होती है।
उत्तर देने से पहले स्पष्ट करने योग्य प्रश्न
- क्या व्यावसायिक मीट्रिक घटना (occurrence) या आगमन (arrival) पर आधारित है? यह समस्या खरीदारी
event_timeका उपयोग करती है। प्रोसेसिंग समय केवल तभी उपयुक्त होता है जब मीट्रिक विशेष रूप से "सिस्टम द्वारा अभी संसाधित किए गए अनुरोध" हो। - क्या एक मिनट का परिणाम एक अनुमान है या अंतिम संख्या? यह यहाँ एक अनुमान है, इसलिए प्रोसेसिंग-समय प्रारंभिक ट्रिगर मान्य है। यदि प्रत्येक प्रदर्शित परिणाम पूर्ण होना चाहिए, तो सिस्टम को अधिक प्रतीक्षा करनी होगी या बैच गणना का उपयोग करना होगा।
- क्या 24 घंटे के बाद रिपोर्ट बदल सकती है? 24 घंटे स्ट्रीमिंग जॉब की स्वचालित सुधार अवधि है। बाद का डेटा समाधान (reconciliation) में जाता है। प्रत्येक देर से आने वाले रिकॉर्ड को सही करने के लिए एक नियामक या निपटान आवश्यकता का अर्थ है कि स्थिति की सफ़ाई डेटा के अंतिम सत्य को परिभाषित नहीं कर सकती है।
- क्या डुप्लिकेट को ID द्वारा पहचाना जा सकता है? यहाँ एक स्थिर
event_idमौजूद है। उपयोगकर्ता, राशि और समय से अनुमान लगाने से वैध खरीदारी भी हट जाएगी और डुप्लिकेट भी छूट जाएंगे; पहले निर्माता अनुबंध को ठीक करें। - क्या होगा यदि एक ID में भिन्न सामग्री हो? पहले-लिखने-वाले-की-जीत (first-write-wins) या अंतिम-लिखने-वाले-की-जीत (last-write-wins) न चुनें। पेलोड-फ़िंगरप्रिंट विरोध रिकॉर्ड करें और इसे क्वारंटाइन करें क्योंकि निर्माता ने इडेम्पोटेंसी अनुबंध का उल्लंघन किया है।
- क्या सिंक स्नैपशॉट, डेल्टा या रिट्रैक्शन स्वीकार करता है? यह डिज़ाइन पूर्ण विंडो स्नैपशॉट उत्सर्जित करता है और
(window_start, dimensions)द्वारा अपसर्ट करता है। केवल-जोड़ने वाले (append-only) सिंक को एक संस्करणित चेंजलॉग की आवश्यकता होती है जिसकी रीड लेयर नवीनतम संशोधन का चयन करती है। - क्या इवेंट टाइमस्टैम्प पर भरोसा किया जा सकता है? भविष्य के टाइमस्टैम्प, रॉ-डेटा प्रतिधारण से पुराने, या सहमत समयक्षेत्र के लिए अमान्य टाइमस्टैम्प को क्वारंटाइन करें। अधिकतम ईवेंट समय में भाग लेने वाला एक खराब भविष्य का टाइमस्टैम्प वॉटरमार्क को बहुत तेजी से आगे बढ़ा सकता है।
30-सेकंड का उत्तर
"मैं event_time द्वारा UTC प्रति घंटा विंडो असाइन करूँगा, प्रत्येक सक्रिय स्रोत पार्टीशन के लिए एक वॉटरमार्क उत्पन्न करूँगा, और न्यूनतम सुरक्षित प्रगति पर आगे बढ़ूँगा। एक प्रोसेसिंग-टाइम ट्रिगर प्रत्येक मिनट अनुमान उत्सर्जित करता है; वॉटरमार्क ऑन-टाइम परिणाम उत्सर्जित करता है, जिसमें सुधार के लिए 24 घंटे तक स्थिति बनी रहती है। event_id और एक पेलोड फ़िंगरप्रिंट इवेंट्स को डिडुप्लिकेट करते हैं, जबकि सिंक विंडो और संशोधन द्वारा अपसर्ट करता है। अत्यधिक देर से आने वाला डेटा साइड आउटपुट और दैनिक समाधान में जाता है। पुनर्प्राप्ति एक रिप्ले करने योग्य स्रोत, चेकपॉइंट्स, और एक ट्रांजेक्शनल या इडेम्पोटेंट सिंक को जोड़ती है।"
चरण-दर-चरण गहन विश्लेषण
परिणाम अनुबंध से प्रारंभ करें। [10:00, 11:00) जैसी अर्ध-खुली (half-open) विंडो का उपयोग करें, और कुंजी में कम से कम window_start और रिपोर्टिंग आयाम शामिल करें। प्रत्येक आउटपुट अपने वर्तमान संस्करण में उस विंडो के लिए पूर्ण राजस्व स्नैपशॉट है:
HourlyRevenue {
window_start
window_end
revenue
revision
result_state // EARLY | ON_TIME | FINAL
}revision एक विंडो के लिए मोनोटोनिक रूप से बढ़ता है, और सिंक केवल एक बड़ा संशोधन स्वीकार करता है। EARLY एक मिनट का रिफ्रेश है और पूर्णता का दावा नहीं करता है। ON_TIME का अर्थ है कि वॉटरमार्क ने विंडो के अंत को पार कर लिया है। FINAL का अर्थ है कि 24 घंटे की स्वचालित सुधार अवधि समाप्त हो गई है। FINAL एक परिचालन अनुबंध है, यह दावा नहीं कि कोई और डेटा मौजूद नहीं है; समाधान अभी भी कटऑफ़ से परे रिकॉर्ड को संभालता है।
बाध्यताओं से समय नीति प्राप्त करें। प्रत्येक सक्रिय स्रोत पार्टीशन एक मान्य event_time निकालता है। एक सामान्य बाउंडेड-आउट-ऑफ़-ऑर्डरनेस रणनीति है:
partition_watermark = max_valid_event_time_seen - out_of_order_bound
operator_watermark = min(active_partition_watermarks)देखे गए आगमन-विलंब वितरण, स्वीकार्य देर-डेटा दर, और ON_TIME विलंबता लक्ष्य से out_of_order_bound चुनें। प्रति-पार्टीशन वॉटरमार्क एक तेज़ पार्टीशन को धीमे पार्टीशन को पूर्ण घोषित करने से रोकते हैं। एक बहु-इनपुट ऑपरेटर न्यूनतम प्रगति लेता है ताकि वह ऐसे इनपुट से आगे न निकल जाए जो अभी भी पुराने रिकॉर्ड उत्पन्न कर सकता है। किसी पार्टीशन को एक सहमत अंतराल के लिए कोई डेटा उत्पन्न न करने के बाद ही निष्क्रिय (idle) चिह्नित करें; अन्यथा यह वैश्विक वॉटरमार्क को हमेशा के लिए रोक सकता है। अत्यधिक छोटा निष्क्रियता थ्रेशोल्ड फिर से शुरू किए गए पार्टीशन से पुराने रिकॉर्ड को देर से आने वाला बना देता है, इसलिए साइड आउटपुट आवश्यक रहता है।
विंडो में दो ट्रिगर परिवार हैं। एक प्रोसेसिंग-टाइम प्रारंभिक ट्रिगर मिनट में एक बार वर्तमान पूर्ण स्नैपशॉट उत्सर्जित करता है। जब वॉटरमार्क window_end को पार कर जाता है, तो ON_TIME संशोधन उत्सर्जित करें। विंडो स्थिति को तब तक बनाए रखें जब तक वॉटरमार्क window_end + 24h को पार न कर जाए; प्रत्येक मान्य देर से आने वाला इवेंट एग्रीगेट को अपडेट करता है और एक उच्च संशोधन उत्पन्न करता है। सफ़ाई के समय, FINAL उत्सर्जित करें और स्थिति छोड़ें। Beam का ट्रिगर और संचय मॉडल दर्शाता है कि "कब उत्सर्जित करना है" और "क्या एक पेन (pane) एक डेल्टा है या संचित स्नैपशॉट" अलग-अलग विकल्प क्यों हैं। यह डिज़ाइन संचित स्नैपशॉट और अपसर्ट का उपयोग करता है, इसलिए डाउनस्ट्रीम सिस्टम को पेन जोड़ने की आवश्यकता नहीं होती है।
एग्रीगेशन से पहले event_id द्वारा डिडुप्लिकेट करें। event_id → (event_time, payload_fingerprint) स्टोर करें। पहले इवेंट को पास होने दें, उसी फ़िंगरप्रिंट वाली दूसरी कॉपी को हटा दें, और भिन्न फ़िंगरप्रिंट वाली समान ID को विरोध स्ट्रीम में भेजें। प्रतिधारण उस सबसे लंबे अंतराल पर निर्भर करता है जिसमें अपस्ट्रीम दूसरी कॉपी को फिर से चला सकता है, न कि केवल विंडो आकार पर। Spark के वॉटरमार्क-आधारित डिडुप्लिकेशन शब्दार्थ को भी सबसे शुरुआती और नवीनतम डुप्लिकेट के बीच टाइमस्टैम्प अंतर से अधिक लंबे विलंब थ्रेशोल्ड की आवश्यकता होती है। बहुत जल्दी होने वाली स्थिति की सफ़ाई देर से आने वाले डुप्लिकेट को फिर से गिनने देती है।
स्थिति का स्पष्ट रूप से अनुमान लगाएं। यदि 20,000 इवेंट प्रति सेकंड का पीक 24 घंटे तक रहता है और प्रत्येक इवेंट अद्वितीय है, तो जॉब 20,000 × 86,400 = 1,728,000,000 IDs याद रखता है। प्रति प्रविष्टि केवल 40 बाइट्स का एक उदाहरण तार्किक पेलोड भी इस ऊपरी सीमा को लगभग 69.12 GB बनाता है। इंजन ऑब्जेक्ट, इंडेक्स, स्टेट बैकएंड, चेकपॉइंट्स और रेप्लिकेट्स भौतिक पदचिह्न (physical footprint) को बढ़ाते हैं। इसलिए स्थिति को कुंजी-आधारित विभाजन और वृद्धिशील चेकपॉइंट की आवश्यकता होती है, जबकि प्रतिधारण रिप्ले अनुबंध का पालन करता है। यदि 99% डुप्लिकेट कम अवधि में आते हैं, तो एक छोटे स्ट्रीमिंग क्षितिज को सिंक व्यावसायिक कुंजी और ऑफ़लाइन समाधान के साथ जोड़ा जा सकता है, लेकिन अवशिष्ट डुप्लिकेट जोखिम को एक TTL के पीछे छिपाने के बजाय परिमाणित किया जाना चाहिए।
आउटपुट पथ को रिप्ले सहन करना चाहिए। आदर्श रूप से, सिंक चेकपॉइंट लेनदेन में भाग लेता है ताकि स्रोत स्थिति, ऑपरेटर स्थिति और विंडो आउटपुट एक साथ कमिट हों। अन्यथा, (window_key, revision) को एक इडेम्पोटेंट अपसर्ट बनाएं: क्रैश के बाद उसी या पुराने संशोधन को फिर से चलाने से नया मान ओवरराइट नहीं हो सकता है। ब्रोकर का एग्जैक्टली-वन्स लेबल तब अपर्याप्त होता है जब बाहरी डेटाबेस राइट चेकपॉइंट पुष्टि के बाहर बैठता है। एंड-टू-एंड शब्दार्थ स्रोत, स्थिति और सिंक की साझा कमिट सीमा पर निर्भर करते हैं।
24 घंटे से परे के इवेंट अत्यधिक-विलंब वाले साइड आउटपुट में जाते हैं और रॉ लॉग में बने रहते हैं। एक दैनिक बैच जॉब उसी ID और राशि नियमों के साथ event_time द्वारा प्रभावित विंडो की पुनर्गणना करता है, फिर स्ट्रीमिंग FINAL स्नैपशॉट के साथ उनकी तुलना करता है। एक महत्वपूर्ण अंतर एक उच्च संशोधन बनाता है या वित्तीय अनुमोदन में प्रवेश करता है। बैच जॉब एक विंडो को परमाणु रूप से (atomically) बदल देता है या एक संस्करणित अपसर्ट करता है। मौजूदा कुल में बैच परिणाम जोड़ने पर बैच के दोबारा चलने पर वही डेटा दो बार गिना जाएगा।
शब्दार्थ को सत्यापित करने के लिए एक न्यूनतम अनुक्रम का उपयोग करें। विंडो [10:00, 11:00) को A=100, B=50, और फिर A की एक समान प्रति प्राप्त होती है। डिडुप्लिकेट किया गया ON_TIME परिणाम 150 है। वॉटरमार्क के 11:00 पार करने के बाद, C=20 अनुमत विलंबता के दौरान आता है, और एक उच्च संशोधन विंडो को 170 में बदल देता है। D वॉटरमार्क द्वारा सफ़ाई बिंदु पार करने के बाद आता है और साफ़ की गई स्थिति को सीधे छूने के बजाय समाधान में जाता है। यदि A के दूसरे पेलोड में राशि 120 है, तो यह विरोध स्ट्रीम में जाती है; परिणाम 170 या 190 नहीं बनना चाहिए।
परीक्षण मैट्रिक्स पार्टीशन में अव्यवस्था, एक निष्क्रिय पार्टीशन जो वॉटरमार्क को रोकता है, भविष्य के टाइमस्टैम्प का क्वारंटाइन, वॉटरमार्क के तुरंत पहले और बाद की सीमाएं, सिंक राइट से पहले और बाद में क्रैश, चेकपॉइंट रिकवरी और रिप्ले, सिंक तक पहुंचने वाले आउट-ऑफ-ऑर्डर संशोधन और दोहराए गए समाधान रन को भी कवर करता है। प्रोडक्शन मेट्रिक्स में p50/p95/p99 और टेल अराइवल डिले, प्रति ऑपरेटर वर्तमान वॉटरमार्क और लैग, अर्ली/ऑन-टाइम/लेट/टू-लेट काउंट और राशियां, डिडुप्लिकेशन हिट्स और फ़िंगरप्रिंट विरोध, स्थिति बाइट्स, चेकपॉइंट अवधि, अस्वीकृत सिंक संशोधन, और बैच-बनाम-स्ट्रीम डेल्टा शामिल हैं।
मजबूत नमूना उत्तर
"मैं एक मिनट की संख्या को एक अनुमान और 24 घंटे को स्वचालित सुधार अवधि के रूप में परिभाषित करूँगा। event_time एक खरीदारी को एक UTC घंटे के लिए असाइन करता है; प्रोसेसिंग समय केवल एक मिनट के शुरुआती ट्रिगर को चलाता है। प्रत्येक सक्रिय इनपुट पार्टीशन अपने सबसे बड़े मान्य इवेंट समय से अव्यवस्था भत्ते को घटाकर एक वॉटरमार्क उत्पन्न करता है, और डाउनस्ट्रीम सबसे धीमे सक्रिय इनपुट का उपयोग करता है। निष्क्रिय पार्टीशन को स्पष्ट रूप से चिह्नित किया जाता है, और अमान्य भविष्य के टाइमस्टैम्प को क्वारंटाइन किया जाता है।
एग्रीगेशन से पहले, मैं ईवेंट समय और एक पेलोड फ़िंगरप्रिंट संग्रहीत करते हुए, event_id द्वारा डिडुप्लिकेट करता हूँ। एक समान पुनः प्रयास को एक बार गिना जाता है; भिन्न सामग्री वाली समान ID विरोध स्ट्रीम में प्रवेश करती है। प्राकृतिक-घंटे की विंडो प्रत्येक मिनट एक पूर्ण EARLY स्नैपशॉट उत्सर्जित करती है, इसका अंत वॉटरमार्क के पीछे होने के बाद एक ON_TIME स्नैपशॉट उत्सर्जित करती है, और 24 घंटे तक स्थिति बनाए रखती है। उस अवधि में देर से आने वाले इवेंट कुल को अपडेट करते हैं और संशोधन को बढ़ाते हैं; सफ़ाई के समय FINAL उत्सर्जित होता है।
सिंक प्रत्येक स्नैपशॉट को डेल्टा के रूप में जोड़ने के बजाय विंडो कुंजी और संशोधन द्वारा अपसर्ट करता है। स्रोत रिप्ले करने योग्य है और ऑपरेटर स्थिति चेकपॉइंट से पुनर्स्थापित होती है। यदि सिंक किसी चेकपॉइंट लेनदेन में शामिल नहीं हो सकता है, तो संस्करणित राइट्स रिकवरी रिप्ले को एक नए परिणाम को बदलने से रोकते हैं। 24 घंटे से अधिक देर के इवेंट एक साइड आउटपुट में जाते हैं, और एक दैनिक बैच जॉब रॉ लॉग की पुनर्गणना करता है और FINAL के साथ इसकी तुलना करता है।
यदि 20,000 इवेंट प्रति सेकंड का पीक 24 घंटे तक रहता है और सभी ID अद्वितीय हैं, तो ऊपरी सीमा 1.728 बिलियन डिडुप्लिकेशन प्रविष्टियां हैं। प्रति प्रविष्टि 40 तार्किक बाइट्स भी 69.12 GB है, इसलिए मापे गए विलंब वितरण को प्रतिधारण को उचित ठहराना चाहिए, और स्थिति तथा चेकपॉइंट लागत की निगरानी की आवश्यकता है। मैं डुप्लिकेट, अव्यवस्था, विलंबता, भविष्य के टाइमस्टैम्प, निष्क्रिय पार्टीशन और रिकवरी के लिए फॉल्ट इंजेक्शन के साथ समाप्त करूँगा, यह दावा करते हुए कि अंतिम विंडो ऑफ़लाइन पुनर्गणना के बराबर है।"
सामान्य गलतियाँ
- प्रोसेसिंग समय द्वारा व्यावसायिक विंडो असाइन करना → ऐतिहासिक रिप्ले के दौरान वही ईवेंट एक अलग घंटे में आ सकता है → सदस्यता के लिए इवेंट समय और केवल प्रारंभिक आउटपुट के लिए प्रोसेसिंग समय का उपयोग करें।
- वॉटरमार्क को "कोई पुराना डेटा नहीं आ सकता" के रूप में वर्णित करना → यह आमतौर पर एक अनुमानी प्रगति अनुमान है, इसलिए एक पुराना ईवेंट अभी भी दिखाई दे सकता है → अनुमत विलंबता, एक टू-लेट साइड आउटपुट और समाधान को परिभाषित करें।
- वॉटरमार्क विलंब, अनुमत विलंबता और डिडुप्लिकेशन TTL के लिए एक अस्पष्ट मान का उपयोग करना → वे क्रमशः फायरिंग, विंडो स्थिति और डुप्लिकेट पहचान को नियंत्रित करते हैं → उन्हें विलंबता SLO, सुधार अवधि और अपस्ट्रीम रिप्ले अनुबंध से प्राप्त करें।
- प्रत्येक फायरिंग पर एक कुल जोड़ना → EARLY, ON_TIME, और देर से होने वाली फायरिंग बार-बार जोड़ी जाती हैं → एक पूर्ण स्नैपशॉट उत्सर्जित करें और विंडो और संशोधन द्वारा अपसर्ट करें, या एक वापस लेने योग्य-डेल्टा प्रोटोकॉल परिभाषित करें।
- वैश्विक अधिकतम इवेंट समय से आगे बढ़ना → एक तेज़ पार्टीशन या खराब भविष्य का टाइमस्टैम्प धीमे-पार्टीशन डेटा को बहुत जल्दी देर से आने वाला बना देता है → पार्टीशन प्रगति उत्पन्न करें, सक्रिय-इनपुट न्यूनतम लें, और टाइमस्टैम्प मान्य करें।
- न्यूनतम में स्थायी रूप से शांत पार्टीशन रखना → वॉटरमार्क रुक जाता है, जिससे विंडो और डिडुप्लिकेशन स्थिति साफ़ होने से रुक जाती है → अवलोकन योग्य निष्क्रियता हैंडलिंग का उपयोग करें और फिर से शुरू किए गए देर के डेटा को सही ढंग से रूट करें।
- केवल ईवेंट ID संग्रहीत करना → निर्माता द्वारा किसी ID का पुन: उपयोग चुपचाप निगल लिया जाता है → एक पेलोड फ़िंगरप्रिंट भी संग्रहीत करें और विरोधों को क्वारंटाइन करें।
- "एग्जैक्टली-वन्स सक्षम है" कहना → एक बाहरी सिंक चेकपॉइंट सीमा साझा नहीं कर सकता है → एंड-टू-एंड कमिट पथ को ट्रेस करें और एक ट्रांजेक्शनल सिंक या संस्करणित इडेम्पोटेंट राइट्स का उपयोग करें।
- स्थिति की सफ़ाई के बाद डेटा छोड़ना → स्ट्रीम स्थिर दिखती है जबकि वित्तीय परिणाम अस्पष्ट हो जाते हैं → रॉ इवेंट्स को बनाए रखें और साइड आउटपुट, बैच पुनर्गणना और विसंगति समीक्षा लागू करें।
अनुवर्ती प्रश्न और उत्तर
अनुवर्ती 1: वॉटरमार्क का अव्यवस्था भत्ता कितना बड़ा होना चाहिए?
उत्पादन में processing_time - event_time को मापें और इसे स्रोत, क्लाइंट संस्करण और क्षेत्र द्वारा विभाजित करें। पहले नवीनतम स्वीकार्य ON_TIME परिणाम और सुधार पथ में अनुमत इवेंट्स का हिस्सा चुनें; फिर एक ऐसा पर्सेंटाइल चुनें जो दोनों को संतुष्ट करता हो। p50, p95, p99 और टेल की निगरानी जारी रखें। वितरण परिवर्तन को नीति को स्वचालित रूप से स्थानांतरित करने की अनुमति देने के बजाय कॉन्फ़िगरेशन रिलीज़ के माध्यम से जाना चाहिए। वॉटरमार्क समय पर निष्पादन को नियंत्रित करता है, जबकि 24 घंटे की सुधार अवधि अभी भी एक लंबी पूंछ (tail) को कवर करती है।
अनुवर्ती 2: एक शांत Kafka पार्टीशन विंडो को क्यों रोक सकता है?
एक बहु-इनपुट ऑपरेटर की सुरक्षित प्रगति न्यूनतम इनपुट वॉटरमार्क है। एक पार्टीशन जो आगे नहीं बढ़ता है वह उस न्यूनतम को अपरिवर्तित रखता है। यह पुष्टि करने के बाद कि उसने निष्क्रियता थ्रेशोल्ड के लिए कोई ईवेंट उत्पन्न नहीं किया है, इसे निष्क्रिय चिह्नित करें ताकि यह अस्थायी रूप से न्यूनतम से बाहर निकल जाए। थ्रेशोल्ड को बहुत छोटा सेट न करें: फिर से शुरू किए गए पार्टीशन के रिकॉर्ड नए वॉटरमार्क से पुराने हो सकते हैं। उन्हें अनुमत-देर प्रसंस्करण या साइड आउटपुट में प्रवेश करना चाहिए, और निष्क्रिय-स्थिति परिवर्तनों को अपने स्वयं के मीट्रिक की आवश्यकता होती है।
अनुवर्ती 3: 24 घंटे की डिडुप्लिकेशन स्थिति यह वादा क्यों नहीं करती कि दोहराव हमेशा के लिए असंभव है?
चौबीस घंटे केवल इस समस्या की स्वचालित सुधार अवधि को कवर करते हैं। यदि अपस्ट्रीम 25वें घंटे में एक ID को फिर से चलाता है, तो स्ट्रीमिंग स्थिति चली जाती है और ईवेंट एग्रीगेट में फिर से प्रवेश कर सकता है। स्थायी डिडुप्लिकेशन के लिए एक लंबे समय तक चलने वाले व्यावसायिक अद्वितीय इंडेक्स, एक क्वेरी योग्य ईवेंट रजिस्ट्री, या पूर्ण ऑफ़लाइन डिडुप्लिकेशन की आवश्यकता होती है; प्रत्येक स्टोरेज और राइट लागत जोड़ता है। गारंटी को "घोषित रिप्ले क्षितिज के भीतर डिडुप्लिकेट किया गया" के रूप में बताएं, फिर बाद के रिकॉर्ड का मिलान करें।
अनुवर्ती 4: क्या होगा यदि डाउनस्ट्रीम डेटाबेस INSERT का समर्थन करता है लेकिन अपसर्ट का नहीं?
प्रत्येक परिणाम को विंडो और संशोधन द्वारा अनुक्रमित एक अपरिवर्तनीय चेंजलॉग में लिखें। रीड लेयर प्रत्येक विंडो के लिए अधिकतम संशोधन से वर्तमान दृश्य बनाती है। उपभोक्ताओं को पता होना चाहिए कि प्रत्येक रिकॉर्ड एक पूर्ण स्नैपशॉट है, डेल्टा नहीं, और सभी संशोधनों का योग नहीं करना चाहिए। यदि क्वेरी लागत बहुत अधिक है, तो ऑडिट और रिकवरी के लिए संस्करण लॉग को बनाए रखते हुए परमाणु प्रतिस्थापन का समर्थन करने वाली सर्विंग टेबल में एसिंक्रोनस रूप से कॉम्पैक्ट करें।
अनुवर्ती 5: आप कैसे साबित करते हैं कि रिकवरी में न तो अधिक गिनती होती है और न ही कम गिनती?
एक निश्चित इनपुट के लिए, अपेक्षित डिडुप्लिकेशन सेट और विंडो संशोधन रिकॉर्ड करें। कार्य को तीन सीमाओं पर समाप्त करें: स्रोत पढ़ने के बाद लेकिन स्थिति चेकपॉइंट से पहले, स्थिति चेकपॉइंटिंग के बाद लेकिन सिंक पुष्टि से पहले, और सिंक राइट के बाद लेकिन चेकपॉइंट पुष्टि से पहले। पुनर्प्राप्ति के बाद फिर से चलाएं और दावा करें कि प्रत्येक ईवेंट ID एक बार योगदान देता है, सिंक केवल सबसे बड़ा संशोधन रखता है, और अंतिम विंडो ऑफ़लाइन पुनर्गणना के बराबर होती है। इन कमिट सीमाओं पर फॉल्ट इंजेक्शन के बिना RUNNING स्थिति एंड-टू-एंड शुद्धता का प्रमाण नहीं है।