Job 任务调度
job 模块提供基于 Kotlin 协程的进程内定时任务。它支持注解任务、DSL 动态任务、固定间隔、指定时刻、同步/异步执行、取消和异常事件。
它不是分布式调度系统:没有持久化、集群选主、错过补偿、执行记录和跨节点互斥。
适合 Job 的场景是单实例应用内的缓存刷新、临时数据清理和状态轮询。任务必须跨节点只执行一次、宕机后补跑或保留审计记录时,应使用专门的分布式调度系统。
最短使用路径是:添加依赖 → 给 Bean 标记 @CronTab → 给无参数方法标记 @Cron。需要运行时创建或取消任务时再使用 JobBuilder 与 JobManager。
添加依赖
implementation("com.IceCreamQAQ.Rain:job:1.0.0-DEV12")JobCenter 可选依赖 EventBus;同时引入 Event 模块后,任务异常会发布 JobRunExceptionEvent。
注解任务
@CronTab
class CacheJobs(
private val cacheService: CacheService,
) {
@Cron("5m")
fun refresh() {
cacheService.refresh()
}
}- 类使用
@CronTab,触发JobLoader。 - 类必须位于扫描包且可由 DI 创建。
- 方法使用
@Cron。 - 普通 Kotlin 函数、Java 方法与
suspend函数会包装为不同CronInvoker。 - 任务方法应为无参数方法;Loader 不提供任务参数解析。
@Cron 参数
@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:
class ReportScheduler(
private val jobManager: JobManager,
) {
fun schedule(): String = JobBuilder("daily-report")
.every("10m")
.task {
reportService.generate()
}
.build()
.let(jobManager::registerJob)
}registerJob 返回随机 ID,后续可取消:
val removed = jobManager.deleteTimer(jobId)取消会从 JobCenter Map 移除任务并 cancel 调度协程。已经以 async 启动的独立任务体是否及时停止,取决于其协程父子关系和任务是否响应取消。
固定间隔 every
JobBuilder("heartbeat")
.every(5_000L)
.task { heartbeat() }
.build()也可使用字符串:
.every("5s")EveryNext 根据上次开始时间与结束时间计算下一次 delay,使同步任务尽量落在固定时间网格上:若间隔 5 秒,任务耗时 7 秒,会跳过已经错过的 5 秒点,在约第 10 秒开始下一次。
指定时刻 at
JobBuilder("hourly-cleanup")
.at("32:00")
.alaways()
.task { cleanup() }
.build()格式:
mm:ss:本小时的某分某秒;已过则移动到下一小时。HH:mm:ss:当天某时刻;已过则移动到下一天。
alaways()(源码当前拼写)让 at 任务按一小时或一天重复。没有调用时只执行一次。
JobBuilder("one-shot")
.at("23:30:00")
.task { snapshot() }
.build()时区使用 ZoneId.systemDefault(),部署容器的系统时区会影响执行时刻。
first 与 next
可直接组合首次延迟和后续间隔:
JobBuilder("delayed-poll")
.first(10_000L)
.next(60_000L)
.task { poll() }
.build()- 只有
first:延迟后执行一次。 first == next:退化为固定间隔。- 两者不同:先等待 first,以后按 next 周期。
自定义 NextTime
NextTime 接收上次开始和结束时间,返回“从现在起还需等待的毫秒数”;返回负数结束任务:
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 的重叠执行。
@Cron("5s", async = false)
fun reconcile() = Unitasync = true 时,JobRuntime 在 scope 中 launch 任务体,然后调度循环立即记录结束时间并继续等待。若执行时间超过间隔,多个任务体可以同时运行:
@Cron("5s", async = true)
suspend fun fetchRemote() = Unit只有任务天然可并发且依赖线程安全时才开启。数据库批处理、结算等任务通常应保持同步,或自行实现分布式锁/互斥。
runWithStart
注解 Loader 在注册周期任务前可立即调用一次:
@Cron("1h", runWithStart = true)
fun warmCache() = Unit这次调用发生在 Loader 加载期间,而周期 Job 随后注册到 JobCenter。若启动执行很慢且是同步调用,会延长应用启动时间;启动必须成功的初始化更适合 ApplicationService.start()。
任务线程池与生命周期
JobCenter 创建:
coreNumThreadPool("Job")
+ SupervisorJob
+ CoroutineScope每个 JobRuntime 有一个调度 Job。JobCenter 也是 ApplicationService:停止时取消所有调度任务、取消 Scope 并关闭线程池。
SupervisorJob 防止单个任务异常直接取消整个调度 Scope。
异常处理
JobRuntime 捕获任务体异常、记录日志并调用 error callback。存在 EventBus 时发布:
@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 Job | Spring Scheduling |
|---|---|
@Cron("5s") | @Scheduled(fixedRate = 5000) 更接近 |
async | 配合 @Async / 自定义 TaskScheduler |
| JobBuilder | TaskScheduler.schedule* |
| JobManager ID 取消 | ScheduledFuture.cancel |
| JobRunExceptionEvent | ErrorHandler / 应用事件 |
Spring 的 cron 属性是真正 Cron 表达式;Rain 当前 @Cron 是间隔 DSL。Rain 实现更短:JobRuntime 就是一个 delay 循环加 invoker,容易理解和定制,但没有 Spring 调度器的 Trigger、线程池配置和成熟监控生态。
生产使用检查表
- 明确任务是单机一次还是每节点一次。
- 多实例部署时为非幂等任务增加分布式锁。
- 不把标准 Cron 表达式写进当前
@Cron。 - async 任务自己捕获异常和限制并发。
- 指定部署时区,避免 at 任务随环境变化。
- 保存动态 Job ID,服务停止或业务删除时取消。
- 任务体支持协程取消,不阻塞 JobCenter 关闭。
- 对重要任务记录开始、结束、耗时和结果。