You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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中间操作实现任务管控,同时调整批量任务的发射逻辑,确保每个任务都能被精准控制:

  1. 添加暂停状态控制器
    用MutableStateFlow维护实时暂停状态,支持状态变更的即时响应:

    private val isPaused = MutableStateFlow(false)
    
  2. 修改Flow处理逻辑
    在任务处理前插入暂停检查,利用awaitUntil挂起协程直到暂停状态解除,确保只有非暂停时才执行任务:

    init {
        viewModelScope.launch(IO) {
            faceProcessorFlow
                .transform { item ->
                    // 等待暂停状态解除后再发射任务
                    isPaused.awaitUntil { !it }
                    emit(item)
                }
                .onEach { performEnrollment(it) }
                .collect()
        }
    }
    
  3. 重构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()
                }
            }
        }
    }
    
  4. 对外暴露暂停/恢复方法
    提供控制接口供外部调用:

    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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.13 07:57:07