数据截至 (上游 commit fcda16293bbd)
调度与触发: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(定时) | Schedule、ScheduleOnDates | Scheduler | Scheduler 自己 | 纯 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/src/main/java/io/kestra/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/src/main/java/io/kestra/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:318)。 它做的第一件事很反直觉:先把 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/src/main/java/io/kestra/core/models/triggers/Schedulable.java:14)。因为switch里Schedulable分支排在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/src/main/java/io/kestra/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:332)。这个号用来在结果回流时分辨"这是这一次派发的回执,还是被取代的旧回执"(§4.5 解释为什么需要它)。 - 派不出去就不加锁。 比如没有匹配的 Worker 队列、没 有可用 Worker,
send返回false,状态原样保存,下一拍按nextEvaluationDate重试(TriggerWorkerJobPublisher.java:52的返回语义)。
派发的真正动作在 TriggerWorkerJobPublisher.send(scheduler/pubsub/TriggerWorkerJobPublisher.java:52):把触发器包成 WorkerTrigger,按 WorkerSelector 的标签路由到某个 Worker 队列,emit 出去。
3.6 触发器状态机(TriggerState)
它要解决的小问题: "下次几点评估""是否加锁""是否被禁用""上次谁在跑"——这些必须持久化,不然 Scheduler 一重启全丢。
答案是 TriggerState(core/src/main/java/io/kestra/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:268)会解锁并清掉执行/Worker 关联;若执行的终态在 stopAfter 列表里,则自动禁用该触发器(:268)——这就是"跑失败 N 次后自动停"这类语义的落点。
3.7 补跑漏掉的调度(RecoverMissedSchedules)
它要解决的小问题: Scheduler 停机了两小时,期间有 4 个 9 点/10 点的 cron 点没跑。恢复后,这些漏掉的要不要补?
Kestra 给三种策略(core/src/main/java/io/kestra/core/models/triggers/RecoverMissedSchedules.java):
| 策略 | 行为 |
|---|---|
ALL | 全部补跑(默认) |
LAST | 只补最后一个漏掉的 |
NONE | 一个都不补,从现在起 |
策略在 Scheduler 启动/接管 vNode 时生效,由 TriggerScheduler.onStart 处理(scheduler/TriggerScheduler.java:210-231):
LAST: 用previousEvaluationDate算出"上一个 cron 点",若它比状态里的evaluatedAt晚,就把nextEvaluationDate拨到那个点——下一拍立刻补跑这一个(:208-214)。NONE: 从"现在"重新起算下一个点,直接跳过所有漏掉的(:215-224)。ALL: 什么都不做——因为状态里的nextEvaluationDate还停在很久以前,主循环会自然地一个接一个把它们全部追平(:225-227)。
默认值来自插件配置,取不到就是 ALL(core/src/main/java/io/kestra/core/models/triggers/Schedulable.java:47-52)。
3.8 回填(Backfill)
回填 = 主动为一段过去的时间区间补造执行。 比如"帮我把上个月每天的报表都补跑一遍"。
配置见 Backfill(core/src/main/java/io/kestra/core/models/triggers/Backfill.java):start/end 是区间,currentDate 是"回填进行到哪了",paused 可暂停,还能带 inputs/labels。
工作方式(inferred,综合 TriggerEventHandler.onCreateBackfill 与 TriggerState):创建回填时,先把当前的 nextEvaluationDate 存进 previousNextExecutionDate 备份(TriggerState.java:229-241),然后把评估时间点拨到回填区间里,让主循环像追赶漏掉的调度一样逐点补跑。删除回填时(onDeleteBackfillTrigger,TriggerEventHandler.java:189)恢复备份的 nextEvaluationDate,回到正常节奏。每往前走一格,getBackFillForNextEvaluationDate(TriggerState.java:338)推进 currentDate,越过 end 就把回填清空,回填结束。
4. 深入实现
4.1 分布式分片:vNode 是怎么回事
问题: 多个 Scheduler 实例同时在线,同一个触发器不能被两个实例都评估(会造出两个执行)。
Kestra 的解法是虚拟节点(vNode)分片。 每个触发器按 Flow 哈希到一个 vNode(VNodes.computeVNodeFromFlow,TriggerScheduler.java:173);所有 vNode 通过一致性哈希环分给在线的 Scheduler 实例,一个 vNode 同一时刻只属于一个实例。于是"同一个触发器只被一个实例评估"就有了保证。
DefaultScheduler 订阅 vNode 的分配与撤销(scheduler/DefaultScheduler.java:156-194):
- 被分配 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 缓存:
partitionedCache是vNode -> 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,差别是"知不知道条件":
| 重载 | 签名 | 用在哪 | 会不会考虑触发器条件 |
|---|---|---|---|
带 ConditionContext | get(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:104)统一处理,内部 switch 分派(:116-135):
| 事件 | 含义 | 处理要点 |
|---|---|---|
TriggerCreated | 新触发器 | 建状态,算首个 nextEvaluationDate(:506) |
TriggerUpdated | 定义变了 | 杀掉在跑的实例,用新定义重算(:429) |
TriggerEvaluated | Worker 评估回执 | 解锁/写回,有执行就发出(:312,见 §4.5) |
TriggerExecutionTerminated | 执行结束 | 解锁,按终态判断是否 stopAfter 禁用(:253) |
TriggerReceived | Worker 已接手 | 记 workerId(:354) |
TriggerWorkerLost | 持有触发器的 Worker 没了 | 解锁,让触发器重新可被派发(:386) |
CreateBackfill/SetPauseBackfill/DeleteBackfill | 回填命令 | 见 §3.8 |
ResetTrigger/SetDisableTrigger | 重置/启停 | 重算或切换禁用 |
去重是这里的关键。 大多数队列是"至少一次"投递,同一事件可能来两遍。findTriggerState(:557)在应用事件前比对 lastEventId:只有当前事件"比上次应用的更新"才处理,否则丢弃(:566-573):
// scheduler/TriggerEventHandler.java:635-642(节选)
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/src/main/java/io/kestra/worker/processors/internals/WorkerTriggerCallable.java:33)直接调pollingTrigger.eval(conditionContext, triggerContext)——真正的"查 S3、发 HTTP"在这里发生,在 Worker 线程上,不占用 Scheduler。 - 结果回流: Worker 把
WorkerTriggerResult发回,worker-controller 把它翻译成事件(worker-controller/src/main/java/io/kestra/controller/grpc/services/GrpcWorkerControllerService.java:247):有评估结果就发TriggerEvaluated,评估失败(启动不了)就发TriggerExecutionTerminated(FAILED)。 - 回执处理:
onTriggerEvaluated(TriggerEventHandler.java:383)——有执行则写回状态并triggerExecutionPublisher.send(execution)(:342-345);没有执行(轮询什么都没匹配到,或作业在派发前被拒)则解锁(:330-337),否则这个触发器会永远卡在锁里再也不被评估。
为什么 Realtime 的解锁要靠 dispatchEpoch 而不是 lastEventId? 一个在跑的 Realtime 触发器会吐出大量执行,它们的"终止"信号不能释放触发器的锁(否则会把还在 Worker 上跑的触发器重新派一遍)。Realtime 唯一合法的解锁信号是"启动失败"那个 FAILED 执行(TriggerEventHandler.java:351-376)。而这些回执来自多个独立的生产者、事件 id 不保证有序,所以这里改用派发世代 dispatchEpoch 做栅栏:回执的世代比当前状态旧就丢弃(:297-298)。
4.6 Flow 触发:被"别的执行"触发
Flow 触发器根本不经过 Scheduler。 它由"某个执行进入某状态"驱动,发生在 Executor 一侧,逻辑在 FlowTriggerService(executor/FlowTriggerService.java:39)。
每当一个执行状态变化,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/src/main/java/io/kestra/core/models/conditions。
5. 巧妙之处(可借鉴的技术)
- "先排好下一拍,再评估" —— 无论评估成功、失败还是没产出,
nextEvaluationDate都已在评估前写回(TriggerScheduler.java:371-372)。触发器永远有"下一次",不会因一次异常永久卡死。 - 加锁 +
locked=false过滤,天然防重叠 —— 不允许并发的触发器一旦派出去就加锁,而"挑到期"的查询只取未加锁的(DefaultSchedulableTriggerFetcher.java:68)。锁的存在与否直接实现了allowConcurrent语义,无需额外协调。 dispatchEpoch世代栅栏 —— 面对"至 少一次"投递 + 多生产者 + 无序回执,用一个单调递增的派发号区分"当前派发的回执"和"被取代的旧回执"(TriggerEventHandler.java:368、TriggerState.java:332),比依赖事件顺序更稳。lastEventId事件去重 —— 状态里记住最后应用的事件 id,重复/过期事件直接丢(TriggerEventHandler.java:636),让整个事件处理具备幂等性。- 按 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:19-22)。 - Realtime 靠 Worker 常驻,Worker 掉了要靠事件恢复。 Worker 丢失通过
TriggerWorkerLost解锁重派(TriggerEventHandler.java:456),但这依赖丢失检测与事件投递,不是瞬时的。 - 杀不掉的 Realtime 会拖到重启。 若无法向 Realtime 触发器发送 kill(队列异常),它会一直跑到 Kestra 重启(
TriggerEventHandler.java:636)。 - Schedule 的评估占用调度线程。 Schedule 在 Scheduler 内评估(§3.4),虽然轻,但极端复杂的条件渲染仍会挤占调度节奏——这是"不外派"换来的代价。
- 回填/补跑本质是"把时间往回拨让主循环追", 区间很大时会产生大量执行,需注意对下游的压力。
7. 横向对比
本章是 Kestra 子库的一章,和同组其它章的边界:
- 触发器产出的
Execution之后如何推进,见执行引擎:Executor 状态机。本章只负责"造出来"。 - Polling/Realtime 的评估在 Worker 上跑,Worker 与
RunContext的机制见任务执行:Worker 与 RunContext。 Flow/Trigger/Execution等领域模型的定义见领域模型。- 触发器事件、Worker 作业都走队列传递,队列抽象与持久化见消息骨架:Queue 抽象与持久化;触发器状态的 JDBC 落库属于其中。
Schedule/Webhook等具体 触发器是插件,插件体系见扩展与接入:插件系统与 Web 层。
8. 代码地图(导航索引)
| 主题 | 文件路径 | 关键符号 |
|---|---|---|
| 调度服务本体、vNode 分配 | scheduler/src/main/java/io/kestra/scheduler/DefaultScheduler.java | DefaultScheduler、start、onVNodesAssigned |
| 每秒调度主循环 | scheduler/src/main/java/io/kestra/scheduler/TriggerSchedulingLoop.java | run、processTriggerEvents、SCHEDULE_INTERVAL_MILLIS |
| 调度核心:挑选/评估/派发 | scheduler/src/main/java/io/kestra/scheduler/TriggerScheduler.java | onSchedule、evaluate、processSchedulableTrigger、processWorkerTrigger、onStart |
| 挑出到期且有效的触发器 | scheduler/src/main/java/io/kestra/scheduler/internals/DefaultSchedulableTriggerFetcher.java | getSchedulableTriggers |
| Schedule 类在 Scheduler 内评估 | scheduler/src/main/java/io/kestra/scheduler/internals/SchedulableEvaluator.java | SchedulableEvaluator.evaluate |
| 算下一个评估时间点 | scheduler/src/main/java/io/kestra/scheduler/internals/NextEvaluationDate.java | NextEvaluationDate.get |
| 触发器事件处理与去重 | scheduler/src/main/java/io/kestra/scheduler/TriggerEventHandler.java | handle、onTriggerEvaluated、findTriggerState |
| 派发触发器到 Worker 队列 | scheduler/src/main/java/io/kestra/scheduler/pubsub/TriggerWorkerJobPublisher.java | TriggerWorkerJobPublisher.send |
| 触发器状态(不可变) | core/src/main/java/io/kestra/core/scheduler/model/TriggerState.java | TriggerState、updateForNextEvaluationDate、nextDispatchEpoch、updateOnExecutionTerminated |
| 状态库接口 | core/src/main/java/io/kestra/core/scheduler/store/TriggerStateStore.java | findTriggersEligibleForScheduling |
| 分片缓存装饰器 | scheduler/src/main/java/io/kestra/scheduler/stores/CachedTriggerStateStore.java | CachedTriggerStateStore、partitionedCache、init |
| 触发器基类与语义 | core/src/main/java/io/kestra/core/models/triggers/AbstractTrigger.java | AbstractTrigger、when、allowConcurrent、stopAfter |
| 轮询触发器接口 | core/src/main/java/io/kestra/core/models/triggers/PollingTriggerInterface.java | eval、getInterval、nextEvaluationDate |
| 实时触发器接口 | core/src/main/java/io/kestra/core/models/triggers/RealtimeTriggerInterface.java | RealtimeTriggerInterface.eval |
| 定时触发器接口 | core/src/main/java/io/kestra/core/models/triggers/Schedulable.java | Schedulable、previousEvaluationDate、defaultRecoverMissedSchedules |
| 触发上下文 | core/src/main/java/io/kestra/core/models/triggers/TriggerContext.java | TriggerContext、nextExecutionDate、backfill |
| 回填模型 | core/src/main/java/io/kestra/core/models/triggers/Backfill.java | Backfill、currentDate、previousNextExecutionDate |
| 补跑策略 | core/src/main/java/io/kestra/core/models/triggers/RecoverMissedSchedules.java | RecoverMissedSchedules(ALL/LAST/NONE) |
| 执行/回执构建工具 | core/src/main/java/io/kestra/core/models/triggers/TriggerService.java | generateEvaluationResult、buildLabels |
| Worker 侧轮询评估 | worker/src/main/java/io/kestra/worker/processors/internals/WorkerTriggerCallable.java | WorkerTriggerCallable.doCall |
| Worker 回执→事件翻译 | worker-controller/src/main/java/io/kestra/controller/grpc/services/GrpcWorkerControllerService.java | sendWorkerTriggerResults |
| Flow 触发(由执行驱动) | executor/src/main/java/io/kestra/executor/FlowTriggerService.java | computeExecutionsFromFlowTriggerConditions、computeExecutionsFromFlowTriggerDependsOn |
| Schedule 插件(Schedulable 实例) | core/src/main/java/io/kestra/plugin/core/trigger/Schedule.java | nextEvaluationDate、previousEvaluationDate、eval |