Kotlin Academy · บทเรียน

flowOn และ buffer: บริบทและแรงดันย้อนกลับ

เปลี่ยนบริบทการปล่อยค่าด้วย flowOn และพักค่าที่ปล่อยด้วย buffer เพื่อรองรับแรงดันย้อนกลับ

บทเรียน 4 จาก 413 ขั้นตอน

flowOn และ buffer: บริบทและแรงดันย้อนกลับ เป็นบทเรียน Kotlin Academy ฟรีบน CoddyKit นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน คุณสามารถอ่านบทเรียนทั้งหมดด้านล่างฟรี — จากนั้นลองปฏิบัติด้วยตัวคุณเองในเบราว์เซอร์พร้อมตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7 บทเรียนนี้เป็นส่วนหนึ่งของเส้นทางการเรียน Kotlin Academy และความก้าวหน้าของคุณจะซิงค์ข้ามเว็บและแอป CoddyKit คอร์ส Kotlin Academy มีบทเรียนทั้งหมด 4 บทเรียน

บริบทของโฟลว์

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

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    flow {
        println("emit on: ${Thread.currentThread().name}")
        emit(1)
    }.collect {
        println("collect on: ${Thread.currentThread().name}")
    }
}

flowOn เปลี่ยนบริบทต้นทาง

flowOn(dispatcher) ทำให้โฟลว์ต้นทาง ซึ่งหมายถึงทุกอย่างที่อยู่เหนือโอเปอเรเตอร์นี้ในสายการทำงาน ทำงานบนตัวจัดส่งงานที่ระบุ ขณะที่การรวบรวมยังคงทำงานบนตัวจัดส่งงานของผู้เรียก

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    flow {
        println("emit: ${Thread.currentThread().name}")
        emit(1)
    }.map {
        println("map: ${Thread.currentThread().name}")
        it * 2
    }.flowOn(Dispatchers.Default)  // above runs on Default
    .collect {
        println("collect: ${Thread.currentThread().name}")
    }
}

ใช้ flowOn หลายครั้งในสายการทำงาน

คุณสามารถใช้ flowOn ได้หลายครั้ง แต่ละครั้งจะมีผลต่อโอเปอเรเตอร์ที่อยู่เหนือโอเปอเรเตอร์นั้นโดยตรง จนถึง flowOn ก่อนหน้า

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    flow { emit(1) }
        .map { it + 1 }.flowOn(Dispatchers.IO)      // map runs on IO
        .map { it * 2 }.flowOn(Dispatchers.Default)  // this map runs on Default
        .collect { println("Result: $it") }           // collect on Main (runBlocking)
}

ปัญหาแรงดันย้อนกลับ

เมื่อผู้ผลิตส่งค่าเร็วกว่าที่ผู้รวบรวมจะประมวลผลได้ ค่าเหล่านั้นจะถูกต่อคิว หากไม่มีการบัฟเฟอร์ ผู้ผลิตจะต้องรอ

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val time = System.currentTimeMillis()
    flow {
        repeat(3) { i ->
            delay(100)  // fast producer
            emit(i)
        }
    }.collect {
        delay(300)      // slow consumer
        println("Got $it in ${System.currentTimeMillis() - time}ms")
    }
}

โอเปอเรเตอร์ buffer()

buffer() ทำให้ผู้ผลิตและผู้ใช้ทำงานพร้อมกันในโครูทีนแยกกัน โดยบัฟเฟอร์ค่าที่ส่งออกไว้ในแชนเนล ผู้ผลิตจึงไม่ต้องรอผู้ใช้

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    val time = System.currentTimeMillis()
    flow {
        repeat(3) { i -> delay(100); emit(i) }
    }.buffer()  // producer and consumer run concurrently
    .collect {
        delay(300)
        println("Got $it in ${System.currentTimeMillis() - time}ms")
    }
}

ความจุของ buffer

