Skip to content

Job 任务调度

job 模块提供基于 Kotlin 协程的进程内定时任务。它支持注解任务、DSL 动态任务、固定间隔、指定时刻、同步/异步执行、取消和异常事件。

它不是分布式调度系统:没有持久化、集群选主、错过补偿、执行记录和跨节点互斥。

适合 Job 的场景是单实例应用内的缓存刷新、临时数据清理和状态轮询。任务必须跨节点只执行一次、宕机后补跑或保留审计记录时,应使用专门的分布式调度系统。

最短使用路径是:添加依赖 → 给 Bean 标记 @CronTab → 给无参数方法标记 @Cron。需要运行时创建或取消任务时再使用 JobBuilderJobManager

添加依赖

kotlin
implementation("com.IceCreamQAQ.Rain:job:1.0.0-DEV12")

JobCenter 可选依赖 EventBus;同时引入 Event 模块后,任务异常会发布 JobRunExceptionEvent

注解任务

kotlin
@CronTab
class CacheJobs(
    private val cacheService: CacheService,
) {
    @Cron("5m")
    fun refresh() {
        cacheService.refresh()
    }
}
  • 类使用 @CronTab,触发 JobLoader
  • 类必须位于扫描包且可由 DI 创建。
  • 方法使用 @Cron
  • 普通 Kotlin 函数、Java 方法与 suspend 函数会包装为不同 CronInvoker
  • 任务方法应为无参数方法;Loader 不提供任务参数解析。

@Cron 参数

kotlin
@Cron(
    value = "30s",
    async = false,
    runWithStart = false,
)
参数说明
value当前版本解析为时间间隔字符串
async是否允许任务体并发运行
runWithStart注册时是否先立即执行一次

不是标准 Cron 表达式

虽然注解名是 Cron,当前 JobLoader 使用 JobBuilder.every(cron.value),字符串经 Rain toTime() 转为间隔。JobBuilder.cron() 仍是 TODO。不要填写 0 0 * * * 等标准 Cron 表达式。

常见时间字符串由 toTime() 支持,例如测试中的 1s。采用更多单位前应以 function 模块当前解析规则为准并增加测试。

动态任务 DSL

注入 JobManager,或直接使用 JobCenter

kotlin
class ReportScheduler(
    private val jobManager: JobManager,
) {
    fun schedule(): String = JobBuilder("daily-report")
        .every("10m")
        .task {
            reportService.generate()
        }
        .build()
        .let(jobManager::registerJob)
}

registerJob 返回随机 ID,后续可取消:

kotlin
val removed = jobManager.deleteTimer(jobId)

取消会从 JobCenter Map 移除任务并 cancel 调度协程。已经以 async 启动的独立任务体是否及时停止,取决于其协程父子关系和任务是否响应取消。

固定间隔 every

kotlin
JobBuilder("heartbeat")
    .every(5_000L)
    .task { heartbeat() }
    .build()

也可使用字符串:

kotlin
.every("5s")

EveryNext 根据上次开始时间与结束时间计算下一次 delay,使同步任务尽量落在固定时间网格上:若间隔 5 秒,任务耗时 7 秒,会跳过已经错过的 5 秒点,在约第 10 秒开始下一次。

指定时刻 at

kotlin
JobBuilder("hourly-cleanup")
    .at("32:00")
    .alaways()
    .task { cleanup() }
    .build()

格式:

  • mm:ss:本小时的某分某秒;已过则移动到下一小时。
  • HH:mm:ss:当天某时刻;已过则移动到下一天。

alaways()(源码当前拼写)让 at 任务按一小时或一天重复。没有调用时只执行一次。

kotlin
JobBuilder("one-shot")
    .at("23:30:00")
    .task { snapshot() }
    .build()

时区使用 ZoneId.systemDefault(),部署容器的系统时区会影响执行时刻。

first 与 next

可直接组合首次延迟和后续间隔:

