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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 00:37:29