buffer(capacity) กำหนดขนาดบัฟเฟอร์ของแชนเนล เมื่อบัฟเฟอร์เต็ม ผู้ผลิตจะถูกพักการทำงาน ซึ่งก่อให้เกิดแรงดันย้อนกลับ ความจุเริ่มต้นคือ 64

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
import kotlinx.coroutines.channels.Channel
fun main() = runBlocking {
    flow { repeat(5) { emit(it) } }
        .buffer(Channel.RENDEZVOUS)    // 0: producer waits
        // .buffer(Channel.BUFFERED)   // default: 64
        // .buffer(Channel.UNLIMITED)  // unbounded
        .collect { delay(50); println(it) }
}

conflate() สำหรับเก็บเฉพาะค่าล่าสุด

conflate() ทิ้งค่าระหว่างกลางเมื่อผู้รวบรวมทำงานช้า โดยเก็บไว้เฉพาะค่าที่ส่งออกล่าสุด เหมาะสำหรับสถานะของส่วนติดต่อผู้ใช้

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    flow {
        repeat(5) { i -> emit(i); delay(50) }
    }.conflate()
    .collect { i ->
        delay(150)
        println("Collected: $i")  // skips some values
    }
}

collectLatest สำหรับผู้รวบรวมที่ทำงานช้า

collectLatest ยกเลิกบล็อกการรวบรวมปัจจุบันเมื่อมีค่าใหม่มาถึง แล้วเริ่มทำงานใหม่ด้วยค่าล่าสุด

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun main() = runBlocking {
    flow {
        emit(1); delay(50)
        emit(2); delay(50)
        emit(3)
    }.collectLatest { value ->
        println("Processing $value")
        delay(100)  // gets cancelled if new value arrives
        println("Done $value")  // only prints for last value
    }
}

รูปแบบ flowOn + buffer

ใช้ flowOn และ buffer ร่วมกัน โดยให้ผู้ผลิตทำงานบน IO เช่น เครือข่ายหรือดิสก์ บัฟเฟอร์ผลลัพธ์ แล้วรวบรวมบนตัวจัดส่งงานหลัก

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun fetchItems(): Flow<String> = flow {
    repeat(3) { i ->
        delay(100)  // simulate IO
        emit("item-$i")
    }
}.flowOn(Dispatchers.IO).buffer(10)
fun main() = runBlocking {
    fetchItems().collect { println("UI: $it") }
}

channelFlow สำหรับผู้ผลิตที่ทำงานพร้อมกัน

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

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
fun concurrentFlow(): Flow<Int> = channelFlow {
    launch { send(1) }
    launch { send(2) }
    launch { send(3) }
}
fun main() = runBlocking {
    concurrentFlow().collect { println(it) }
}

การเลือกกลยุทธ์ที่เหมาะสม

สรุป: ใช้ flowOn สำหรับการเปลี่ยนบริบท ใช้ buffer เพื่อเพิ่มอัตราการประมวลผล ใช้ conflate เพื่อเก็บเฉพาะค่าล่าสุดของส่วนติดต่อผู้ใช้ และใช้ collectLatest เพื่อยกเลิกการประมวลผลที่ล้าสมัย

import kotlinx.coroutines.*
import kotlinx.coroutines.flow.*
// Guidelines:
// CPU-heavy production  -> flowOn(Dispatchers.Default)
// IO-heavy production   -> flowOn(Dispatchers.IO)
// Slow consumer         -> buffer()
// UI state updates      -> conflate() or StateFlow
// Search/autocomplete   -> collectLatest or flatMapLatest
fun main() = runBlocking { println("Choose the right strategy!") }

ตรวจสอบความเข้าใจอย่างรวดเร็ว

flowOn มีผลต่อส่วนใดของสายการทำงานของโฟลว์

สรุปทบทวน