kotlin
JobBuilder("delayed-poll")
    .first(10_000L)
    .next(60_000L)
    .task { poll() }
    .build()
  • 只有 first:延迟后执行一次。
  • first == next:退化为固定间隔。
  • 两者不同:先等待 first,以后按 next 周期。

自定义 NextTime

NextTime 接收上次开始和结束时间,返回“从现在起还需等待的毫秒数”;返回负数结束任务:

kotlin
val threeTimes = object : NextTime {
    var count = 0
    override fun invoke(invokeTime: Long, endTime: Long): Long =
        if (count++ >= 3) -1 else 1_000
}

JobBuilder 当前没有公开直接设置任意 NextTime 的 DSL 方法,框架扩展可直接构造 JobRuntime,但这会依赖非稳定构造细节。

同步与 async 语义

默认 async = false:调度协程等待任务完成,再计算下一次 delay。不会发生同一 Job 的重叠执行。

kotlin
@Cron("5s", async = false)
fun reconcile() = Unit

async = true 时,JobRuntime 在 scope 中 launch 任务体,然后调度循环立即记录结束时间并继续等待。若执行时间超过间隔,多个任务体可以同时运行:

kotlin
@Cron("5s", async = true)
suspend fun fetchRemote() = Unit

只有任务天然可并发且依赖线程安全时才开启。数据库批处理、结算等任务通常应保持同步,或自行实现分布式锁/互斥。

runWithStart

注解 Loader 在注册周期任务前可立即调用一次:

kotlin
@Cron("1h", runWithStart = true)
fun warmCache() = Unit

这次调用发生在 Loader 加载期间,而周期 Job 随后注册到 JobCenter。若启动执行很慢且是同步调用,会延长应用启动时间;启动必须成功的初始化更适合 ApplicationService.start()

任务线程池与生命周期

JobCenter 创建:

text
coreNumThreadPool("Job")
  + SupervisorJob
  + CoroutineScope

每个 JobRuntime 有一个调度 Job。JobCenter 也是 ApplicationService:停止时取消所有调度任务、取消 Scope 并关闭线程池。

SupervisorJob 防止单个任务异常直接取消整个调度 Scope。

异常处理

JobRuntime 捕获任务体异常、记录日志并调用 error callback。存在 EventBus 时发布:

kotlin
@EventListener
class JobErrorListener {
    @SubscribeEvent
    fun onError(event: JobRunExceptionEvent) {
        alertService.report(event.job, event.error)
    }
}

异常不会终止周期调度;下一时间点仍会继续执行。需要连续失败熔断时,可在任务体或监听器中记录失败次数并调用 deleteTimer

async 异常边界

当前 JobRuntime.invoke()runCatching 中调用 scope.launch { invoker() }。异步协程内部稍后抛出的异常不一定被外层 runCatching 捕获,因此 JobRunExceptionEvent 对 async 任务未必完整。关键任务建议保持同步,或在任务体内部 try/catch 并主动上报。

与 Spring @Scheduled 对照

Rain JobSpring Scheduling
@Cron("5s")@Scheduled(fixedRate = 5000) 更接近
async配合 @Async / 自定义 TaskScheduler
JobBuilderTaskScheduler.schedule*
JobManager ID 取消ScheduledFuture.cancel
JobRunExceptionEventErrorHandler / 应用事件

Spring 的 cron 属性是真正 Cron 表达式;Rain 当前 @Cron 是间隔 DSL。Rain 实现更短:JobRuntime 就是一个 delay 循环加 invoker,容易理解和定制,但没有 Spring 调度器的 Trigger、线程池配置和成熟监控生态。

生产使用检查表

  • 明确任务是单机一次还是每节点一次。
  • 多实例部署时为非幂等任务增加分布式锁。
  • 不把标准 Cron 表达式写进当前 @Cron
  • async 任务自己捕获异常和限制并发。
  • 指定部署时区,避免 at 任务随环境变化。
  • 保存动态 Job ID,服务停止或业务删除时取消。
  • 任务体支持协程取消,不阻塞 JobCenter 关闭。
  • 对重要任务记录开始、结束、耗时和结果。

基于 Apache License 2.0 发布