Flow — এটি কী, Kotlin করুটিনে cold এবং hot স্ট্রিম

লেখক: IT Sectr প্রকাশিত: 2026-03-17 পড়ার সময়: 9 মিনিট

Flow হল Kotlin Coroutines লাইব্রেরির একটি অ্যাসিঙ্ক্রোনাস ডেটা স্ট্রিম টাইপ, যা cold সিম্যান্টিক্স বাস্তবায়ন করে। Kotlin Documentation, 2025 অনুসারে, Flow map, filter, catch এবং collect অপারেটরের মাধ্যমে মানের একটি ক্রম নির্গত করতে দেয়। LiveData-এর বিপরীতে, Flow করুটিনের উপর নির্মিত এবং ব্যাকপ্রেশার সমর্থন করে।

মূল পয়েন্ট

  • Flow — Kotlin Coroutines-এ cold অ্যাসিঙ্ক্রোনাস ডেটা স্ট্রিম, সংগ্রহ না করা পর্যন্ত মান নির্গত করে না
  • Cold স্ট্রিম — প্রতিটি সাবস্ক্রাইবার শুরু থেকে নিজস্ব স্বাধীন নির্গমন শুরু করে
  • Hot স্ট্রিম (SharedFlow, StateFlow) — সাবস্ক্রাইবার নির্বিশেষে মান নির্গত করে
  • অপারেটর map, filter, catch, debounce, flatMapLatest ব্লকিং ছাড়াই স্ট্রিম রূপান্তর করে
  • Flow StateFlow এবং collectAsState()-এর মাধ্যমে Jetpack Compose-এর সাথে সম্পূর্ণ সামঞ্জস্যপূর্ণ

Kotlin-এ Flow কী?

Flow হল kotlinx.coroutines.flow প্যাকেজের একটি টাইপ, যা একটি cold অ্যাসিঙ্ক্রোনাস ডেটা স্ট্রিম উপস্থাপন করে। এর মূলে, Flow হল একটি করুটিন ক্রম যা emit() ফাংশনের মাধ্যমে মান নির্গত করে এবং সফলভাবে বা ব্যতিক্রমের সাথে শেষ হয়। স্ট্রিম সংগ্রহ টার্মিনাল অপারেটর collect()-এর মাধ্যমে করা হয়, যা একটি suspend ফাংশন।

Cold সিম্যান্টিক্স

Cold স্ট্রিম মানে হল flow বিল্ডারের ভিতরের কোড প্রতিটি সাবস্ক্রাইবারের জন্য নতুন করে চলে। RxJava-তে Observable.fromIterable একইভাবে আচরণ করে: একটি নতুন সাবস্ক্রাইবার শুরু থেকে সমস্ত মান পায়। Flow-এ, এটি suspend ফাংশন collect-এর মাধ্যমে বাস্তবায়িত, যা ডেটা সংগ্রহের পুরো সময়ের জন্য করুটিনকে ব্লক করে।

Flow বিল্ডার

Kotlin Flow তৈরি করার বেশ কয়েকটি উপায় প্রদান করে: flow { } — emit()-সহ মৌলিক নির্মাণ, flowOf(vararg values) — মানের একটি নির্দিষ্ট সেটের জন্য, .asFlow() — সংগ্রহ এবং Sequence-এর জন্য এক্সটেনশন। সমস্ত বিল্ডার cold — টার্মিনাল অপারেটর কল করলেই কেবল ডেটা উৎপন্ন হয়।

Cold এবং Hot স্ট্রিম

Cold এবং hot স্ট্রিমে বিভাজন রিঅ্যাকটিভ প্রোগ্রামিংয়ের একটি মূল ধারণা। Cold স্ট্রিম (Flow, Observable) সাবস্ক্রিপশনে ডেটা জেনারেশন শুরু করে। Hot স্ট্রিম (Channel, SharedFlow) স্বাধীনভাবে ডেটা নির্গত করে — সাবস্ক্রাইবার কেবল সাবস্ক্রিপশনের পরে যা ঘটে তা পায়, ক্রমের শুরু ছাড়া।