flowOn ย้ายงานต้นทางไปยังตัวจัดส่งงานอื่น buffer แยกผู้ผลิตออกจากผู้ใช้เพื่อเพิ่มอัตราการประมวลผล conflate ทิ้งค่าระหว่างกลาง และ collectLatest ยกเลิกการประมวลผลที่ช้าเมื่อมีค่าใหม่มาถึง

เริ่มต้นได้ฟรี

เรียนรู้ Kotlin ด้วย AI tutor — ฟรี

เขียนและเรียกใช้โค้ดจริงในเบราว์เซอร์ของคุณ รับความช่วยเหลือทันทีจาก AI tutor 24/7 และเรียนรู้ต่อจากที่คุณหยุดบนเว็บหรือในแอป

คอร์ส
51
บทเรียน
203

คำถามที่พบบ่อย

บทเรียน “flowOn และ buffer: บริบทและแรงดันย้อนกลับ” ฟรีหรือไม่

ใช่ — ข้อความเต็มของ “flowOn และ buffer: บริบทและแรงดันย้อนกลับ” ฟรีให้อ่านที่นี่บนเว็บ เพื่อปฏิบัติแบบโต้ตอบ (ตัวแก้ไขโค้ดในตัวและติวเตอร์ AI ตลอด 24/7) และปลดล็อคส่วนที่เหลือของคอร์ส Kotlin Academy ให้อัปเกรดเป็น CoddyKit PRO คอร์ส Kotlin Academy มีบทเรียนทั้งหมด 4 บทเรียน

คุณจะเรียนรู้อะไรในบทเรียน “flowOn และ buffer: บริบทและแรงดันย้อนกลับ”

เปลี่ยนบริบทการปล่อยค่าด้วย flowOn และพักค่าที่ปล่อยด้วย buffer เพื่อรองรับแรงดันย้อนกลับ คุณปฏิบัติ Kotlin Academy ด้วยโค้ดที่ใช้งานได้จริงที่คุณเรียกใช้โดยตรงในเบราว์เซอร์ และติวเตอร์ AI ตลอด 24/7 ตอบคำถามของคุณขณะที่คุณไปผ่านบทเรียน

คุณต้องมีประสบการณ์ก่อนที่จะเริ่มเรียน Kotlin Academy หรือไม่

ไม่จำเป็นต้องมีประสบการณ์มาก่อน Kotlin Academy บน CoddyKit ออกแบบมาสำหรับผู้เริ่มต้นไปจนถึงผู้เรียนขั้นสูง คุณสามารถเริ่มต้นที่นี่หรือเริ่มจากตัวแรกและเรียนด้วยความเร็วของคุณเอง นี่คือบทเรียนที่ 4 จากทั้งหมด 4 บทเรียน

บทเรียน “flowOn และ buffer: บริบทและแรงดันย้อนกลับ” ใช้เวลานานแค่ไหน

บทเรียน CoddyKit ส่วนใหญ่ใช้เวลาประมาณ 5–10 นาที แต่ละบทเรียนจึงสั้นและเป็นแบบโต้ตอบ คุณสามารถก้าวหน้าอย่างต่อเนื่องและกลับมาเรียนต่อจากตรงที่เพิ่งหยุดบนเว็บและแอปได้เลย

ฉันเขียนและรันโค้ดในบทเรียน Kotlin Academy นี้ได้ไหม

ได้ บทเรียน Kotlin Academy ทุกบทมีตัวแก้ไขโค้ดในตัว คุณจึงเขียนและรันโค้ดจริงได้เลยในเบราว์เซอร์ และได้รับข้อเสนอแนะจาก AI ในทันที — ไม่ต้องติดตั้งในเครื่องของคุณ

บทเรียนทั้งหมดในหลักสูตรนี้

  1. ตัวดำเนินการ Flow: map, filter, transform และ take
  2. catch และ onCompletion: การจัดการข้อผิดพลาดใน Flow
  3. combine และ zip: รวม Flow หลายสาย
  4. flowOn และ buffer: บริบทและแรงดันย้อนกลับ
← กลับไปที่ Kotlin Academy