جریان‌های Kotlin در Android

در روتین‌های همکار، جریان نوعی است که می‌تواند چندین مقدار را به‌صورت متوالی منتشر کند، برخلاف توابع تعلیق که فقط یک مقدار را برمی‌گردانند. برای مثال، می‌توانید از یک جریان برای دریافت به‌روزرسانی‌های زنده از پایگاه داده استفاده کنید.

جریان‌ها براساس روتین‌های همکار ساخته می‌شوند و می‌توانند مقادیر متعددی ارائه دهند. جریان ازنظر مفهومی جریانی از داده‌ها است که می‌تواند به‌صورت ناهمزمان محاسبه شود. مقادیر منتشرشده باید از یک نوع باشند. برای مثال، 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 عنصر مشخص‌شده را بلافاصله به کانال اضافه می‌کند، فقط درصورتی‌که این کار محدودیت‌های ظرفیت آن را نقض نکند، و سپس نتیجه موفقیت‌آمیز را برمی‌گرداند.

منابع جریان اضافی