Flow, Kotlin Coroutines लाइब्रेरी से एक एसिंक्रोनस डेटा स्ट्रीम प्रकार है, जो cold सेमैंटिक्स लागू करता है। Kotlin Documentation, 2025 के अनुसार, Flow map, filter, catch और collect ऑपरेटरों के साथ मानों का अनुक्रम उत्सर्जित करने की अनुमति देता है। LiveData के विपरीत, Flow कोरूटीन पर बनाया गया है और बैकप्रेशर का समर्थन करता है।
मुख्य बिंदु
Flow kotlinx.coroutines.flow पैकेज का एक प्रकार है, जो एक cold एसिंक्रोनस डेटा स्ट्रीम का प्रतिनिधित्व करता है। इसके मूल में, Flow एक कोरूटीन अनुक्रम है जो emit() फ़ंक्शन के माध्यम से मान उत्सर्जित करता है और या तो सफलतापूर्वक या अपवाद के साथ समाप्त होता है। स्ट्रीम संग्रहण टर्मिनल ऑपरेटर collect() के माध्यम से किया जाता है, जो एक suspend फ़ंक्शन है।
Cold स्ट्रीम का अर्थ है कि flow बिल्डर के अंदर का कोड प्रत्येक सब्सक्राइबर के लिए फिर से चलता है। RxJava में Observable.fromIterable समान व्यवहार करता है: एक नया सब्सक्राइबर शुरुआत से सभी मान प्राप्त करता है। Flow में, यह suspend फ़ंक्शन collect के माध्यम से लागू किया गया है, जो डेटा संग्रहण की पूरी अवधि के लिए कोरूटीन को ब्लॉक करता है।
Kotlin Flow बनाने के कई तरीके प्रदान करता है: flow { } — emit() के साथ मूल निर्माण, flowOf(vararg values) — मानों के एक निश्चित सेट के लिए, .asFlow() — संग्रह और Sequence के लिए एक्सटेंशन। सभी बिल्डर cold हैं — डेटा केवल टर्मिनल ऑपरेटर को कॉल करने पर उत्पन्न होता है।
Cold और hot स्ट्रीम में विभाजन रिएक्टिव प्रोग्रामिंग की एक मुख्य अवधारणा है। Cold स्ट्रीम (Flow, Observable) सब्सक्रिप्शन पर डेटा जनरेशन शुरू करता है। Hot स्ट्रीम (Channel, SharedFlow) स्वतंत्र रूप से डेटा उत्सर्जित करता है — सब्सक्राइबर केवल वही प्राप्त करता है जो सब्सक्रिप्शन के बाद होता है, अनुक्रम की शुरुआत के बिना।
SharedFlow एक hot Flow है जिसके कई सब्सक्राइबर हो सकते हैं और replay कॉन्फ़िगर होने पर हाल के मानों को फिर से चला सकता है। SharedFlow इवेंट्स (एक बार की सूचनाओं) के लिए उपयुक्त है। StateFlow इसका एक प्रकार है जिसमें एक निश्चित स्थिति मान होता है, जो नए सब्सक्राइबर के लिए नवीनतम मान कैश करता है।
ChannelFlow अंदरूनी रूप से Channel का उपयोग करता है, Flow और Channel के गुणों को जोड़ता है। यह capacity के माध्यम से बफरिंग और बैकप्रेशर का समर्थन करता है। ChannelFlow कॉलबैक API को रिएक्टिव स्ट्रीम में बदलने में उपयोगी है, जहां मान विभिन्न कोरूटीन से उत्सर्जित होते हैं।
Cold Flow को hot SharedFlow में बदलने के लिए shareIn(scope, started, replay) ऑपरेटर का उपयोग किया जाता है। पैरामीटर started शुरुआत के क्षण को नियंत्रित करता है: SharingStarted.WhileSubscribed() — तब तक सक्रिय जब तक सब्सक्राइबर हैं, Lazily — पहले सब्सक्राइबर पर शुरुआत, Eagerly — तत्काल शुरुआत। उल्टा रूपांतरण — hot से cold: StateFlow.asFlow() एक cold Flow लौटाता है जो collect पर StateFlow का वर्तमान मान उत्सर्जित करता है। यह परीक्षण के लिए सुविधाजनक है।
Flow ऑपरेटरों का एक समृद्ध सेट प्रदान करता है जो कोरूटीन के अंदर suspend फ़ंक्शन के रूप में काम करते हैं। ऑपरेटर स्टेटलेस होते हैं और एक नया Flow लौटाते हैं — मूल स्ट्रीम अपरिवर्तित रहती है। यह बिना साइड इफेक्ट के सुरक्षित रूपांतरण श्रृंखला बनाने की अनुमति देता है।
map ऑपरेटर प्रत्येक स्ट्रीम मान को एसिंक्रोनस या सिंक्रोनस रूपांतरण के माध्यम से रूपांतरित करता है। filter केवल उन मानों को पास करता है जो शर्त को पूरा करते हैं। catch टर्मिनल ऑपरेटर से पहले अपवादों को पकड़ता है और स्ट्रीम को पुनर्प्राप्त करने की अनुमति देता है। flatMapLatest नए मान आने पर पिछले उत्सर्जन को रद्द करता है — Rx में switchMap के समान।
Flow में debounce ऑपरेटर मान प्रकाशन को निर्दिष्ट टाइमआउट तक विलंबित करता है। यदि इस समय के दौरान कोई नया मान आता है, तो टाइमर रीसेट हो जाता है। Android में, debounce का उपयोग खोज के लिए किया जाता है: अनुरोध केवल 300-400 ms के ठहराव के बाद भेजा जाता है, जिससे API कॉल 3-5 गुना कम हो जाती हैं।
collect() के अलावा, Flow अन्य टर्मिनल ऑपरेटरों का समर्थन करता है: toList() सभी मानों को एक सूची में एकत्र करता है — परीक्षणों के लिए उपयोगी, first() पहला तत्व लौटाता है और स्ट्रीम को रद्द करता है, single() बिल्कुल एक तत्व की अपेक्षा करता है। fold(initial) पारित फ़ंक्शन के माध्यम से मानों को संचित करता है। सभी टर्मिनल ऑपरेटर suspend फ़ंक्शन हैं और उन्हें कोरूटीन या किसी अन्य suspend फ़ंक्शन के अंदर कॉल किया जाना चाहिए।
पहला उदाहरण — map ऑपरेटर के माध्यम से रूपांतरण के साथ संख्याएँ उत्पन्न करने वाला मूल Flow:
val numberFlow = flow {
for (i in 1..5) {
delay(500)
emit(i)
}
}
scope.launch {
numberFlow
.map { "संख्या: $it" }
.collect { value ->
println(value)
}
}
दूसरा उदाहरण — catch के माध्यम से फ़िल्टरिंग और त्रुटि प्रबंधन के साथ स्ट्रीम रूपांतरण:
flow {
emit("data1")
emit("data2")
throw RuntimeException("network error")
}
.catch { e ->
emit("fallback_data")
}
.collect { value ->
println(value)
}
तीसरा उदाहरण — Jetpack Compose में रिएक्टिव UI के लिए ViewModel में StateFlow का उपयोग:
class SearchViewModel : ViewModel() {
private val _query = MutableStateFlow("")
val results: StateFlow<List<Result>> = _query
.debounce(300)
.flatMapLatest { query ->
repository.search(query)
}
.catch { emit(emptyList()) }
.stateIn(viewModelScope, SharingStarted.WhileSubscribed(5000), emptyList())
fun onQueryChanged(query: String) {
_query.value = query
}
}
StateFlow एक ही वर्तमान मान वाला hot Flow है। यह नवीनतम मान कैश करता है और इसे तुरंत नए सब्सक्राइबर को देता है। StateFlow स्थिति के लिए एक अवलोकनीय कंटेनर है, equals तुलना का समर्थन करता है — यदि नया मान वर्तमान से मेल खाता है, तो कोई उत्सर्जन नहीं होता है। Jetpack Compose collectAsState() के माध्यम से StateFlow का उपयोग करता है।
SharedFlow बिना अनिवार्य प्रारंभिक मान के अधिक लचीला hot Flow है। SharedFlow को replay (नए सब्सक्राइबर के लिए मानों की संख्या), extraBufferCapacity (replay से परे बफर), और onBufferOverflow (अतिप्रवाह पर रणनीति) के माध्यम से कॉन्फ़िगर किया जाता है। SharedFlow एक बार के इवेंट के लिए आदर्श है: नेविगेशन, Snackbar, विश्लेषिकी।
Android आर्किटेक्चर में Flow को Google द्वारा प्राथमिक डेटा स्रोत (परत: Repository → UseCase → ViewModel) के रूप में अनुशंसित किया गया है। LiveData लचीलेपन में Flow से कमतर है: Flow कोरूटीन, ऑपरेटर, बैकप्रेशर का समर्थन करता है और UI परत के बाहर काम करता है। LiveData से Flow में माइग्रेशन आधुनिक Android प्रोजेक्ट में एक मानक अभ्यास है।
ViewModel में Flow का उपयोग करते समय सही प्रकार चुनना महत्वपूर्ण है। StateFlow UI स्थिति के लिए आदर्श है जो स्क्रीन रोटेशन से बचनी चाहिए। SharedFlow उन इवेंट के लिए उपयुक्त है जहां पुनः प्रसंस्करण अस्वीकार्य है — उदाहरण के लिए, नेविगेशन। lifecycleScope में collect() के साथ Flow निष्पादन संदर्भ पर अधिकतम नियंत्रण देता है लेकिन स्क्रीन छोड़ने पर मैन्युअल रद्दीकरण की आवश्यकता होती है।
Flow का परीक्षण kotlinx-coroutines-test के माध्यम से किया जाता है। लाइब्रेरी TestDispatcher प्रदान करती है — वर्चुअल समय जो देरी (delay) को तेज करने और कोरूटीन निष्पादन क्रम को नियंत्रित करने की अनुमति देता है। TestScope.runTest { } Flow परीक्षण के लिए एक पृथक वातावरण बनाता है। toList() ऑपरेटर का उपयोग अक्सर परीक्षणों में टाइमआउट के साथ सभी flow मानों को एकत्र करने के लिए किया जाता है, यह सत्यापित करने के लिए कि स्ट्रीम ने सही डेटा अनुक्रम उत्सर्जित किया।
Flow Room (Android DB लाइब्रेरी) के साथ अच्छी तरह से एकीकृत होता है: DAO विधियाँ Flow<List<Entity>> लौटा सकती हैं। Room किसी भी तालिका परिवर्तन पर स्वचालित रूप से एक नया मान उत्सर्जित करता है — UI बिना मैन्युअल ट्रिगर के अपडेट होता है। यह InvalidationTracker के माध्यम से कार्यान्वित किया जाता है, जो अंदरूनी रूप से callbackFlow के साथ Flow का उपयोग करता है। यह दृष्टिकोण LiveData की आवश्यकता को समाप्त करता है और डेटा परत को पूरी तरह से कोरूटीन-उन्मुख बनाता है। Jetpack Compose collectAsState() के माध्यम से StateFlow की सदस्यता लेता है और केवल उन घटकों को फिर से खींचता है जिनका डेटा बदल गया है — यह LiveData-उन्मुख आर्किटेक्चर के साथ अप्राप्य प्रदर्शन देता है। DataStore (SharedPreferences का प्रतिस्थापन) भी Flow<Preferences> लौटाता है, जो मैन्युअल अपडेट ट्रिगर के बिना एप्लिकेशन सेटिंग्स का रिएक्टिव रीडिंग प्रदान करता है।
Flow बिना अतिरिक्त लाइब्रेरी के JVM पर kotlinx-coroutines-core के माध्यम से इंटर-प्रोसेस संचार का समर्थन करता है। उदाहरण के लिए, Ktor पर सर्वर एप्लिकेशन में, Flow आने वाले WebSocket संदेश स्ट्रीम का प्रतिनिधित्व कर सकता है। प्रत्येक संदेश स्ट्रीम में उत्सर्जित होता है, ऑपरेटरों के माध्यम से फ़िल्टरिंग और एकत्रीकरण से गुजरता है, और परिणाम क्लाइंट को भेजा जाता है। यह दृष्टिकोण Kotlin प्रोजेक्ट में Reactor या RxJava जैसी रिएक्टिव लाइब्रेरी को बदलता है।
मौजूदा RxJava कोड के साथ Flow संगतता kotlinx-coroutines-rx3 मॉड्यूल द्वारा प्रदान की जाती है। एक्सटेंशन फ़ंक्शन Flow.asObservable() Flow को RxJava 3 के Observable में बदलता है। उल्टा रूपांतरण — CompletableSource.asFlow(), Observable.asFlow()। यह RxJava से कोरूटीन में माइग्रेशन को सरल बनाता है: प्रोजेक्ट को चरणों में फिर से लिखा जा सकता है, कुछ परतों को RxJava पर छोड़ते हुए। रूपांतरण करते समय, cold/hot सेमैंटिक्स में अंतर पर विचार करना होगा: Observable cold और hot दोनों हो सकता है, Flow सामान्य Flow के लिए हमेशा cold और SharedFlow के लिए hot होता है।
Flow में त्रुटि प्रबंधन की एक विशिष्टता है: यदि टर्मिनल ऑपरेटर से पहले flow बिल्डर के अंदर अपवाद होता है, तो यह catch में प्रचारित होता है। यदि बिल्डर के बाद किसी ऑपरेटर में अपवाद होता है, तो उसे उस ऑपरेटर के बाद catch पकड़ता है। retryWhen शर्त के साथ सब्सक्रिप्शन को पुनः प्रयास करने की अनुमति देता है: नेटवर्क त्रुटि पर 3 बार तक पुनः प्रयास करें, लेकिन CancellationException पर पुनः प्रयास न करें। Flow स्थिति-निर्भर त्रुटियों को समाप्त करता है क्योंकि यह स्थिति संग्रहीत नहीं करता — यह Observable की तुलना में डिबगिंग को सरल बनाता है, जहां Subject आंतरिक स्थिति संग्रहीत करता है।
kotlinx-coroutines-test के साथ Flow परीक्षण देरी को अनुकरण करने के लिए TestDispatcher का उपयोग करता है। Turbine Flow परीक्षण के लिए एक लोकप्रिय सामुदायिक लाइब्रेरी है: test { } Flow शुरू करता है, awaitItem() अगले मान की प्रतीक्षा करता है, awaitComplete() पूर्णता की प्रतीक्षा करता है। Turbine एक डिफ़ॉल्ट टाइमआउट जोड़ता है, जो परीक्षणों को हैंग होने से रोकता है। StateFlow के परीक्षण के लिए, कालानुक्रमिक क्रम में मान सत्यापन के साथ .testIn(scope) का उपयोग करें।
अक्सर पूछे जाने वाले प्रश्न
Flow कोरूटीन, ऑपरेटर और बैकप्रेशर के समर्थन वाला एक एसिंक्रोनस स्ट्रीम है, जो किसी भी आर्किटेक्चर परत पर काम करता है। LiveData केवल UI परत के लिए एक lifecycle-aware घटक है। Google व्यावसायिक तर्क और रिपॉजिटरी के लिए Flow की सिफारिश करता है, ViewModel में सरल अवलोकन के लिए LiveData की।
StateFlow — जब UI स्थिति (कार्य सूची, खोज पाठ, लोडिंग फ़्लैग) संग्रहीत करनी हो — प्रत्येक सब्सक्राइबर को वर्तमान मान मिलता है। SharedFlow — एक बार के इवेंट (नेविगेशन, Snackbar) के लिए। StateFlow का उपयोग इवेंट के लिए नहीं किया जाना चाहिए क्योंकि नया मान फिर से संसाधित हो सकता है।
Flow में, बैकप्रेशर suspend तंत्र के माध्यम से लागू किया गया है: यदि कलेक्टर पिछले मान को संसाधित कर रहा है तो emit() कोरूटीन को निलंबित करता है। ChannelFlow में चैनल (Channel) में capacity आकार का बफर होता है। अतिप्रवाह पर: suspending (प्रतीक्षा), drop (छोड़ना) या conflate (अंतिम से बदलना)।
callbackFlow का उपयोग करें — कॉलबैक API के लिए Flow बिल्डर। अंदर, कॉलबैक के अंदर emit(value) के साथ registerCallback() को कॉल करें। awaitClose कोरूटीन रद्द होने पर unregisterCallback() को कॉल करने की गारंटी देता है। callbackFlow अंदरूनी रूप से Channel(UNLIMITED) के माध्यम से बफरिंग का समर्थन करता है।
हाँ, कनवर्टर के माध्यम से: Flow.asObservable() kotlinx-coroutines-rx3 पैकेज से Flow को RxJava 3 Observable में बदलता है। उल्टा — CompletableSource.asFlow() Single/Completable/Maybe के लिए। यह बड़े प्रोजेक्ट में RxJava से कोरूटीन में माइग्रेट करते समय उपयोगी है।
सारांश
हम एक मोबाइल एप्लिकेशन टर्नकी विकसित करेंगे
IT Sectr 2017 से स्टार्टअप और व्यवसायों के लिए iOS और Android एप्लिकेशन बनाता है। हम आपको सलाह देंगे और सर्वोत्तम समाधान प्रस्तावित करेंगे।
यह भी पढ़ें