ในโครูทีน Flow เป็นประเภทที่สามารถปล่อยค่าหลายค่า ตามลำดับ ซึ่งต่างจากฟังก์ชันระงับที่แสดงผลค่าเดียวเท่านั้น เช่น คุณสามารถใช้โฟลว์เพื่อรับข้อมูลอัปเดตแบบเรียลไทม์จากฐานข้อมูลได้
Flow สร้างขึ้นบนโครูทีนและสามารถให้ค่าได้หลายค่า
โดยแนวคิดแล้ว Flow คือสตรีมข้อมูลที่สามารถคำนวณได้
แบบไม่พร้อมกัน ค่าที่ปล่อยออกมาต้องเป็นประเภทเดียวกัน ตัวอย่างเช่น Flow<Int> คือโฟลว์ที่ปล่อยค่าจำนวนเต็ม
Flow คล้ายกับ Iterator มาก ซึ่งสร้างลำดับของค่า แต่ใช้ฟังก์ชันระงับเพื่อสร้างและใช้ค่าแบบไม่พร้อมกัน ซึ่งหมายความว่าตัวอย่างเช่น โฟลว์สามารถส่งคำขอเครือข่ายเพื่อสร้างค่าถัดไปได้อย่างปลอดภัยโดยไม่บล็อกเทรดหลัก
มีเอนทิตี 3 รายการที่เกี่ยวข้องกับสตรีมข้อมูล ได้แก่
- ผู้ผลิตสร้างข้อมูลที่จะเพิ่มลงในสตรีม Flow ยังสร้างข้อมูลแบบไม่พร้อมกันได้ด้วย Coroutine
- (ไม่บังคับ) ตัวกลางสามารถแก้ไขค่าแต่ละค่าที่ปล่อยลงใน สตรีมหรือสตรีมเองได้
- ผู้ใช้จะใช้ค่าจากสตรีม
ใน 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 จะเพิ่มองค์ประกอบที่ระบุลงในช่องทันที
ก็ต่อเมื่อการดำเนินการนี้ไม่ละเมิดข้อจำกัดด้านความจุ จากนั้นจะแสดงผลลัพธ์ที่
สำเร็จ