Kotlin ทำงานอย่างลื่นไหลใน Android

ในโครูทีน Flow เป็นประเภทที่สามารถปล่อยค่าหลายค่า ตามลำดับ ซึ่งต่างจากฟังก์ชันระงับที่แสดงผลค่าเดียวเท่านั้น เช่น คุณสามารถใช้โฟลว์เพื่อรับข้อมูลอัปเดตแบบเรียลไทม์จากฐานข้อมูลได้

Flow สร้างขึ้นบนโครูทีนและสามารถให้ค่าได้หลายค่า โดยแนวคิดแล้ว Flow คือสตรีมข้อมูลที่สามารถคำนวณได้ แบบไม่พร้อมกัน ค่าที่ปล่อยออกมาต้องเป็นประเภทเดียวกัน ตัวอย่างเช่น Flow<Int> คือโฟลว์ที่ปล่อยค่าจำนวนเต็ม

Flow คล้ายกับ Iterator มาก ซึ่งสร้างลำดับของค่า แต่ใช้ฟังก์ชันระงับเพื่อสร้างและใช้ค่าแบบไม่พร้อมกัน ซึ่งหมายความว่าตัวอย่างเช่น โฟลว์สามารถส่งคำขอเครือข่ายเพื่อสร้างค่าถัดไปได้อย่างปลอดภัยโดยไม่บล็อกเทรดหลัก

มีเอนทิตี 3 รายการที่เกี่ยวข้องกับสตรีมข้อมูล ได้แก่

  • ผู้ผลิตสร้างข้อมูลที่จะเพิ่มลงในสตรีม Flow ยังสร้างข้อมูลแบบไม่พร้อมกันได้ด้วย Coroutine
  • (ไม่บังคับ) ตัวกลางสามารถแก้ไขค่าแต่ละค่าที่ปล่อยลงใน สตรีมหรือสตรีมเองได้
  • ผู้ใช้จะใช้ค่าจากสตรีม

เอนทิตีที่เกี่ยวข้องในสตรีมข้อมูล ผู้บริโภค ตัวกลาง และผู้ผลิต (ไม่บังคับ)
              ตัวกลาง และผู้ผลิต
รูปที่ 1 เอนทิตีที่เกี่ยวข้องในสตรีมข้อมูล ผู้บริโภค ตัวกลาง (ไม่บังคับ) และผู้ผลิต

ใน 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 แบบไม่พร้อมกันเดียวกัน แต่มีข้อจำกัดบางประการดังนี้

  • โฟลว์เป็นแบบตามลำดับ เนื่องจากโปรดิวเซอร์อยู่ในโครูทีน เมื่อเรียกใช้ฟังก์ชันระงับ โปรดิวเซอร์จะระงับจนกว่าฟังก์ชันระงับจะส่งคืน ในตัวอย่างนี้ โปรดิวเซอร์จะระงับจนกว่าfetchLatestNews คำขอเครือข่ายจะเสร็จสมบูรณ์ จากนั้นระบบจะส่งผลลัพธ์ไปยังสตรีม
  • ด้วยเครื่องมือสร้าง flow ผู้ผลิตจึงไม่สามารถ emit ค่าจาก CoroutineContext อื่น ดังนั้น อย่าเรียกใช้ emit ในCoroutineContextอื่นโดยการสร้างโครูทีนใหม่หรือใช้บล็อกโค้ด withContext คุณสามารถใช้เครื่องมือสร้างโฟลว์อื่นๆ เช่น callbackFlow ในกรณีเหล่านี้ได้

การแก้ไขสตรีม

ตัวกลางสามารถใช้ตัวดำเนินการกลางเพื่อแก้ไขสตรีมของ ข้อมูลโดยไม่ต้องใช้ค่า ตัวดำเนินการเหล่านี้เป็นฟังก์ชันที่เมื่อ ใช้กับสตรีมข้อมูล จะตั้งค่าเชนของการดำเนินการที่ จะไม่ดำเนินการจนกว่าจะมีการใช้ค่าในอนาคต ดูข้อมูลเพิ่มเติมเกี่ยวกับ โอเปอเรเตอร์ระดับกลางได้ใน เอกสารอ้างอิง Flow

ในตัวอย่างด้านล่าง เลเยอร์ที่เก็บจะใช้อันดับกลาง 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 เป็นฟังก์ชันระงับ จึงต้องเรียกใช้ภายใน โครูทีน โดยจะใช้ lambda เป็นพารามิเตอร์ที่เรียกใช้กับค่าใหม่ทุกค่า เนื่องจากเป็นฟังก์ชันระงับ ดังนั้น โครูทีนที่เรียก 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
            }
        }
    }
}

การรวบรวมโฟลว์จะทริกเกอร์โปรดิวเซอร์ที่รีเฟรชข่าวสารล่าสุด และปล่อยผลลัพธ์ของคำขอเครือข่ายในช่วงเวลาที่กำหนด เนื่องจาก Producer จะยังคงใช้งานได้เสมอด้วยลูป while(true) สตรีม ของข้อมูลจะปิดเมื่อ ViewModel ถูกล้างและ viewModelScope ถูกยกเลิก

