Android-এ Kotlin flows

করুটিনে, ফ্লো হল এমন একটি ধরন যা একাধিক ভ্যালু ধারাবাহিকভাবে এমিট করতে পারে, এটি সাসপেন্ড ফাংশনের বিপরীত, যা শুধুমাত্র একটি ভ্যালু রিটার্ন করে। যেমন, আপনি কোনও ডেটাবেস থেকে লাইভ আপডেট পেতে ফ্লো ব্যবহার করতে পারেন।

ফ্লো কোরাউটিনের উপরে তৈরি করা হয় এবং একাধিক ভ্যালু প্রদান করতে পারে। ফ্লো হল কনসেপচুয়ালি ডেটার স্ট্রিম যা অ্যাসিঙ্ক্রোনাস কম্পিউট করা যায়। ইমিট করা মান একই ধরনের হতে হবে। যেমন, 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 চ্যানেলে নির্দিষ্ট এলিমেন্টটি সঙ্গে সঙ্গে যোগ করে, যদি এটি চ্যানেলের ধারণক্ষমতা সংক্রান্ত বিধিনিষেধ লঙ্ঘন না করে এবং তারপরে সফল ফলাফল রিটার্ন করে।

ফ্লো সংক্রান্ত অতিরিক্ত রিসোর্স