跳到主要内容

调度与触发:Scheduler 与 Trigger

30 秒导读: 执行(Execution)不会凭空产生。要么有人在 UI 点了运行,要么就是本章讲的两条路:到点了(定时,cron)或发生了某件事(轮询到外部变化、别的流程跑完了)。Scheduler 是那个"每秒看一眼表、看谁该跑"的组件;Trigger 是挂在 Flow 上、描述"什么时候该跑"的声明。这一章讲清楚:一个 Execution 是怎么从"时间到了 / 事件来了"被造出来的。

本章只讲执行从哪来。Execution 被造出来之后怎么一步步跑完,是执行引擎:Executor 状态机的事,这里不碰。


1. 这是什么(零基础也能懂)

一句话定义

Scheduler(调度器)是一个后台服务,它按固定节奏(每秒一次)检查所有触发器,把"到期该评估"的触发器挑出来评估;评估通过就生成一个 Execution 扔进队列。

解决什么问题

假设你写了个 Flow,想让它:

  • 每天早上 9 点跑一次 —— 这是 Schedule 触发器(定时,cron)。
  • 每 30 秒查一下某个 S3 桶,有新文件就跑 —— 这是 Polling 触发器(轮询)。
  • 某个上游 Flow 一旦成功就立刻跑 —— 这是 Flow 触发器(被别的执行触发)。
  • 盯着一个 Kafka topic,来一条消息就跑一次 —— 这是 Realtime 触发器(实时流)。

这四种"什么时候跑",都由触发器声明,由 Scheduler(和 Worker)负责兑现。

用起来什么样

用户只在 Flow 的 YAML 里写触发器,剩下的全是后台的事:

id: daily_report
namespace: company.team
tasks:
- id: run
type: io.kestra.plugin.core.log.Log
message: "hello"
triggers:
- id: every_morning
type: io.kestra.plugin.core.trigger.Schedule
cron: "0 9 * * *" # 每天 9:00

写完这段,Scheduler 就会在每天 9:00 自动造出一个 daily_report 的 Execution。用户不需要写任何调度代码。

一句话直觉

把 Scheduler 想成一个尽职的门卫,手里拿着一张"谁几点该进来"的名单(触发器状态表)。他每秒看一次表和名单,到点的名字就放行(生成执行),然后把这个名字的下次时间往后推一格。 名单存在数据库里,所以门卫换班(Scheduler 重启或多实例)也不会丢。


2. 顶层全景(它大概怎么转)

两条主线,先分清

Kestra 的触发分两个世界,由谁评估是最关键的分界线:

触发器类型典型例子谁算"下次时间"谁做"真正的评估"为什么这么分
Schedule(定时)ScheduleScheduleOnDatesSchedulerScheduler 自己纯 cron 数学 + 造执行,不碰外部系统,没必要外派
Polling(轮询)HTTP、S3、SQL 等Scheduler(now+interval)Worker要访问外部系统、可能很慢,不能堵住调度循环
Realtime(实时)Kafka/Pulsar 实时消费Scheduler(立即)Worker(长驻)一个连接常驻 Worker,来一条消息发一个执行
Flow(流触发)上游 Flow 跑完不涉及时间Executor由"别的执行结束"驱动,不在 Scheduler 里

这一章的核心分工,记住这句: Scheduler 负责"算时间 + 挑出到期的",Schedule 类的评估它顺手自己做了,Polling/Realtime 类它派给 Worker 做;Flow 触发压根不经过 Scheduler,由 Executor 处理。

部件一句话职责