SharedFlow হল একটি hot Flow যার একাধিক সাবস্ক্রাইবার থাকতে পারে এবং replay কনফিগার করা থাকলে সাম্প্রতিক মানগুলি পুনরায় চালাতে পারে। SharedFlow ইভেন্টের জন্য উপযুক্ত (একবারের বিজ্ঞপ্তি)। StateFlow হল এর একটি বৈকল্পিক যার একটি নির্দিষ্ট অবস্থার মান রয়েছে, নতুন সাবস্ক্রাইবারের জন্য সর্বশেষ মান ক্যাশ করে।

ChannelFlow ভিতরে Channel ব্যবহার করে, Flow এবং Channel-এর বৈশিষ্ট্য একত্রিত করে। এটি capacity-র মাধ্যমে বাফারিং এবং ব্যাকপ্রেশার সমর্থন করে। ChannelFlow কলব্যাক API-কে রিঅ্যাকটিভ স্ট্রিমে রূপান্তর করার সময় দরকারী, যেখানে মান বিভিন্ন করুটিন থেকে নির্গত হয়।

Cold এবং hot-এর মধ্যে রূপান্তর

cold Flow-কে hot SharedFlow-এ রূপান্তর করতে shareIn(scope, started, replay) অপারেটর ব্যবহার করা হয়। প্যারামিটার started শুরুর মুহূর্ত নিয়ন্ত্রণ করে: SharingStarted.WhileSubscribed() — যতক্ষণ সাবস্ক্রাইবার আছে ততক্ষণ সক্রিয়, Lazily — প্রথম সাবস্ক্রাইবারে শুরু, Eagerly — তাৎক্ষণিক শুরু। বিপরীত রূপান্তর — hot থেকে cold: StateFlow.asFlow() একটি cold Flow ফেরত দেয় যা collect-এ StateFlow-এর বর্তমান মান নির্গত করে। এটি পরীক্ষণের জন্য সুবিধাজনক।

Flow অপারেটর

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 ফাংশনের ভিতরে কল করতে হবে।

Flow কোড উদাহরণ

প্রথম উদাহরণ — map অপারেটরের মাধ্যমে রূপান্তর সহ সংখ্যা উৎপন্নকারী মৌলিক Flow:

kotlin
val numberFlow = flow {
    for (i in 1..5) {
        delay(500)
        emit(i)
    }
}

scope.launch {
    numberFlow
        .map { "সংখ্যা: $it" }
        .collect { value ->
            println(value)
        }
}

দ্বিতীয় উদাহরণ — catch-এর মাধ্যমে ফিল্টারিং এবং ত্রুটি ব্যবস্থাপনা সহ স্ট্রিম রূপান্তর:

kotlin
flow {
    emit("data1")
    emit("data2")
    throw RuntimeException("network error")
}
    .catch { e ->
        emit("fallback_data")
    }
    .collect { value ->
        println(value)
    }

তৃতীয় উদাহরণ — Jetpack Compose-এ রিঅ্যাকটিভ UI-এর জন্য ViewModel-এ StateFlow ব্যবহার:

kotlin
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 এবং SharedFlow

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-এ ত্রুটি ব্যবস্থাপনার একটি বিশেষত্ব রয়েছে: যদি টার্মিনাল অপারেটরের আগে 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-র মধ্যে পার্থক্য কী?

Flow হল করুটিন, অপারেটর এবং ব্যাকপ্রেশার সমর্থন-সহ একটি অ্যাসিঙ্ক্রোনাস স্ট্রিম, যা যেকোনো আর্কিটেকচার স্তরে কাজ করে। LiveData কেবল UI স্তরের জন্য একটি lifecycle-aware উপাদান। Google ব্যবসায়িক যুক্তি এবং রিপোজিটরির জন্য Flow সুপারিশ করে, ViewModel-এ সহজ পর্যবেক্ষণের জন্য LiveData।

