করুটিনে, ফ্লো হল এমন একটি ধরন যা একাধিক ভ্যালু ধারাবাহিকভাবে এমিট করতে পারে, এটি সাসপেন্ড ফাংশনের বিপরীত, যা শুধুমাত্র একটি ভ্যালু রিটার্ন করে। যেমন, আপনি কোনও ডেটাবেস থেকে লাইভ আপডেট পেতে ফ্লো ব্যবহার করতে পারেন।
ফ্লো কোরাউটিনের উপরে তৈরি করা হয় এবং একাধিক ভ্যালু প্রদান করতে পারে।
ফ্লো হল কনসেপচুয়ালি ডেটার স্ট্রিম যা অ্যাসিঙ্ক্রোনাস
কম্পিউট করা যায়। ইমিট করা মান একই ধরনের হতে হবে। যেমন, Flow<Int> হল এমন একটি ফ্লো যা পূর্ণসংখ্যা ভ্যালু নির্গমন করে।
ফ্লো হল Iterator-এর মতো যা ভ্যালুর একটি সিকোয়েন্স তৈরি করে, তবে এটি অ্যাসিঙ্ক্রোনাসভাবে ভ্যালু তৈরি ও গ্রহণ করতে
সাসপেন্ড ফাংশন ব্যবহার করে।
এর অর্থ হল, উদাহরণস্বরূপ, ফ্লো নিরাপদে একটি নেটওয়ার্ক অনুরোধ করতে পারে
যা মূল থ্রেডকে ব্লক না করেই পরবর্তী ভ্যালু তৈরি করে।
ডেটা স্ট্রিম করার ক্ষেত্রে তিনটি এন্টিটি জড়িত থাকে:
- প্রডিউসার এমন ডেটা তৈরি করে যা স্ট্রিমের সাথে যোগ করা হয়। কোরুটিনের জন্য, ফ্লো অ্যাসিঙ্ক্রোনাসভাবে ডেটা তৈরি করতে পারে।
- (ঐচ্ছিক) মধ্যস্থতাকারী স্ট্রিম বা স্ট্রিমের মধ্যে নির্গত প্রতিটি ভ্যালু পরিবর্তন করতে পারে।
- কনজিউমার স্ট্রিম থেকে ভ্যালু কনজিউ করে।
Android-এ, রেপোজিটরি হল সাধারণত UI ডেটার একটি প্রোডিউসার যার কনজিউমার হিসেবে ইউজার ইন্টারফেস (UI) থাকে যা শেষ পর্যন্ত ডেটা দেখায়। অন্যান্য ক্ষেত্রে, UI লেয়ার হল ব্যবহারকারীর ইনপুট ইভেন্টের প্রযোজক এবং হায়ারার্কির অন্যান্য লেয়ার সেগুলি ব্যবহার করে। প্রডিউসার ও কনজিউমারের মধ্যেকার লেয়ারগুলি সাধারণত মধ্যস্থতাকারী হিসেবে কাজ করে এবং পরের লেয়ারের প্রয়োজনীয়তা অনুযায়ী ডেটা স্ট্রিমকে পরিবর্তন করে।
ফ্লো তৈরি করা
ফ্লো তৈরি করতে, ফ্লো বিল্ডার
API ব্যবহার করুন। flow বিল্ডার ফাংশন একটি নতুন ফ্লো তৈরি করে যেখানে আপনি ম্যানুয়ালি
emit
ফাংশন ব্যবহার করে ডেটার স্ট্রিমে নতুন ভ্যালু
এমিট করতে পারেন।
নিম্নলিখিত উদাহরণে, একটি ডেটা সোর্স নির্দিষ্ট সময় অন্তর অটোমেটিক লেটেস্ট খবর সংগ্রহ করে। যেহেতু সাসপেন্ড ফাংশন একাধিক ধারাবাহিক ভ্যালু রিটার্ন করতে পারে না, তাই ডেটা সোর্স এই প্রয়োজনীয়তা পূরণ করার জন্য একটি ফ্লো তৈরি করে ও রিটার্ন করে। এই ক্ষেত্রে, ডেটা সোর্সটি প্রডিউসার হিসেবে কাজ করে।
class NewsRemoteDataSource( private val newsApi: NewsApi, private val refreshIntervalMs: Long = 5000 ) { val latestNews: Flow<List<ArticleHeadline>> = flow { while (true) { val latestNews = newsApi.fetchLatestNews() emit(latestNews) // Emits the result of the request to the flow delay(refreshIntervalMs) // Suspends the coroutine for some time } } } // Interface that provides a way to make network requests with suspend functions interface NewsApi { suspend fun fetchLatestNews(): List<ArticleHeadline> }
flow বিল্ডার একটি কোরাউটিনের মধ্যে এক্সিকিউট করা হয়। তাই, এটি একই অ্যাসিঙ্ক্রোনাস API থেকে সুবিধা পায়, তবে কিছু বিধিনিষেধ প্রযোজ্য:
- Flow হল ক্রমিক। প্রডিউসার যেহেতু একটি কোরাউটিনে আছে, তাই
সাসপেন্ড ফাংশনকে কল করলে, সাসপেন্ড ফাংশন
রিটার্ন না করা পর্যন্ত প্রডিউসার সাসপেন্ড হয়ে যায়। এই উদাহরণে,
fetchLatestNewsনেটওয়ার্ক অনুরোধ সম্পূর্ণ না হওয়া পর্যন্ত প্রোডিউসর সাসপেন্ড করা হয়। তবেই ফলাফল স্ট্রিম করা হয়। flowবিল্ডারের সাহায্যে, প্রযোজকemitভ্যালুগুলি একটি আলাদাCoroutineContextথেকে নিতে পারবেন না। তাই, নতুন কোরাউটিন তৈরি করে বাwithContextকোডের ব্লক ব্যবহার করেemit-কে অন্যCoroutineContext-এ কল করবেন না। এইসব ক্ষেত্রে আপনিcallbackFlowএর মতো অন্যান্য ফ্লো বিল্ডার ব্যবহার করতে পারেন।
স্ট্রিম পরিবর্তন করা
ভ্যালু ব্যবহার না করেই ডেটার স্ট্রিম পরিবর্তন করতে, মধ্যস্থতাকারীরা ইন্টারমিডিয়েট অপারেটর ব্যবহার করতে পারেন। এইসব অপারেটর হল এমন ফাংশন যা ডেটার স্ট্রিমে প্রয়োগ করা হলে, অপারেশনের একটি চেইন সেট-আপ করে যা ভবিষ্যতে ভ্যালু ব্যবহার না করা পর্যন্ত এক্সিকিউট করা হয় না। ফ্লো রেফারেন্স ডকুমেন্টেশন থেকে ইন্টারমিডিয়েট অপারেটর সম্পর্কে আরও জানুন।
নিচের উদাহরণে, রিপোজিটরি লেয়ার, View-এ দেখানো ডেটা ট্রান্সফর্ম করতে ইন্টারমিডিয়েট অপারেটর
map
ব্যবহার করে:
class NewsRepository( private val newsRemoteDataSource: NewsRemoteDataSource, private val userData: UserData ) { /** * Returns the favorite latest news applying transformations on the flow. * These operations are lazy and don't trigger the flow. They just transform * the current value emitted by the flow at that point in time. */ val favoriteLatestNews: Flow<List<ArticleHeadline>> = newsRemoteDataSource.latestNews // Intermediate operation to filter the list of favorite topics .map { news -> news.filter { userData.isFavoriteTopic(it) } } // Intermediate operation to save the latest news in the cache .onEach { news -> saveInCache(news) } }
ইন্টারমিডিয়েট অপারেটররা একে একে প্রয়োগ করা যায়, এর ফলে অপারেশনের একটি চেইন তৈরি হয় যা কোনও আইটেম ফ্লোতে এমিট করা হলে অলসভাবে এক্সিকিউট করা হয়। মনে রাখবেন, কোনও স্ট্রিমে শুধু ইন্টারমিডিয়েট অপারেটর প্রয়োগ করলেই ফ্লো সংগ্রহ করা শুরু হয় না।
ফ্লো থেকে সংগ্রহ করা
ভ্যালু শোনার জন্য ফ্লো ট্রিগার করতে টার্মিনাল অপারেটর ব্যবহার করুন।
স্ট্রিম থেকে সব ভ্যালু যেভাবে এমিট করা হয় সেইভাবে পেতে, collect ব্যবহার করুন।
আপনি
অফিসিয়াল ফ্লো ডকুমেন্টেশন থেকে টার্মিনাল অপারেটর সম্পর্কে আরও জানতে পারবেন।
collect একটি সাসপেন্ড ফাংশন হওয়ায়, এটি একটি
করুটিনের মধ্যে এক্সিকিউট করতে হবে। এটি প্যারামিটার হিসেবে একটি ল্যাম্বডা নেয় যা
প্রতিটি নতুন ভ্যালুর উপর কল করা হয়। যেহেতু এটি একটি সাসপেন্ড ফাংশন, তাই যে কোরাউটিন
কল করে collect সেটি ফ্লো বন্ধ না হওয়া পর্যন্ত সাসপেন্ড হতে পারে।
আগের উদাহরণটি চালিয়ে গেলে, এখানে
ViewModel রেপোজিটরি লেয়ার থেকে ডেটা কনজিউ করার একটি সহজ ইমপ্লিমেন্টেশন দেওয়া হল:
class LatestNewsViewModel( private val newsRepository: NewsRepository ) : ViewModel() { init { viewModelScope.launch { // Trigger the flow and consume its elements using collect newsRepository.favoriteLatestNews.collect { favoriteNews -> // Update UI with the latest favorite news } } } }
ফ্লো সংগ্রহ করা হলে, তা প্রোডিউসারকে ট্রিগার করে, যা লেটেস্ট খবর রিফ্রেশ করে
এবং নির্দিষ্ট সময় অন্তর নেটওয়ার্ক অনুরোধের ফলাফল ইমিট করে। যেহেতু
প্রডিউসার while(true) লুপের সাথে সবসময় অ্যাক্টিভ থাকে, তাই ViewModel মুছে ফেলা হলে এবং
viewModelScope বাতিল করা হলে ডেটার স্ট্রিম
বন্ধ হয়ে যাবে।
নিম্নলিখিত কারণে ফ্লো সংগ্রহ করা বন্ধ হয়ে যেতে পারে:
- আগের উদাহরণে দেখানো হয়েছে যে, সংগ্রহকারী কোরাউটিন বাতিল করা হয়েছে। এটি অন্তর্নিহিত প্রযোজককেও বন্ধ করে দেয়।
- প্রোডিউসার আইটেম এমিট করা শেষ করে। এই ক্ষেত্রে, ডেটার স্ট্রিম
বন্ধ হয়ে যায় এবং যে কোরাউটিন
collectকল করেছে সেটি আবার এক্সিকিউট করা শুরু করে।
অন্যান্য ইন্টারমিডিয়েট
অপারেটর নির্দিষ্ট করা না থাকলে, ফ্লো কোল্ড ও লেজি হয়। এর অর্থ হল, ফ্লোতে কোনও
টার্মিনাল অপারেটর কল করা হলে, প্রতিবারই প্রযোজক কোড এক্সিকিউট করা হয়। আগের উদাহরণে,
একাধিক ফ্লো কালেক্টর থাকার ফলে ডেটা সোর্সকে
বিভিন্ন নির্দিষ্ট সময় অন্তর একাধিকবার লেটেস্ট খবর ফেচ করতে হয়। একাধিক কনজিউমার একই সময়ে ডেটা সংগ্রহ করলে, ফ্লো অপ্টিমাইজ ও
শেয়ার করতে, shareIn
অপারেটর ব্যবহার করুন।
অপ্রত্যাশিত ব্যতিক্রম ধরা
প্রযোজকের প্রয়োগ থার্ড-পার্টি লাইব্রেরি থেকে আসতে পারে।
এর অর্থ হল, এটি অপ্রত্যাশিত ব্যতিক্রম দেখাতে পারে। এইসব
ব্যতিক্রম হ্যান্ডেল করতে, catch
ইন্টারমিডিয়েট অপারেটর ব্যবহার করুন।
class LatestNewsViewModel( private val newsRepository: NewsRepository ) : ViewModel() { init { viewModelScope.launch { newsRepository.favoriteLatestNews // Intermediate catch operator. If an exception is thrown, // catch and update the UI .catch { exception -> notifyError(exception) } .collect { favoriteNews -> // Update UI with the latest favorite news } } } }
আগের উদাহরণে, কোনও ব্যতিক্রম ঘটলে, collect
ল্যাম্বডাকে কল করা হয় না, কারণ কোনও নতুন আইটেম পাওয়া যায়নি।
এছাড়াও, catch ফ্লোতে emit আইটেম যোগ করতে পারবেন। উদাহরণস্বরূপ, রিপোজিটরি
লেয়ারের পরিবর্তে emit ক্যাশে করা ভ্যালু ব্যবহার করা যেতে পারে:
class NewsRepository( // ... ) { val favoriteLatestNews: Flow<List<ArticleHeadline>> = newsRemoteDataSource.latestNews .map { news -> news.filter { userData.isFavoriteTopic(it) } } .onEach { news -> saveInCache(news) } // If an error happens, emit the last cached values .catch { exception -> emit(lastCachedNews()) } }
এই উদাহরণে, কোনও ব্যতিক্রম ঘটলে, collect lambda-কে
কল করা হয়, কারণ ব্যতিক্রমের কারণে স্ট্রিমকে একটি নতুন আইটেম
এমনিট করা হয়েছে।
অন্য CoroutineContext-এ এক্সিকিউট করা
ডিফল্ট হিসেবে, flow বিল্ডারের প্রযোজক, এটি থেকে সংগ্রহ করা কোরাউটিনের
CoroutineContext-এ এক্সিকিউট করে এবং
আগে উল্লেখ করা হয়েছে, এটি অন্য
CoroutineContext থেকে ভ্যালু emit করতে পারে না। কিছু ক্ষেত্রে এই আচরণ অবাঞ্ছিত হতে পারে।
যেমন, এই বিষয় জুড়ে ব্যবহৃত উদাহরণে, Dispatchers.Main-এর উপর
রেপোজিটরি
লেয়ারের অপারেশন পারফর্ম করা উচিত নয়, যা viewModelScope-এর দ্বারা ব্যবহৃত হয়।
ফ্লোয়ের CoroutineContext পরিবর্তন করতে, ইন্টারমিডিয়েট অপারেটর
flowOn ব্যবহার করুন।
flowOn আপস্ট্রিম ফ্লো-এর CoroutineContext পরিবর্তন করে, অর্থাৎ
প্রযোজক ও যেকোনও মধ্যবর্তী অপারেটর আগে (বা উপরে)
flowOn প্রয়োগ করে। ডাউনস্ট্রিম ফ্লো (মধ্যবর্তী অপারেটর পরে flowOn
গ্রাহক সহ) প্রভাবিত হয় না এবং ফ্লো থেকে collect করতে ব্যবহৃত
CoroutineContext-এ এক্সিকিউট হয়। একাধিক flowOn অপারেটর থাকলে, প্রত্যেকটি অপারেটর তার
বর্তমান লোকেশন থেকে আপস্ট্রিম পরিবর্তন করে।
class NewsRepository( private val newsRemoteDataSource: NewsRemoteDataSource, private val userData: UserData, private val defaultDispatcher: CoroutineDispatcher ) { val favoriteLatestNews: Flow<List<ArticleHeadline>> = newsRemoteDataSource.latestNews .map { news -> // Executes on the default dispatcher news.filter { userData.isFavoriteTopic(it) } } .onEach { news -> // Executes on the default dispatcher saveInCache(news) } // flowOn affects the upstream flow ↑ .flowOn(defaultDispatcher) // the downstream flow ↓ is not affected .catch { exception -> // Executes in the consumer's context emit(lastCachedNews()) } }
এই কোডের মাধ্যমে, onEach ও map অপারেটররা defaultDispatcher ব্যবহার করেন,
অন্যদিকে catch অপারেটর ও কনজিউমারকে viewModelScope-এর
Dispatchers.Main-এ এক্সিকিউট করা হয়।
ডেটা সোর্স লেয়ার যেহেতু I/O কাজ করছে, তাই আপনাকে এমন ডিসপ্যাচার ব্যবহার করতে হবে যা I/O অপারেশনের জন্য অপ্টিমাইজ করা হয়েছে:
class NewsRemoteDataSource( // ... private val ioDispatcher: CoroutineDispatcher ) { val latestNews: Flow<List<ArticleHeadline>> = flow { // Executes on the IO dispatcher // ... } .flowOn(ioDispatcher) }
Jetpack লাইব্রেরিতে ফ্লো
Flow অনেক Jetpack লাইব্রেরিতে ইন্টিগ্রেট করা আছে এবং এটি Android থার্ড-পার্টি লাইব্রেরির মধ্যে জনপ্রিয়। লাইভ ডেটা আপডেট ও ডেটার অবিরাম স্ট্রিমের জন্য Flow খুব উপযুক্ত।
ডেটাবেসে কোনও পরিবর্তন হলে সেই বিষয়ে বিজ্ঞপ্তি পেতে আপনি
রুমের সাথে ফ্লো
ব্যবহার করতে পারেন। ডেটা অ্যাক্সেস অবজেক্ট (DAO)
ব্যবহার করার সময়,
লাইভ আপডেট পেতে Flow টাইপ রিটার্ন করুন।
@Dao abstract class ExampleDao { @Query("SELECT * FROM Example") abstract fun getExamples(): Flow<List<Example>> }
Example টেবিলে কোনও পরিবর্তন হলে, প্রতিবারই নতুন আইটেম সহ একটি নতুন তালিকা
ডাটাবেসে যোগ করা হয়।
কলব্যাক-ভিত্তিক API-কে ফ্লোতে কনভার্ট করা
callbackFlow
হল একটি ফ্লো বিল্ডার যা আপনাকে কলব্যাক-ভিত্তিক API-কে ফ্লোতে কনভার্ট করতে দেয়।
যেমন, Firebase Firestore
Android API কলব্যাক ব্যবহার করে।
এইসব API-কে ফ্লোতে কনভার্ট করতে এবং Firestore ডেটাবেস আপডেট শোনার জন্য, আপনি নিচের কোড ব্যবহার করতে পারেন:
class FirestoreUserEventsDataSource(
private val firestore: FirebaseFirestore
) {
// Method to get user events from the Firestore database
fun getUserEvents(): Flow<UserEvents> = callbackFlow {
// Reference to use in Firestore
var eventsCollection: CollectionReference? = null
try {
eventsCollection = FirebaseFirestore.getInstance()
.collection("collection")
.document("app")
} catch (e: Throwable) {
// If Firebase cannot be initialized, close the stream of data
// flow consumers will stop collecting and the coroutine will resume
close(e)
}
// Registers callback to firestore, which will be called on new events
val subscription = eventsCollection?.addSnapshotListener { snapshot, _ ->
if (snapshot == null) { return@addSnapshotListener }
// Sends events to the flow! Consumers will get the new events
try {
trySend(snapshot.getEvents())
} catch (e: Throwable) {
// Event couldn't be sent to the flow
}
}
// The callback inside awaitClose will be executed when the flow is
// either closed or cancelled.
// In this case, remove the callback from Firestore
awaitClose { subscription?.remove() }
}
}
flow বিল্ডারের মতো নয়, callbackFlow
send
ফাংশন ব্যবহার করে অন্য CoroutineContext থেকে ভ্যালু এমিট করার
অনুমতি দেয় অথবা trySend
ফাংশন ব্যবহার করে কোরাউটিনের বাইরে থেকে ভ্যালু এমিট করার
অনুমতি দেয়।
ইন্টার্নালি, callbackFlow একটি
চ্যানেল ব্যবহার করে,
যা একটি ব্লকিং
কিউ-এর সাথে কনসেপচুয়ালি খুব মিল আছে।
কোনও চ্যানেল ক্ষমতা অনুযায়ী কনফিগার করা হয়, এটি হল সর্বাধিক সংখ্যক এলিমেন্ট
যা বাফার করা যেতে পারে। callbackFlow-এ তৈরি করা চ্যানেলের ডিফল্ট
ক্যাপাসিটি হল ৬৪টি এলিমেন্ট। আপনি কোনও সম্পূর্ণ
চ্যানেলে নতুন এলিমেন্ট যোগ করার চেষ্টা করলে, send নতুন
এলিমেন্টের জন্য জায়গা না হওয়া পর্যন্ত প্রযোজককে সাসপেন্ড করে, অন্যদিকে trySend চ্যানেলে এলিমেন্ট যোগ করে না এবং
false সাথে সাথেই রিটার্ন করে।
trySend চ্যানেলে নির্দিষ্ট এলিমেন্টটি সঙ্গে সঙ্গে যোগ করে,
যদি এটি চ্যানেলের ধারণক্ষমতা সংক্রান্ত বিধিনিষেধ লঙ্ঘন না করে এবং তারপরে
সফল ফলাফল রিটার্ন করে।