การรวบรวมโฟลว์อาจหยุดลงด้วยเหตุผลต่อไปนี้

  • ระบบจะยกเลิกโครูทีนที่รวบรวม ดังที่แสดงในตัวอย่างก่อนหน้า ซึ่งจะหยุดผู้ผลิตที่เกี่ยวข้องด้วย
  • ผู้ผลิตส่งรายการเสร็จสิ้น ในกรณีนี้ สตรีมข้อมูล จะปิดและโครูทีนที่เรียก collect จะกลับมาดำเนินการต่อ

Flow จะเป็นแบบเย็นและขี้เกียจ เว้นแต่จะระบุด้วยตัวดำเนินการขั้นกลางอื่นๆ ซึ่งหมายความว่าระบบจะเรียกใช้รหัสผู้ผลิตทุกครั้งที่มีการเรียกใช้ ผู้ให้บริการเทอร์มินัลในโฟลว์ ในตัวอย่างก่อนหน้า การมีตัวรวบรวมโฟลว์หลายตัวจะทำให้แหล่งข้อมูลดึงข้อมูล ข่าวสารล่าสุดหลายครั้งในช่วงเวลาที่กำหนดต่างๆ หากต้องการเพิ่มประสิทธิภาพและ แชร์โฟลว์เมื่อผู้บริโภคหลายรายรวบรวมข้อมูลพร้อมกัน ให้ใช้ตัวดำเนินการ shareIn

การตรวจหาข้อยกเว้นที่ไม่คาดคิด

การติดตั้งใช้งาน Producer อาจมาจากไลบรารีของบุคคลที่สาม ซึ่งหมายความว่าฟังก์ชันนี้อาจทำให้เกิดข้อยกเว้นที่ไม่คาดคิด หากต้องการจัดการข้อยกเว้นเหล่านี้ ให้ใช้ตัวดำเนินการกลาง 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 Lambda เนื่องจากยังไม่ได้รับสินค้าใหม่

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 ของโครูทีนที่รวบรวมจากโปรดิวเซอร์นั้น และตามที่ กล่าวไว้ก่อนหน้านี้ โปรดิวเซอร์จะemitค่าจาก CoroutineContext อื่นไม่ได้ ลักษณะการทำงานนี้อาจไม่เป็นที่ต้องการในบางกรณี เช่น ในตัวอย่างที่ใช้ตลอดทั้งหัวข้อนี้ เลเยอร์ที่เก็บไม่ควรดำเนินการกับ 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 ใช้

เนื่องจากเลเยอร์แหล่งข้อมูลกำลังทำงาน I/O คุณจึงควรใช้ Dispatcher ที่เพิ่มประสิทธิภาพสำหรับการดำเนินการ I/O

class NewsRemoteDataSource(
    // ...
    private val ioDispatcher: CoroutineDispatcher
) {

    val latestNews: Flow<List<ArticleHeadline>> = flow {
        // Executes on the IO dispatcher
        // ...
    }
        .flowOn(ioDispatcher)
}

Flow ในไลบรารี Jetpack

Flow ผสานรวมอยู่ในไลบรารี Jetpack หลายรายการ และเป็นที่นิยมในหมู่ไลบรารีของบุคคลที่สามใน Android Flow เหมาะสำหรับการอัปเดตข้อมูลแบบเรียลไทม์ และสตรีมข้อมูลที่ไม่มีที่สิ้นสุด

คุณใช้ Flow with Room เพื่อรับการแจ้งเตือนเกี่ยวกับการเปลี่ยนแปลงในฐานข้อมูลได้ เมื่อใช้ออบเจ็กต์การเข้าถึงข้อมูล (DAO) ให้ส่งคืนประเภท Flow เพื่อรับข้อมูลอัปเดตแบบเรียลไทม์

@Dao
abstract class ExampleDao {
    @Query("SELECT * FROM Example")
    abstract fun getExamples(): Flow<List<Example>>
}

ทุกครั้งที่มีการเปลี่ยนแปลงในExample ตาราง ระบบจะสร้างรายการใหม่ พร้อมกับรายการใหม่ในฐานข้อมูล

แปลง API ที่อิงตามการเรียกกลับเป็น Flow

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 Builder callbackFlow ต่างจาก ตรงที่อนุญาตให้ส่งค่าจาก CoroutineContext อื่นด้วยฟังก์ชัน send หรือภายนอกโครูทีนด้วยฟังก์ชัน trySend

ภายใน callbackFlow จะใช้แชแนล ซึ่งมีแนวคิดคล้ายกับคิวการบล็อกเป็นอย่างมาก ช่องจะได้รับการกำหนดค่าด้วยความจุ ซึ่งเป็นจำนวนสูงสุดขององค์ประกอบ ที่บัฟเฟอร์ได้ แชแนลที่สร้างใน callbackFlow มีความจุเริ่มต้น 64 องค์ประกอบ เมื่อคุณพยายามเพิ่มองค์ประกอบใหม่ลงในช่องที่เต็มแล้ว send จะระงับโปรดิวเซอร์จนกว่าจะมีพื้นที่สำหรับองค์ประกอบใหม่ ส่วน trySend จะไม่เพิ่มองค์ประกอบลงในช่องและจะแสดงข้อผิดพลาด false ทันที

trySend จะเพิ่มองค์ประกอบที่ระบุลงในช่องทันที ก็ต่อเมื่อการดำเนินการนี้ไม่ละเมิดข้อจำกัดด้านความจุ จากนั้นจะแสดงผลลัพธ์ที่ สำเร็จ

แหล่งข้อมูลเพิ่มเติมเกี่ยวกับโฟลว์