Kotlin协程实现多任务调度全局暂停与恢复方案咨询
Kotlin协程实现可暂停/恢复的任务调度方案
核心思路
基于Kotlin协程构建调度器,核心是状态控制+任务队列管理:
- 用线程安全的状态变量控制调度器的运行/暂停状态
- 维护按时间排序的待执行任务队列,暂停后可快速筛选出错过的任务
- 利用协程的
delay()实现精准延迟,错误触发暂停后等待用户输入,恢复时批量执行所有未完成(含错过)的任务
关键组件与实现步骤
1. 定义任务模型与状态管理
先封装任务数据,用AtomicBoolean做线程安全的状态开关(也可用MutableStateFlow实现响应式状态):
import java.time.Instant import java.util.concurrent.atomic.AtomicBoolean import kotlinx.coroutines.* // 调度任务:包含目标执行时间与具体操作 data class ScheduledTask(val targetTime: Instant, val action: suspend () -> Unit) // 调度器核心类 class TaskScheduler(private val coroutineScope: CoroutineScope = CoroutineScope(Dispatchers.IO)) { // 运行状态:true=正常调度,false=暂停 private val isRunning = AtomicBoolean(true) // 待执行任务队列(按时间升序排序) private lateinit var tasks: List<ScheduledTask> // 初始化:从文件读取并解析任务 fun loadTasksFromFile(filePath: String) { tasks = java.io.File(filePath).readLines() .map { line -> // 假设每行是ISO格式的时间,比如"2024-05-20T14:30:00Z" val targetTime = Instant.parse(line.trim()) ScheduledTask(targetTime) { // 这里替换为实际任务逻辑,比如打印日志、调用业务方法 println("Executing task at ${Instant.now()} (scheduled for $targetTime)") } } .sortedBy { it.targetTime } // 必须按时间排序,保证调度顺序 }
2. 调度与错误处理逻辑
实现循环调度,遇到无法处理的错误时触发暂停,等待用户恢复:
// 启动调度 fun start() { coroutineScope.launch { tasks.forEach { task -> // 检查是否处于暂停状态,若暂停则阻塞等待用户恢复 while (!isRunning.get()) { println("Scheduler paused. Enter 'resume' to continue...") val input = readLine()?.trim() if (input.equals("resume", ignoreCase = true)) { isRunning.set(true) println("Scheduler resumed. Starting to catch up on missed tasks...") } } try { val now = Instant.now() when { // 任务已错过,立即执行补跑 now.isAfter(task.targetTime) -> { println("Catching up missed task (scheduled for ${task.targetTime})") task.action() } // 未到执行时间,延迟等待 else -> { val delayMillis = task.targetTime.toEpochMilli() - now.toEpochMilli() delay(delayMillis) task.action() } } } catch (e: Exception) { // 捕获无法自行处理的错误,触发暂停 println("Unrecoverable error: ${e.message}. Pausing scheduler...") isRunning.set(false) } } println("All tasks completed.") } } // 手动触发暂停(可选) fun pause() { isRunning.set(false) } }
3. 控制台程序入口
fun main() = runBlocking { val scheduler = TaskScheduler() scheduler.loadTasksFromFile("tasks.txt") scheduler.start() // 保持主线程存活,直到所有任务完成 scheduler.coroutineScope.coroutineContext.job.join() }
关键注意事项
- 任务排序:必须按目标时间升序排列,确保调度顺序正确,补跑时不会遗漏
- 线程安全:状态变量用
AtomicBoolean保证多协程下的状态一致性,避免竞态条件 - 协程作用域:自定义
CoroutineScope而非使用GlobalScope,方便管理生命周期 - 补跑逻辑:恢复后遍历剩余任务,直接执行所有已错过的任务,未错过的继续按原时间调度
- 错误边界:仅在捕获到无法处理的异常时触发暂停,可根据业务需求调整异常捕获范围
内容的提问来源于stack exchange,提问作者Gregor Lang
相关产品推荐
相关产品推荐

