如何在Kotlin中暂停与恢复协程?附具体业务代码场景
给批量人脸录入任务添加暂停/恢复功能
问题场景
现有如下Kotlin代码,调用start()后会批量处理itemsToDetect中的人脸录入任务,耗时随数据规模变化,需要实现按需暂停、恢复的功能:
private val faceProcessorFlow = MutableSharedFlow<Pair<EnrollmentModel.Waiting, Int>>() init { viewModelScope.launch(IO) { faceProcessorFlow .onEach { performEnrollment(it) }.collect() } } fun start(){ itemsToDetect.forEach { (index, itemToDetect) -> viewModelScope.launch(IO) { val item = (itemToDetect as? EnrollmentModel.Waiting) if (item != null) faceProcessorFlow.emit(item to index) else updateUI() } } }
解决思路与修改方案
核心思路是通过协程状态控制+Flow中间操作实现任务管控,同时调整批量任务的发射逻辑,确保每个任务都能被精准控制:
添加暂停状态控制器
用MutableStateFlow维护实时暂停状态,支持状态变更的即时响应:private val isPaused = MutableStateFlow(false)修改Flow处理逻辑
在任务处理前插入暂停检查,利用awaitUntil挂起协程直到暂停状态解除,确保只有非暂停时才执行任务:init { viewModelScope.launch(IO) { faceProcessorFlow .transform { item -> // 等待暂停状态解除后再发射任务 isPaused.awaitUntil { !it } emit(item) } .onEach { performEnrollment(it) } .collect() } }重构
start()的任务发射逻辑
原代码通过forEach启动多协程并发发射任务,无法逐个管控。改为单个协程串行遍历,确保每个任务发射前都检查暂停状态:fun start() { viewModelScope.launch(IO) { itemsToDetect.forEach { (index, itemToDetect) -> // 每次处理前先确认是否处于暂停状态 isPaused.awaitUntil { !it } val item = (itemToDetect as? EnrollmentModel.Waiting) if (item != null) { faceProcessorFlow.emit(item to index) } else { updateUI() } } } }对外暴露暂停/恢复方法
提供控制接口供外部调用:fun pauseTask() { isPaused.value = true } fun resumeTask() { isPaused.value = false }
完整修改后代码
import kotlinx.coroutines.flow.MutableSharedFlow import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.flow.transform import kotlinx.coroutines.flow.awaitUntil import kotlinx.coroutines.launch import androidx.lifecycle.viewModelScope import kotlin.coroutines.CoroutineContext // 假设IO是自定义的CoroutineContext private val IO: CoroutineContext = TODO() private val faceProcessorFlow = MutableSharedFlow<Pair<EnrollmentModel.Waiting, Int>>() private val isPaused = MutableStateFlow(false) init { viewModelScope.launch(IO) { faceProcessorFlow .transform { item -> isPaused.awaitUntil { !it } emit(item) } .onEach { performEnrollment(it) } .collect() } } fun start() { viewModelScope.launch(IO) { itemsToDetect.forEach { (index, itemToDetect) -> isPaused.awaitUntil { !it } val item = (itemToDetect as? EnrollmentModel.Waiting) if (item != null) { faceProcessorFlow.emit(item to index) } else { updateUI() } } } } fun pauseTask() { isPaused.value = true } fun resumeTask() { isPaused.value = false } // 原有业务方法保留 private fun performEnrollment(task: Pair<EnrollmentModel.Waiting, Int>) { // 人脸录入核心逻辑 } private fun updateUI() { // UI更新逻辑 } // 假设的模型类 sealed class EnrollmentModel { data class Waiting(val data: Any) : EnrollmentModel() }
关键说明
MutableStateFlow确保暂停状态的变更能即时通知所有监听的协程,无延迟awaitUntil替代轮询,仅在状态满足时才恢复协程,避免资源浪费- 串行化的任务发射逻辑确保每个任务的执行都受暂停状态管控,不会出现并发失控的情况
内容的提问来源于stack exchange,提问作者saas
相关产品推荐
相关产品推荐

