در روتینهای همکار، جریان نوعی است که میتواند چندین مقدار را بهصورت متوالی منتشر کند، برخلاف توابع تعلیق که فقط یک مقدار را برمیگردانند. برای مثال، میتوانید از یک جریان برای دریافت بهروزرسانیهای زنده از پایگاه داده استفاده کنید.
جریانها براساس روتینهای همکار ساخته میشوند و میتوانند مقادیر متعددی ارائه دهند.
جریان ازنظر مفهومی جریانی از دادهها است که میتواند بهصورت ناهمزمان محاسبه شود. مقادیر منتشرشده باید از یک نوع باشند. برای
مثال، Flow<Int> جریانی است که مقادیر صحیح را منتشر میکند.
جریان بسیار شبیه به Iterator است که دنبالهای از مقادیر تولید میکند، اما از توابع تعلیق برای تولید و مصرف مقادیر بهصورت ناهمزمان استفاده میکند. این یعنی، برای مثال، جریان میتواند بدون مسدود کردن رشته اصلی، درخواست شبکه را بهطور ایمن برای تولید مقدار بعدی انجام دهد.
سه نهاد در جاریسازی دادهها دخیل هستند:
- تولیدکننده دادههایی تولید میکند که به جاریسازی اضافه میشود. بهلطف روالهای همکار، جاریسازیها میتوانند دادهها را بهصورت ناهمزمان نیز تولید کنند.
- (اختیاری) واسطهها میتوانند هر مقدار منتشرشده در جاریسازی یا خود جاریسازی را تغییر دهند.
- مصرفکننده مقادیر را از جاریسازی مصرف میکند.
در Android، یک مخزن معمولاً دادههای رابط کاربری را از پایگاههای داده محلی (Room) یا زیرینههای ابری (مثل Firebase) با رابط کاربری که بهعنوان مصرفکننده دادهها عمل میکند جاریسازی میکند. برعکس، میانای کاربر میتواند رویدادهای ورودی کاربر را برای لایههای دیگر تولید کند تا مصرف کنند. لایههای میانی (مثل نمونههای
ViewModel) جاریسازی را تبدیل میکنند تا با الزامات هر لایه مطابقت داشته باشد.
درحال ایجاد جریان
برای ایجاد جریانهای کار، از
سازنده جریان کار
میاناهای برنامهسازی کاربردی استفاده کنید. تابع سازنده 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 در یک روال همکار اجرا میشود. بنابراین، از همان میاناهای برنامهسازی کاربردی ناهمزمان بهره میبرد، اما برخی محدودیتها اعمال میشود:
- جریانها ترتیبی هستند. ازآنجاییکه تولیدکننده در یک روتین همکار است، هنگام فراخوانی
یک تابع تعلیق، تولیدکننده تا زمانی که تابع تعلیق
برگردد تعلیق میشود. در این مثال، تولیدکننده تا زمان تکمیل درخواست شبکه
fetchLatestNewsتعلیق میشود. فقط در این صورت نتیجه به جاریسازی ارسال میشود. - با سازنده
flow، تهیهکننده نمیتواند مقادیرemitرا ازCoroutineContextمتفاوتی دریافت کند. بنابراین، با ایجاد روتینهای همزمان جدید یا بااستفاده از بلوکهای کدwithContext، درCoroutineContextمتفاوتemitرا فراخوانی نکنید. در این موارد میتوانید از سازندگان جریان دیگر مثلcallbackFlowاستفاده کنید.
درحال اصلاح جاریسازی
واسطهها میتوانند از عملگرهای واسطه برای اصلاح جاریسازی دادهها بدون مصرف مقادیر استفاده کنند. این عملگرها توابعی هستند که وقتی روی جاریسازی دادهها اعمال میشوند، زنجیرهای از عملیات راهاندازی میکنند که تا زمانی که مقادیر در آینده مصرف نشوند اجرا نمیشوند. در اسناد مرجع «جریان» درباره کاربران واسطه بیشتر بدانید.
در مثال زیر، لایه مخزن از عامل واسطه
map
برای تبدیل دادهها بهمنظور نمایش در View استفاده میکند:
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 فراخوانده میشود، زیرا بهدلیل استثنا، عنصر جدیدی به جاریسازی ارسال شده است.
اجرا در CoroutineContext متفاوت
بهطور پیشفرض، تولیدکننده سازنده flow در CoroutineContext روتین فرعی که از آن جمعآوری میکند اجرا میشود، و همانطور که قبلاً ذکر شد، نمیتواند مقادیر را از CoroutineContext دیگری emit کند. این رفتار ممکن است در برخی موارد نامطلوب باشد.
برای مثال، در نمونههایی که در سراسر این موضوع استفاده شده است، لایه مخزن نباید عملیاتی را روی Dispatchers.Main انجام دهد که viewModelScope از آن استفاده میکند.
برای تغییر CoroutineContext جاریسازی، از عملگر میانی
flowOn استفاده کنید.
flowOn CoroutineContext جریان بالادستی را تغییر میدهد، یعنی
تولیدکننده و هرگونه عامل واسطهای که قبلاز (یا بالای)
flowOn اعمال شده است. جریان پاییندستی (عملگرهای واسطه بعداز flowOn
همراه با مصرفکننده) تحتتأثیر قرار نمیگیرد و در CoroutineContext استفادهشده برای collect از جریان اجرا میشود. اگر چندین اپراتور 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 و مصرفکننده در
Dispatchers.Main که viewModelScope استفاده میکند اجرا میشوند.
ازآنجاییکه لایه منبع داده کار ورودی/خروجی انجام میدهد، باید از توزیعکنندهای استفاده کنید که برای عملیات ورودی/خروجی بهینهسازی شده است:
class NewsRemoteDataSource( // ... private val ioDispatcher: CoroutineDispatcher ) { val latestNews: Flow<List<ArticleHeadline>> = flow { // Executes on the IO dispatcher // ... } .flowOn(ioDispatcher) }
جریانها در کتابخانههای Jetpack و منابع داده
«جریان» در بسیاری از کتابخانههای Jetpack ادغام شده است و در کتابخانههای طرف سوم Android و کیتهای توسعه نرمافزار ابری محبوب است. Flow برای بهروزرسانیهای داده زنده و جاریسازیهای بیپایان داده بسیار مناسب است.
میتوانید از
جریان با اتاق
برای دریافت اعلان تغییرات در پایگاه داده استفاده کنید. هنگام استفاده از
اشیاء دسترسی به داده (DAO)،
برای دریافت بهروزرسانیهای زنده، نوع Flow را برگردانید.
@Dao abstract class ExampleDao { @Query("SELECT * FROM Example") abstract fun getExamples(): Flow<List<Example>> }
هر بار که تغییری در جدول Example ایجاد میشود، فهرست جدیدی با موارد جدید در پایگاه داده منتشر میشود.
بههمین ترتیب، منابع داده ابری همزمان مانند Cloud Firestore افزونههایی
Flow ارائه میدهند که بهروزرسانیهای داده زنده را مستقیماً به لایه مخزن شما جاریسازی میکنند:
class RemoteMessagesDataSource(
private val firestore: FirebaseFirestore
) {
// Returns a Flow emitting updated message lists when collection changes
val messages: Flow<List<Message>> =
firestore.collection("messages").dataObjects()
}
این ویژگی بهطور طبیعی با فضای ذخیرهسازی محلی جفت میشود: مخزن شما میتواند جاریسازیها را از هر دو Room (حافظه نهان محلی) و Firebase (دادههای ابری از دور) ترکیب کند تا تجربه کاربری بیدرنگ و اول آفلاین را ارائه دهد.
تبدیل میاناهای برنامهسازی کاربردی مبتنی بر برگشتی به جریان
callbackFlow
سازنده جریانی است که به شما امکان میدهد میاناهای برنامهسازی کاربردی مبتنی بر برگشت تماس را به جریان تبدیل کنید. این ویژگی بهویژه هنگام کار با شنوندگان چند رویدادی در کیتهای توسعه نرمافزار مانند Cloud Firestore یا پیکربندی از دور Firebase مفید است.
برای تبدیل میانای برنامهسازی کاربردی شنونده مبتنی بر برگشتی به جاریسازی—برای مثال، گوش دادن به بهروزرسانیهای پایگاه داده 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 عنصر مشخصشده را بلافاصله به کانال اضافه میکند،
فقط درصورتیکه این کار محدودیتهای ظرفیت آن را نقض نکند، و سپس نتیجه
موفقیتآمیز را برمیگرداند.