কখন SharedFlow-এর পরিবর্তে StateFlow ব্যবহার করবেন?

StateFlow — যখন UI অবস্থা (টাস্ক তালিকা, অনুসন্ধান টেক্সট, লোডিং ফ্ল্যাগ) সংরক্ষণ করতে হবে — প্রতিটি সাবস্ক্রাইবার বর্তমান মান পায়। SharedFlow — একবারের ইভেন্টের জন্য (নেভিগেশন, Snackbar)। StateFlow ইভেন্টের জন্য ব্যবহার করা উচিত নয় কারণ নতুন মান পুনরায় প্রক্রিয়া করা যেতে পারে।

Flow-এ ব্যাকপ্রেশার কীভাবে কাজ করে?

Flow-এ, ব্যাকপ্রেশার suspend প্রক্রিয়ার মাধ্যমে বাস্তবায়িত: যদি কালেক্টর পূর্ববর্তী মান প্রক্রিয়া করছে তবে emit() করুটিনকে স্থগিত করে। ChannelFlow-এ চ্যানেলের (Channel) capacity আকারের বাফার থাকে। ওভারফ্লোতে: suspending (অপেক্ষা), drop (বাদ দেওয়া) বা conflate (সর্বশেষ দিয়ে প্রতিস্থাপন)।

কিভাবে কলব্যাককে Flow-এ রূপান্তর করবেন?

callbackFlow ব্যবহার করুন — কলব্যাক API-র জন্য Flow বিল্ডার। ভিতরে, কলব্যাকের ভিতরে emit(value)-সহ registerCallback() কল করুন। awaitClose করুটিন বাতিল হলে unregisterCallback() কল নিশ্চিত করে। callbackFlow ভিতরে Channel(UNLIMITED)-এর মাধ্যমে বাফারিং সমর্থন করে।

Flow কি RxJava-র সাথে ব্যবহার করা যায়?

হ্যাঁ, কনভার্টারের মাধ্যমে: Flow.asObservable() kotlinx-coroutines-rx3 প্যাকেজ থেকে Flow-কে RxJava 3 Observable-এ রূপান্তর করে। বিপরীত — CompletableSource.asFlow() Single/Completable/Maybe-র জন্য। এটি বড় প্রকল্পে RxJava থেকে করুটিনে মাইগ্রেট করার সময় দরকারী।

সারাংশ

  • Flow — suspend ফাংশন collect-সহ Kotlin Coroutines-এ cold অ্যাসিঙ্ক্রোনাস ডেটা স্ট্রিম
  • Cold স্ট্রিম প্রতিটি সাবস্ক্রাইবারের জন্য পুনরায় নির্গমন শুরু করে
  • StateFlow — সর্বশেষ মান ক্যাশিং-সহ hot অবস্থা পাত্র
  • SharedFlow — replay এবং বাফার কনফিগারেশন-সহ ইভেন্টের জন্য hot স্ট্রিম
  • অপারেটর map, filter, debounce, catch, flatMapLatest — স্ট্রিম রূপান্তরের ভিত্তি
  • Google আধুনিক Android আর্কিটেকচারে Flow-কে প্রাথমিক ডেটা উৎস হিসাবে সুপারিশ করে
  • LiveData কেবল UI স্তরের জন্য উপযুক্ত, Flow — অ্যাপ্লিকেশনের সমস্ত স্তরের জন্য

আমরা একটি মোবাইল অ্যাপ্লিকেশন টার্নকি তৈরি করব

IT Sectr 2017 সাল থেকে স্টার্টআপ এবং ব্যবসার জন্য iOS এবং Android অ্যাপ্লিকেশন তৈরি করে। আমরা আপনাকে পরামর্শ দেব এবং সেরা সমাধান প্রস্তাব করব।

প্রকল্প নিয়ে আলোচনা করুন

আরও পড়ুন