部件干什么在哪个文件
DefaultScheduler调度服务本体:起线程、分配 vNode、起停调度循环scheduler/DefaultScheduler.java
TriggerSchedulingLoop调度主循环:每秒 tick 一次,调 onSchedule,顺带处理触发器事件scheduler/TriggerSchedulingLoop.java
TriggerScheduler调度核心逻辑:挑出到期触发器 → 评估 → 派发scheduler/TriggerScheduler.java
DefaultSchedulableTriggerFetcher从状态库拉出"到期且有效"的触发器scheduler/internals/DefaultSchedulableTriggerFetcher.java
SchedulableEvaluator在 Scheduler 内评估 Schedule 类触发器scheduler/internals/SchedulableEvaluator.java
NextEvaluationDate算"下一个评估时间点"的工具scheduler/internals/NextEvaluationDate.java
TriggerEventHandler处理触发器状态变更事件(创建/更新/评估回执/回填…)scheduler/TriggerEventHandler.java
TriggerStateStore / CachedTriggerStateStore触发器状态的读写与按 vNode 分片缓存core/.../scheduler/store/scheduler/stores/
TriggerWorkerJobPublisher把 Polling/Realtime 触发器打包发到 Worker 队列scheduler/pubsub/TriggerWorkerJobPublisher.java
FlowTriggerService由"某个执行结束"计算要触发的下游 Flow 执行executor/FlowTriggerService.java

主线走一遍(高层,不进代码)

以一个 Schedule 触发器为例,从时间到点到执行生成:

每秒一次


TriggerSchedulingLoop.run() ── tick ──▶ TriggerScheduler.onSchedule(now)


Fetcher: 从状态库挑出 nextEvaluationDate <= now 的触发器


对每个到期触发器 evaluate():
① 先算下一个 cron 时间点,写回状态(往后推一格)
② 判断类型:
Schedule ─▶ 自己评估(SchedulableEvaluator)─▶ 生成 Execution
Polling ─▶ 打包发 Worker 队列(评估在 Worker 上跑)
Realtime ─▶ 打包发 Worker 队列(长驻消费)


Execution ──▶ 执行队列 ──▶ Executor(第2章)

Polling/Realtime 走 Worker 那条,评估结果会以 WorkerTriggerResult 回流,变成一个 TriggerEvaluated 事件,再由 TriggerEventHandler 消化(解锁、写回、必要时发执行)。这条回路见 §4.5。


3. 核心原理(逐个机制,由浅入深)

3.1 每秒一拍的调度循环

它要解决的小问题: 怎么做到"到点就跑",又不至于疯狂空转烧 CPU?

思路: 一个固定节奏的循环,每秒醒一次,看看有没有到期的。节奏由一个常量定死——1 秒:

// scheduler/TriggerSchedulingLoop.java:38
private static final long SCHEDULE_INTERVAL_MILLIS = Duration.ofSeconds(1).toMillis();

主循环长这样(简化,示意非源码):

// 示意,非源码:每一拍做三件事
while (running) {
processTriggerEvents(); // 1) 先消化排队的触发器事件(状态变更)
if (now >= nextScheduleTime) {
triggerScheduler.onSchedule(now); // 2) 到拍点了,挑出到期触发器并评估
nextScheduleTime += 1s; // 3) 把下一拍往后推 1 秒
}
waitForNextIterationOrNewEvent(...); // 睡到下一拍,或被新事件叫醒
}

真实主体在 TriggerSchedulingLoop.run()(scheduler/TriggerSchedulingLoop.java:115)。几个值得注意的细节:

  • 拍点用挂钟时间对齐,循环间隔用单调时钟测量。 nextScheduleTime 基于 clock.instant() 推进(:171),而"这一拍花了多久"用 System.nanoTime() 测(:126),避免系统时间跳变把节奏带偏。
  • 落后了会告警。 如果一拍耗时超过 1.1 秒,打 WARN:线程饥饿或触发器太多算不过来(:132-135)。
  • 能被事件叫醒。 如果这一拍没到,但来了新的触发器事件,waitForNextIterationOrNewEvent 会被 signal 提前唤醒(:200-211、:362-368),不必干等满 1 秒。

3.2 挑出"到期"的触发器

它要解决的小问题: 每秒扫一遍,怎么快速知道"谁到期了"?

思路: 每个触发器在状态库里都有一个 nextEvaluationDate(下次该评估的时间)。"到期"就是 nextEvaluationDate <= now。这一步由 Fetcher 完成:

// scheduler/internals/DefaultSchedulableTriggerFetcher.java:65-66
public List<TriggerEvaluationContext> getSchedulableTriggers(final Clock clock, final ZonedDateTime now, final Set<Integer> assignments) {
List<TriggerState> triggers = this.triggerStateStore.findTriggersEligibleForScheduling(now, assignments, false);

findTriggersEligibleForScheduling(now, vNodes, locked=false) 的语义(core/.../scheduler/store/TriggerStateStore.java:31):在给定 vNode 范围内,拿出所有 nextEvaluationDate <= now未加锁的触发器。locked=false 很关键——加了锁的触发器(上一次的执行还没结束、且不允许并发)会被跳过。

Fetcher 拿到候选后还要二次校验,把"陈旧或失效"的滤掉(:82-101):

情况处理
Flow 已被删除删掉该触发器状态,跳过
Flow 被禁用跳过
触发器已从 Flow 里移除跳过
触发器自身被禁用跳过

为什么要二次校验?因为 Flow 的更新是通过事件异步同步到 Scheduler 的,状态库里的触发器状态可能比 Flow 定义"旧一拍"(:93-95 注释)。

3.3 评估:先算下一拍,再决定谁来评估

它要解决的小问题: 一个到期的触发器,拿到手要做什么?

核心方法是 TriggerScheduler.evaluate()(scheduler/TriggerScheduler.java:277)。 它做的第一件事很反直觉:先把 evaluatedAt 设成本次的 nextEvaluationDate(:283)——即"我现在正在评估的是这个时间点"。然后按类型分流:

// scheduler/TriggerScheduler.java:291-301(节选)
switch (trigger) {
case Schedulable schedulableTrigger ->
processSchedulableTrigger(...); // Schedule:自己评估
case PollingTriggerInterface pollingTrigger ->
processPollingTrigger(...); // Polling:派给 Worker
case RealtimeTriggerInterface realtimeTrigger ->
processWorkerTrigger(...); // Realtime:派给 Worker
...
}

注意分派顺序。 Schedule 实现了 Schedulable,而 Schedulable extends PollingTriggerInterface(core/.../models/triggers/Schedulable.java:14)。因为 switchSchedulable 分支排在 PollingTriggerInterface 前面,Schedule 类会走"自己评估"这条,不会被误当成普通 Polling 派出去。

关键设计:下一拍在评估之前就算好、写回。 无论哪条分支,都会先调 NextEvaluationDate 算出下一个评估时间点,updateForNextEvaluationDate 写回状态(如 processSchedulableTrigger 的 :335-336)。这样即使本次评估失败或没产生执行,触发器也已经"排好了下次",不会卡死。

3.4 Schedule 为什么在 Scheduler 里评估

Schedule 触发器的评估很轻:算出这次的 cron 时间点,然后造一个 Execution。 不碰任何外部系统,所以没必要外派,直接在调度线程里做:

// scheduler/TriggerScheduler.java:339
Optional<TriggerEvaluationResult> evaluationResult = schedulableEvaluator.evaluate(trigger, triggerContext, ...);

SchedulableEvaluator.evaluate(scheduler/internals/SchedulableEvaluator.java:36)内部调触发器的 eval(),对 Schedule 而言就是"这个 cron 时间点该不该产生执行、产生什么"。产生了就:

  • 更新状态:updateOnExecutionCreated,若不允许并发则加锁并记下 executionId(:343-346)。
  • 把 Execution 发到执行队列:triggerExecutionSender.send(execution)(:358)。

Schedule 的"下一个时间点"怎么算?Schedule.nextEvaluationDate(core/.../plugin/core/trigger/Schedule.java:239):用 cron 表达式算下一个触发时间,并把结果 truncatedTo(ChronoUnit.SECONDS) 截到秒(:262)。它还配合 previousEvaluationDate(:311)用于"补跑漏掉的调度"(见 §3.7)。

3.5 Polling / Realtime 为什么派给 Worker

这两类的评估可能很慢(网络 I/O)或长驻(实时流),绝不能堵在每秒一拍的调度线程里。 所以 Scheduler 只做两件事:算下次时间、把活儿打包扔给 Worker。

// scheduler/TriggerScheduler.java:375-384(processWorkerTrigger 节选)
TriggerState dispatched = state.nextDispatchEpoch(clock);
if (this.triggerWorkerJobPublisher.send(dispatched, ...)) {
triggerStateStore.save(dispatched
.lastTriggeredDate(clock)
.locked(clock, mustBeLocked)); // 派发成功:加锁,防止重复派
} else {
triggerStateStore.save(state); // 没派出去:不加锁,下拍重试
}

几个要点:

  • 是否加锁看类型。 Realtime 一定加锁;Polling 看 allowConcurrent——不允许并发才加锁(:376)。加锁后这个触发器在结果回来前不会被再次挑中(§3.2 里 locked=false 的过滤)。
  • 派发用 dispatchEpoch 打世代号。 nextDispatchEpoch 把状态的世代 +1(TriggerState.java:330)。这个号用来在结果回流时分辨"这是这一次派发的回执,还是被取代的旧回执"(§4.5 解释为什么需要它)。
  • 派不出去就不加锁。 比如没有匹配的 Worker 队列、没有可用 Worker,send 返回 false,状态原样保存,下一拍按 nextEvaluationDate 重试(TriggerWorkerJobPublisher.java:53 的返回语义)。

派发的真正动作在 TriggerWorkerJobPublisher.send(scheduler/pubsub/TriggerWorkerJobPublisher.java:53):把触发器包成 WorkerTrigger,按 WorkerSelector 的标签路由到某个 Worker 队列,emit 出去。

3.6 触发器状态机(TriggerState)

它要解决的小问题: "下次几点评估""是否加锁""是否被禁用""上次谁在跑"——这些必须持久化,不然 Scheduler 一重启全丢。

答案是 TriggerState(core/.../scheduler/model/TriggerState.java:32),一个不可变记录,每次变更都用 update(clock) 复制出新实例。核心字段:

字段含义
nextEvaluationDate下次该评估的时间(§3.2 的"到期"判据)
evaluatedAt本次正在评估的时间点
locked是否加锁(加锁则本轮不被挑中)
workerId当前持有该触发器的 Worker(派发后)
executionId非并发触发器当前占用锁的执行 id
disabled是否禁用(含 stopAfter 自动禁用)
backfill回填配置(见 §3.8)
dispatchEpoch派发世代号(防陈旧回执)
lastEventId最后应用的事件 id(事件去重,见 §4.4)
vnode该触发器所属的虚拟节点(分片用)

状态怎么随生命周期流转(挑几条关键):

created ──▶ nextEvaluationDate=首个cron时间/now

到期被挑中 evaluate()

├─ Schedule 产生执行 ─▶ updateOnExecutionCreated(可能 locked=true, 记 executionId)

├─ 派给 Worker 成功 ─▶ nextDispatchEpoch + locked=true + 记 workerId

执行结束(TriggerExecutionTerminated)


updateOnExecutionTerminated ─▶ locked=false, executionId=null, workerId=null
若终态命中 stopAfter ─▶ disabled=true

updateOnExecutionTerminated(TriggerState.java:266)会解锁并清掉执行/Worker 关联;若执行的终态在 stopAfter 列表里,则自动禁用该触发器(:268)——这就是"跑失败 N 次后自动停"这类语义的落点。

3.7 补跑漏掉的调度(RecoverMissedSchedules)

它要解决的小问题: Scheduler 停机了两小时,期间有 4 个 9 点/10 点的 cron 点没跑。恢复后,这些漏掉的要不要补?

Kestra 给三种策略(core/.../models/triggers/RecoverMissedSchedules.java):

策略行为
ALL全部补跑(默认)
LAST只补最后一个漏掉的
NONE一个都不补,从现在起

策略在 Scheduler 启动/接管 vNode 时生效,由 TriggerScheduler.onStart 处理(scheduler/TriggerScheduler.java:207-228):

  • LAST:previousEvaluationDate 算出"上一个 cron 点",若它比状态里的 evaluatedAt 晚,就把 nextEvaluationDate 拨到那个点——下一拍立刻补跑这一个(:208-214)。
  • NONE: 从"现在"重新起算下一个点,直接跳过所有漏掉的(:215-224)。
  • ALL: 什么都不做——因为状态里的 nextEvaluationDate 还停在很久以前,主循环会自然地一个接一个把它们全部追平(:225-227)。

默认值来自插件配置,取不到就是 ALL(core/.../models/triggers/Schedulable.java:47-52)。

3.8 回填(Backfill)

回填 = 主动为一段过去的时间区间补造执行。 比如"帮我把上个月每天的报表都补跑一遍"。

配置见 Backfill(core/.../models/triggers/Backfill.java):start/end 是区间,currentDate 是"回填进行到哪了",paused 可暂停,还能带 inputs/labels

工作方式(inferred,综合 TriggerEventHandler.onCreateBackfillTriggerState):创建回填时,先把当前的 nextEvaluationDate 存进 previousNextExecutionDate 备份(TriggerState.java:229-239),然后把评估时间点拨到回填区间里,让主循环像追赶漏掉的调度一样逐点补跑。删除回填时(onDeleteBackfillTrigger,TriggerEventHandler.java:186)恢复备份的 nextEvaluationDate,回到正常节奏。每往前走一格,getBackFillForNextEvaluationDate(TriggerState.java:336)推进 currentDate,越过 end 就把回填清空,回填结束。


4. 深入实现

4.1 分布式分片:vNode 是怎么回事

问题: 多个 Scheduler 实例同时在线,同一个触发器不能被两个实例都评估(会造出两个执行)。

Kestra 的解法是虚拟节点(vNode)分片。 每个触发器按 Flow 哈希到一个 vNode(VNodes.computeVNodeFromFlow,TriggerScheduler.java:170);所有 vNode 通过一致性哈希环分给在线的 Scheduler 实例,一个 vNode 同一时刻只属于一个实例。于是"同一个触发器只被一个实例评估"就有了保证。

DefaultScheduler 订阅 vNode 的分配与撤销(scheduler/DefaultScheduler.java:151-203):

  • 被分配 vNode 时(onVNodesAssigned,:168):预热该 vNode 的触发器状态缓存(:178),按 vNodeId % maxThreads 把 vNode 分给各条调度循环(:190-193),然后开跑。
  • 被撤销时(onVNodesRevoked,:153):停掉队列消费和所有调度循环——这些 vNode 要交给别的实例了。

每条 TriggerSchedulingLoop 只处理分到自己名下的 vNode(assignments),循环开头若 assignments 为空就空转等待(TriggerSchedulingLoop.java:149-156)。

4.2 触发器状态的分片缓存

CachedTriggerStateStore(scheduler/stores/CachedTriggerStateStore.java:26)是一个装饰器,在真正的状态库(JDBC 等)前面套了一层按 vNode 分片的 Caffeine 缓存:

  • partitionedCachevNode -> Cache<uid, TriggerState>(:32),每个 vNode 一块独立缓存,按 cacheMaxSizePerVNode 限容(:41)。
  • save 写穿:先落库,再更新对应 vNode 的缓存(:112-119)。
  • init(vNodes) 在分片变动时被调用:撤销的 vNode 整块缓存丢弃(:146-154),新分到的从库里加载预热(:157-175)。

注意 findTriggersEligibleForScheduling(§3.2 那个"挑到期"的查询)是直接透传给 delegate、不走缓存的(:52-55)——每拍都要最新的"谁到期",缓存反而会给出陈旧结果。

4.3 算下一个评估时间点的两个重载

NextEvaluationDate(scheduler/internals/NextEvaluationDate.java:16)有两个 get,差别是"知不知道条件":

重载签名用在哪会不会考虑触发器条件
ConditionContextget(clock, trigger, triggerContext, conditionContext)(:28)正常路径(能应用 DayWeek 这类条件)
不带get(clock, trigger)(:57)异常兜底不会,只给下一个原始 cron 点

带条件那个还有一处精巧处理:当上下文既没有历史评估日期、又没有回填时,它把 date 播种为"现在"(:35-40),这样 Schedule 会走"考虑条件"的分支,而不是退化成"下一个不带条件的 cron 点"。这解决了"新建触发器时 DayWeek=SUNDAY 之类条件被忽略"的坑。对 Realtime 这种非 Polling 触发器,直接返回 now(:29-31)。

4.4 事件驱动的状态变更与去重

Scheduler 的状态变更不只来自"每秒评估",还来自一堆异步事件:触发器被创建/更新/删除、执行结束、Worker 评估回执、回填命令……这些都是 TriggerEvent,由 TriggerEventHandler.handle(scheduler/TriggerEventHandler.java:101)统一处理,内部 switch 分派(:116-135):

事件含义处理要点
TriggerCreated新触发器建状态,算首个 nextEvaluationDate(:506)
TriggerUpdated定义变了杀掉在跑的实例,用新定义重算(:429)
TriggerEvaluatedWorker 评估回执解锁/写回,有执行就发出(:312,见 §4.5)
TriggerExecutionTerminated执行结束解锁,按终态判断是否 stopAfter 禁用(:253)
TriggerReceivedWorker 已接手workerId(:354)
TriggerWorkerLost持有触发器的 Worker 没了解锁,让触发器重新可被派发(:386)
CreateBackfill/SetPauseBackfill/DeleteBackfill回填命令见 §3.8
ResetTrigger/SetDisableTrigger重置/启停重算或切换禁用

去重是这里的关键。 大多数队列是"至少一次"投递,同一事件可能来两遍。findTriggerState(:557)在应用事件前比对 lastEventId:只有当前事件"比上次应用的更新"才处理,否则丢弃(:566-573):

// scheduler/TriggerEventHandler.java:566-573(节选)
EventId lastEventId = current.getLastEventId();
if (lastEventId == null || event.eventId().isNewerThan(lastEventId)) {
return state; // 是新事件,处理
}
// 否则:更旧或重复,跳过

4.5 Worker 评估回路:一次 Polling 的完整往返

Polling/Realtime 被派到 Worker 后,评估在 Worker 上跑。以 Polling 为例:

Scheduler Worker Scheduler
│ │ │
processWorkerTrigger │ │
│ emit WorkerTrigger ──────────▶ │ │
│ (locked=true, │ WorkerTriggerCallable.doCall │
│ dispatchEpoch=N) │ pollingTrigger.eval(...) │
│ │ 查外部系统,有变化就产出结果 │
│ │ ── WorkerTriggerResult ─────▶ │
│ │ TriggerEvaluated(N)
│ │ │
│ │ TriggerEventHandler.onTriggerEvaluated
│ │ ├─ 算下一个评估时间,写回
│ │ ├─ 有执行:发出 Execution
│ │ └─ 无执行:解锁(下次可再派)
  • Worker 侧评估: WorkerTriggerCallable.doCall(worker/.../internals/WorkerTriggerCallable.java:33)直接调 pollingTrigger.eval(conditionContext, triggerContext)——真正的"查 S3、发 HTTP"在这里发生,在 Worker 线程上,不占用 Scheduler。
  • 结果回流: Worker 把 WorkerTriggerResult 发回,worker-controller 把它翻译成事件(worker-controller/.../GrpcWorkerControllerService.java:247):有评估结果就发 TriggerEvaluated,评估失败(启动不了)就发 TriggerExecutionTerminated(FAILED)。
  • 回执处理: onTriggerEvaluated(TriggerEventHandler.java:312)——有执行则写回状态并 triggerExecutionPublisher.send(execution)(:342-345);没有执行(轮询什么都没匹配到,或作业在派发前被拒)则解锁(:330-337),否则这个触发器会永远卡在锁里再也不被评估。

为什么 Realtime 的解锁要靠 dispatchEpoch 而不是 lastEventId? 一个在跑的 Realtime 触发器会吐出大量执行,它们的"终止"信号不能释放触发器的锁(否则会把还在 Worker 上跑的触发器重新派一遍)。Realtime 唯一合法的解锁信号是"启动失败"那个 FAILED 执行(TriggerEventHandler.java:280-305)。而这些回执来自多个独立的生产者、事件 id 不保证有序,所以这里改用派发世代 dispatchEpoch 做栅栏:回执的世代比当前状态旧就丢弃(:297-298)。

4.6 Flow 触发:被"别的执行"触发

Flow 触发器根本不经过 Scheduler。 它由"某个执行进入某状态"驱动,发生在 Executor 一侧,逻辑在 FlowTriggerService(executor/FlowTriggerService.java:33)。

每当一个执行状态变化,Executor 会问:有没有别的 Flow 声明了"监听这个执行"?分两条计算:

  • 标准条件(computeExecutionsFromFlowTriggerConditions,:64):只看非 dependsOn 的普通条件——先按"监听的状态"过滤(:202),再校验条件(:205),命中就 evaluate 出一个执行。
  • 多重条件 / dependsOn(computeExecutionsFromFlowTriggerDependsOn,:96):处理"上游 A 和 B 都在时间窗内成功了才触发"这类跨执行的累积条件,用 MultipleConditionWindow 把多次执行的结果攒在一个时间窗里(:132-173),窗口内条件全满足才触发。

还有一层防递归保护:computeFlowTriggers(:186)会调 flowService.removeUnwanted 阻止 Flow 触发自己造成的无限链,并滤掉测试类执行(:191-193)。

多重条件的窗口状态存在 MultipleConditionStateStore,过期窗口会被清理(:126-127)。条件模型本身在 core/.../models/conditions


5. 巧妙之处(可借鉴的技术)

  • "先排好下一拍,再评估" —— 无论评估成功、失败还是没产出,nextEvaluationDate 都已在评估前写回(TriggerScheduler.java:335-336)。触发器永远有"下一次",不会因一次异常永久卡死。
  • 加锁 + locked=false 过滤,天然防重叠 —— 不允许并发的触发器一旦派出去就加锁,而"挑到期"的查询只取未加锁的(DefaultSchedulableTriggerFetcher.java:66)。锁的存在与否直接实现了 allowConcurrent 语义,无需额外协调。
  • dispatchEpoch 世代栅栏 —— 面对"至少一次"投递 + 多生产者 + 无序回执,用一个单调递增的派发号区分"当前派发的回执"和"被取代的旧回执"(TriggerEventHandler.java:297TriggerState.java:330),比依赖事件顺序更稳。
  • lastEventId 事件去重 —— 状态里记住最后应用的事件 id,重复/过期事件直接丢(TriggerEventHandler.java:566),让整个事件处理具备幂等性。
  • 按 vNode 分片缓存,但"挑到期"故意绕过缓存 —— 读多的状态走缓存,唯独每秒都要精确结果的到期查询直连库(CachedTriggerStateStore.java:52),在性能与正确性之间做了清醒的取舍。
  • 播种 date=now 让条件生效 —— 一个不起眼但重要的修正:新建触发器时给上下文播种当前时间,避免 DayWeek 等条件在首次评估被跳过(NextEvaluationDate.java:35-40)。

6. 边界与局限

  • 调度粒度是 1 秒。 主循环固定 1 秒一拍(TriggerSchedulingLoop.java:38),达不到亚秒级精度。触发器太多算不过来时,一拍会超过 1 秒并打告警,进一步拉大延迟。
  • Polling 的实际间隔 ≥ interval,不是精确等于。 Scheduler 只保证"不早于 now+interval 再评估",加上派发、Worker 排队、回执往返的开销,实际间隔会略大。文档也建议依赖外部系统的 Polling 至少 PT30S(PollingTriggerInterface.java:18-21)。
  • Realtime 靠 Worker 常驻,Worker 掉了要靠事件恢复。 Worker 丢失通过 TriggerWorkerLost 解锁重派(TriggerEventHandler.java:386),但这依赖丢失检测与事件投递,不是瞬时的。
  • 杀不掉的 Realtime 会拖到重启。 若无法向 Realtime 触发器发送 kill(队列异常),它会一直跑到 Kestra 重启(TriggerEventHandler.java:497)。
  • Schedule 的评估占用调度线程。 Schedule 在 Scheduler 内评估(§3.4),虽然轻,但极端复杂的条件渲染仍会挤占调度节奏——这是"不外派"换来的代价。
  • 回填/补跑本质是"把时间往回拨让主循环追", 区间很大时会产生大量执行,需注意对下游的压力。

7. 横向对比

本章是 Kestra 子库的一章,和同组其它章的边界:


8. 代码地图(导航索引)

主题文件路径关键符号
调度服务本体、vNode 分配scheduler/src/main/java/io/kestra/scheduler/DefaultScheduler.javaDefaultSchedulerstartonVNodesAssigned
每秒调度主循环scheduler/src/main/java/io/kestra/scheduler/TriggerSchedulingLoop.javarunprocessTriggerEventsSCHEDULE_INTERVAL_MILLIS
调度核心:挑选/评估/派发scheduler/src/main/java/io/kestra/scheduler/TriggerScheduler.javaonScheduleevaluateprocessSchedulableTriggerprocessWorkerTriggeronStart
挑出到期且有效的触发器scheduler/src/main/java/io/kestra/scheduler/internals/DefaultSchedulableTriggerFetcher.javagetSchedulableTriggers
Schedule 类在 Scheduler 内评估scheduler/src/main/java/io/kestra/scheduler/internals/SchedulableEvaluator.javaSchedulableEvaluator.evaluate
算下一个评估时间点scheduler/src/main/java/io/kestra/scheduler/internals/NextEvaluationDate.javaNextEvaluationDate.get
触发器事件处理与去重scheduler/src/main/java/io/kestra/scheduler/TriggerEventHandler.javahandleonTriggerEvaluatedfindTriggerState
派发触发器到 Worker 队列scheduler/src/main/java/io/kestra/scheduler/pubsub/TriggerWorkerJobPublisher.javaTriggerWorkerJobPublisher.send
触发器状态(不可变)core/src/main/java/io/kestra/core/scheduler/model/TriggerState.javaTriggerStateupdateForNextEvaluationDatenextDispatchEpochupdateOnExecutionTerminated
状态库接口core/src/main/java/io/kestra/core/scheduler/store/TriggerStateStore.javafindTriggersEligibleForScheduling
分片缓存装饰器scheduler/src/main/java/io/kestra/scheduler/stores/CachedTriggerStateStore.javaCachedTriggerStateStorepartitionedCacheinit
触发器基类与语义core/src/main/java/io/kestra/core/models/triggers/AbstractTrigger.javaAbstractTriggerwhenallowConcurrentstopAfter
轮询触发器接口core/src/main/java/io/kestra/core/models/triggers/PollingTriggerInterface.javaevalgetIntervalnextEvaluationDate
实时触发器接口core/src/main/java/io/kestra/core/models/triggers/RealtimeTriggerInterface.javaRealtimeTriggerInterface.eval
定时触发器接口core/src/main/java/io/kestra/core/models/triggers/Schedulable.javaSchedulablepreviousEvaluationDatedefaultRecoverMissedSchedules
触发上下文core/src/main/java/io/kestra/core/models/triggers/TriggerContext.javaTriggerContextnextExecutionDatebackfill
回填模型core/src/main/java/io/kestra/core/models/triggers/Backfill.javaBackfillcurrentDatepreviousNextExecutionDate
补跑策略core/src/main/java/io/kestra/core/models/triggers/RecoverMissedSchedules.javaRecoverMissedSchedules(ALL/LAST/NONE)
执行/回执构建工具core/src/main/java/io/kestra/core/models/triggers/TriggerService.javagenerateEvaluationResultbuildLabels
Worker 侧轮询评估worker/src/main/java/io/kestra/worker/processors/internals/WorkerTriggerCallable.javaWorkerTriggerCallable.doCall
Worker 回执→事件翻译worker-controller/src/main/java/io/kestra/controller/grpc/services/GrpcWorkerControllerService.javasendWorkerTriggerResults
Flow 触发(由执行驱动)executor/src/main/java/io/kestra/executor/FlowTriggerService.javacomputeExecutionsFromFlowTriggerConditionscomputeExecutionsFromFlowTriggerDependsOn
Schedule 插件(Schedulable 实例)core/src/main/java/io/kestra/plugin/core/trigger/Schedule.javanextEvaluationDatepreviousEvaluationDateeval