跳到主要内容

扩展与接入:插件系统与 Web 层

30 秒导读: Kestra 对外有两张脸。一张朝下,对插件生态——每个 task / trigger / condition 都是一个插件类,YAML 里写 type: io.kestra.plugin.xxx.Foo 时,Kestra 靠一张「插件注册表」把这行字符串解析成真实的 Java 类。另一张朝上,对使用者——一组 REST 控制器(FlowControllerExecutionController 等)把 Flow / Execution 这些模型开成 HTTP API,给 UI 和外部系统调用。本章讲清这两张脸,重点是那条 type 字符串 → 注册表 → 实现类 的接线。

领域模型(Flow / Task / Execution 是什么)在第 1 章已经讲过,本章不重复;这里只关心「这些类型怎么被扩展、怎么被对外操作」。


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

1.1 一句话定义

  • 插件系统: 让 Kestra 的能力(能跑哪些 task、能被哪些事件触发)可以由外部 jar 包扩展,而不必改 Kestra 本体。
  • Web 层: 让用户不必写代码,就能通过 HTTP API(以及背后的 Web UI)去创建 flow、触发 execution、查插件文档。

1.2 为什么需要插件系统

Kestra 是个编排引擎,它自己不该知道「怎么连 BigQuery」「怎么发 Slack 消息」。这些是无穷无尽的第三方集成。

于是 Kestra 定了一个约定:凡是能被 flow 用到的东西,都是插件。Kestra 核心只认识几个基类/接口(TaskAbstractTriggerTaskRunner……),具体的 BigQueryQuerySlackExecution 由独立的插件 jar 提供。用户在 flow 的 YAML 里这样写:

tasks:
- id: hello
type: io.kestra.plugin.core.log.Log # 这一行就是「我要用哪个插件类」
message: "Hello World"

那行 type: 就是插件系统的入口。Kestra 要做的事,是把这个字符串可靠地翻译成一个真实的 Java 类,并且这个类可能来自一个和 Kestra 核心互相隔离的 jar。

1.3 为什么需要 Web 层

Flow 是一份 YAML,Execution 是一次运行记录。用户要能:提交一份 flow、手动触发一次运行、查看某个插件长什么样(它有哪些参数)。这些都通过 REST API 完成——FlowControllerExecutionControllerTriggerControllerPluginController 就是这些 API 的落点。UI 只是这些 API 的一个消费者。

1.4 一句话直觉

插件系统像「USB 接口 + 设备管理器」: Kestra 定义了插槽形状(基类),厂商造出符合形状的设备(插件 jar),设备管理器(注册表)扫描出所有插好的设备,并按型号(type 字符串)找到对应驱动(Java 类)。

Web 层像「前台接待」: 你不用进后厨(核心引擎),在前台(REST 控制器)递单子(HTTP 请求),它转交给厨房(Service 层)去做。


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

两条主线,一条向下(插件解析),一条向上(HTTP 接入),它们在领域模型插件注册表这两个共享物件上交汇。

┌──────────────────────────────────────────────┐
HTTP 请求 ───▶│ Web 层(webserver 模块) │
(UI / curl) │ FlowController / ExecutionController / │
│ TriggerController / PluginController │
│ │ 只做:收参数、鉴权、转交 │
│ ▼ │
│ Service 层(FlowService / ExecutionService…) │
└────────┬─────────────────────────────────────┘
│ 提交的 flow 是 YAML 文本

┌──────────────────────────────────────────────┐
│ 反序列化(JacksonMapper + PluginModule) │
│ 遇到 Task/Trigger 字段 → PluginDeserializer │
└────────┬─────────────────────────────────────┘
│ 拿 YAML 里的 type 字符串

┌──────────────────────────────────────────────┐
│ 插件注册表 PluginRegistry │
│ type 字符串 ──▶ Class<? extends Plugin> │
└────────┬─────────────────────────────────────┘
│ 注册表内容从哪来?

┌──────────────────────────────────────────────┐
│ PluginScanner + PluginClassLoader │
│ 扫描 core + 外部 jar,ServiceLoader 发现类 │
│ (外部 jar 各自一个隔离类加载器) │
└──────────────────────────────────────────────┘

2.1 主要部件一句话职责

部件干什么在哪个文件(相对克隆根)
Plugin所有插件的顶层接口,定义 getType()model/src/main/java/io/kestra/core/models/Plugin.java
PluginRegistry注册表接口:type 字符串 → 类core/src/main/java/io/kestra/core/plugins/PluginRegistry.java
DefaultPluginRegistry注册表默认实现(单例、并发 Map)core/src/main/java/io/kestra/core/plugins/DefaultPluginRegistry.java
PluginScanner扫描类加载器,用 ServiceLoader 发现插件core/src/main/java/io/kestra/core/plugins/PluginScanner.java
PluginClassLoader每个外部 jar 的隔离类加载器(child-first)core/src/main/java/io/kestra/core/plugins/PluginClassLoader.java
RegisteredPlugin一个插件 bundle 的扫描结果(按种类分组的类列表)core/src/main/java/io/kestra/core/plugins/RegisteredPlugin.java
PluginDeserializerJackson 反序列化器,把 type 接到注册表core/src/main/java/io/kestra/core/plugins/serdes/PluginDeserializer.java
PluginManager / LocalPluginManager安装/卸载插件 jar,触发重新注册core/src/main/java/io/kestra/core/plugins/PluginManager.java
MavenPluginDownloader从 Maven 仓库把插件 jar 拉到本地core/src/main/java/io/kestra/core/plugins/MavenPluginDownloader.java
REST 控制器把领域模型对外开成 HTTP APIwebserver/src/main/java/io/kestra/webserver/controllers/api/*.java

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

3.1 关键接线:type 字符串 → 注册表 → 实现类

这是整章最重要的一节。它回答一个问题:YAML 里一行 type: io.kestra.plugin.core.log.Log,是怎么变成一个真实 Java 对象的?

它要解决的小问题

Jackson 反序列化 JSON/YAML 时,通常需要在编译期就知道目标类。但 Kestra 的 task 类编译期未知——它可能来自一个运行时才装上的外部 jar。标准的 @JsonSubTypes 注解列举法在这里行不通,因为你不可能预先列出所有第三方插件。

思路

Kestra 换了个招:不在类型上注解子类,而是给几个基类挂一个自定义反序列化器。这个反序列化器自己读 type 字段,再去问注册表要类。

第一步,PluginModulePluginDeserializer 绑到每个插件基类上:

// PluginModule.java:37-50 —— 一个 Jackson SimpleModule
addDeserializer(Task.class, new PluginDeserializer<>());
addDeserializer(AbstractTrigger.class, new PluginDeserializer<>());
addDeserializer(TaskRunner.class, new PluginDeserializer<>());
addDeserializer(StorageInterface.class, new PluginDeserializer<>());
// …Chart / DataFilter / SecretPluginInterface / LogExporter 等一律如此

这个 module 在全局 JacksonMapper 里被注册(JacksonMapper.java:159 registerModule(new PluginModule())),所以凡是解析到 Task 类型的字段,就走 PluginDeserializer

第二步,PluginDeserializer 读出 type,交给注册表:

// PluginDeserializer.java:85-101 fromObjectNode(...) 核心三行
final String identifier = extractPluginRawIdentifier(node, pluginRegistry.isVersioningSupported());
pluginType = pluginRegistry.findClassByIdentifier(identifier); // 字符串 → Class
// …
return (T) jp.getCodec().treeToValue(node, pluginType); // 再用真实类去反序列化整个节点

extractPluginRawIdentifier(PluginDeserializer.java:156)就是简单地取 node.get("type") 的文本值(支持版本时拼成 type:version)。

第三步,注册表按类名查表:

// DefaultPluginRegistry.java:315-332 findClassByIdentifier / findMetadataByIdentifier
// 内部就是一次并发 Map 查找:
return Optional.ofNullable(pluginClassByIdentifier.get(identifier));

pluginClassByIdentifier 是一个 Map<PluginIdentifier, PluginClassAndMetadata>(DefaultPluginRegistry.java:44),key 用的是 ClassTypeIdentifier——本质就是把全限定类名当 key(DefaultPluginRegistry.java:367 ClassTypeIdentifier)。

整条链一图流

YAML: type: io.kestra.plugin.core.log.Log

▼ Jackson 解析到 Task 字段
PluginDeserializer.deserialize (serdes/PluginDeserializer.java:60)
│ 取出 "io.kestra.plugin.core.log.Log"

pluginRegistry.findClassByIdentifier(...) (DefaultPluginRegistry.java:307)
│ Map 查表:类名 → Class

Class<Log> pluginType
│ jp.getCodec().treeToValue(node, Log.class)

真实的 Log 对象(message 等字段已填好)

关键细节 / 坑

  • 别名(alias)也进同一张表。 插件类可以用 @Plugin(aliases=...) 声明旧类名,扫描时这些别名也被塞进注册表(RegisteredPlugin.java:99-112 hasClass / findClass),所以改了包名的老 flow 仍能解析。
  • @JsonDeserialize 的自反陷阱。 源码注释明确警告(PluginDeserializer.java:137-138):如果具体插件类没有@JsonDeserialize 覆写,treeToValue 会再次进到 PluginDeserializer,导致 StackOverflowException。这是「基类走自定义反序列化、具体类走默认反序列化」这套设计的暗礁。
  • 找不到就报错。 type 查不到类时抛 invalidTypeIdException,消息是 No plugin registered for the defined type(PluginDeserializer.java:147-154)。
  • 这条链是引擎和 Web 层共用的:ExecutorWorker 反序列化 flow 时走的是同一套。

3.2 类加载隔离:child-first 的 PluginClassLoader

它要解决的小问题

不同插件 jar 可能依赖同一个库的不同版本(比如两个插件各自带一个不同版本的某 SDK)。如果都塞进同一个类加载器,就会版本打架。

思路

给每个外部插件 bundle 一个独立的类加载器,并采用 child-first(子优先) 策略:加载类时先在插件自己的 jar 里找,找不到才委托给父加载器。这样插件的私有依赖优先生效,互不干扰。

// PluginClassLoader.java:94-123 loadClass 的委托顺序(简化)
Class<?> loadedClass = findLoadedClass(name); // 1. 已加载?
if (loadedClass == null && shouldLoadFromUrls(name)) { // 2. 不在排除名单?
loadedClass = findClass(name); // 先在插件 jar 里找
}
if (loadedClass == null) {
loadedClass = super.loadClass(name, resolve); // 3. 兜底:委托父加载器
}

关键细节:哪些包不能子优先

child-first 不能一刀切——系统类、Kestra 核心类、日志/Jackson 等必须由父加载器统一提供,否则会出现「同名类来自两个加载器」的 LinkageError,或插件冒充系统类的安全问题。所以有一份排除名单:

// PluginClassLoader.java:29-52 EXCLUDES(节选)
"^(?:java|javax|jakarta"
+ "|io.kestra.core|io.kestra.plugin.core" // 核心必须共享,否则 type 匹配会错乱
+ "|org.slf4j|ch.qos.logback" // 日志要统一
+ "|com.fasterxml.jackson.core|…databind" // Jackson 要统一,否则反序列化跨加载器崩
+ "…)\\..*$"

shouldLoadFromUrls(PluginClassLoader.java:132)就是「不在这份名单里 → 才允许子优先」。为什么 io.kestra.core 必须共享: 3.1 里 baseClass 要判断插件类是不是 Task.class 的子类——如果 Task 在插件加载器里是另一份拷贝,isAssignableFrom 会返回 false,整条接线断掉。

一句话记住这张排除名单:「凡是引擎和插件要交换的类型,都必须来自同一个加载器」。

3.3 扫描与注册:ServiceLoader 发现插件

它要解决的小问题

装上一个插件 jar 后,Kestra 怎么知道里面有哪些 task、哪些 trigger?

思路

用 Java 标准的 ServiceLoader 机制(jar 里 META-INF/services/ 声明实现),把每个实现按它属于哪个基类分门别类。

// PluginScanner.java:143-157 扫描一个类加载器
final ServiceLoader<Plugin> sl = ServiceLoader.load(Plugin.class, classLoader);
for (Plugin plugin : sl) {
switch (plugin) {
case Task task -> tasks.add(task.getClass()); // 归到 tasks
case AbstractTrigger trigger -> triggers.add(trigger.getClass()); // 归到 triggers
case StorageInterface s -> storages.add(s.getClass());
// …十几种插件类别,各归各的 List
}
}

扫描结果打包成一个 RegisteredPlugin(RegisteredPlugin.java:60-79),它就是「一个 bundle 里所有类,按种类分好组」的容器,并从 jar 的 MANIFEST.MF 里读出组名、版本、图标等元数据(RegisteredPlugin.java:262 group():300 version())。

两种扫描入口

入口扫什么用什么类加载器源码
scan()核心内置插件(io.kestra.plugin.core)Kestra 自己的类加载器PluginScanner.java:93
scan(Path)某目录下的外部插件 jar每个 bundle 一个 PluginClassLoaderPluginScanner.java:59

外部扫描是并行的(parallelStream(),PluginScanner.java:63),每个 bundle 独立建加载器、独立扫描,最后过滤掉空 bundle(isValid(),RegisteredPlugin.java:81)。

注册表怎么攒起来

DefaultPluginRegistry 是个懒加载单例(getOrCreate(),DefaultPluginRegistry.java:57):

  • 启动时 init() 先扫核心(DefaultPluginRegistry.java:75-79)。
  • 装插件时 register(Path) 扫目录并逐个登记(DefaultPluginRegistry.java:120-128)。
  • register(RegisteredPlugin)(DefaultPluginRegistry.java:195)把 bundle 里每个类摊平成 类名 → PluginClassAndMetadata 存进那张并发 Map,这一步正是 3.1 查表的数据来源
  • 重复注册用 CRC32 去重(existing.crc32() == plugin.crc32() 就跳过,DefaultPluginRegistry.java:199);替换旧 bundle 时会关掉旧的类加载器释放资源(unregister,DefaultPluginRegistry.java:163-169)。

3.4 动态安装:MavenPluginDownloaderPluginManager

它要解决的小问题

插件 jar 不一定预先放在磁盘上——用户可能只给一个 Maven 坐标(groupId:artifactId:version)。要能从远程仓库把它拉下来再装。

思路

MavenPluginDownloader 用 Maven 的 Aether(org.eclipse.aether)解析坐标、下载 jar 到本地仓库:

// MavenPluginDownloader.java:205-226 resolveArtifact(简化)
DefaultArtifact artifact = new DefaultArtifact(dependency); // 解析坐标
VersionRangeResult version = system.resolveVersionRange(...); // 找最高版本
ArtifactResult result = system.resolveArtifact(session, artifactRequest); // 下载 jar
return new PluginArtifact(..., result.getArtifact().getPath().toUri()); // 返回本地路径

PluginManager(接口在 PluginManager.java,OSS 实现是 LocalPluginManager)负责 install / uninstall,并在装完后触发注册表重新扫描:

// LocalPluginManager.java:140-141 装完就注册
if (installForRegistration && pluginRegistryProvider != null) {
pluginRegistryProvider.get().register(localRepositoryPath); // 回到 3.3 的 register(Path)
}

这就把 3.4(下载安装)接回了 3.3(扫描注册)和 3.1(查表解析)——一个新插件从坐标到能在 YAML 里 type: 用起来的完整闭环


4. Web 层:REST 接口这张脸

4.1 分层铁律:控制器只做「收 + 转」,业务在 Service

Kestra 强制控制器里不写业务规则,控制器只负责:解析请求参数 → 鉴权 → 调 Service → 返回 JSON。业务逻辑一律下沉到 Service 类。

以「创建 flow」为例,控制器方法极短:

// FlowController.java:270-278 createFlow → doCreate → flowService.create
public HttpResponse<FlowWithSource> createFlow(@Body String flow) {
return HttpResponse.ok(doCreate(parseFlowSource(flow)));
}
// doCreate 里真正干活的是:
return flowService.create(flow); // 业务全在 FlowService

控制器通过构造/字段注入拿到一堆 Service。各控制器持有的协作者:

控制器挂载路径主要 Service / 依赖
FlowController/api/v1/{tenant}/flowsFlowServiceFlowRepositoryInterfaceModelValidatorGraphService
ExecutionController/api/v1/{tenant}/executionsExecutionServiceFlowInputOutputExecutionRepositoryInterface、多条 Queue
TriggerController/api/v1/{tenant}/triggersTriggerStateServiceTriggerRepositoryInterface
PluginController/api/v1/plugins/PluginRegistryJsonSchemaGeneratorJsonSchemaCache

依据:FlowController.java:75-91ExecutionController.java:155-183TriggerController.java:71-88PluginController.java:56-74

注意路径里的 {tenant}——除 PluginController 外,资源类 API 都带租户前缀做多租户隔离;PluginController 无租户,因为插件是全局的。

4.2 每个控制器对外提供什么

控制器代表性端点(动作)源码锚点
FlowController建 flow、查 flow、查 revision、搜索、生成拓扑图FlowController.java:268(create)、:172(get)、:117(graph)
ExecutionController触发执行、webhook 触发、重启/重放、改状态、删除ExecutionController.java:676(createExecution)、:524(webhook)、:1022(restart)
TriggerController搜索 trigger、解锁、backfill 建/停、禁用TriggerController.java:93(search)、:162(unlock)、:222(backfill)
PluginController取插件 JSON schema、取文档、按分组列插件、图标PluginController.java:76(schema)、:311(doc)、:360(subgroups)

PluginController 是本章两张脸的交汇点:它直接注入 PluginRegistry,把 3.x 攒出来的注册表内容(插件列表、schema、文档、图标)开放给 UI——UI 的「插件面板」「YAML 自动补全」就是消费这些端点。

4.3 三条横切约定:OpenAPI、DTO、异步 IO

  • OpenAPI 文档化。 每个端点用 @Operation(tags=..., summary=...) 描述,Swagger 注解在编译期生成 OpenAPI 规范(如 PluginController.java:78-82)。UI 和外部 SDK 都靠它。
  • DTO 而非裸模型。 响应用专门的传输对象:分页统一 PagedResults<T>(TriggerController.java:95 返回 PagedResults<ApiTriggerAndState>),插件文档用 DocumentationWithSchema。约定还要求「API 不能直接返回 JSON 数组」(数组无法向后兼容地演进),所以列表也包在对象里。
  • 阻塞操作让出 IO 线程池。 几乎每个端点都挂 @ExecuteOn(TaskExecutors.IO)(如 FlowController.java:105),把阻塞的仓库/文件操作调度到 IO 线程,避免堵住 Micronaut 的事件循环。

鉴权说明:OSS 版这几个控制器上看不到 @Secured / @HasAnyPermission 注解(grep 为空)。细粒度权限(RBAC、tenant 隔离到人)是企业版(EE)的能力;OSS 主要靠 {tenant} 路径做租户维度隔离,方法级权限注解在本克隆里不体现 (inferred,依据:控制器源码中无此类注解)。


5. 前端一句话概览(指路,不深入)

UI 在 ui/ 目录,是 Vue 3 + TypeScript + Vite 的单页应用(依据:ui/package.jsonvue-tsc / vite 脚本、ui/src/App.vue 等),用 Element Plus 封装出自有设计系统。它的核心是可视化的 flow 编辑器 + 拓扑图:用户在编辑器里写 YAML,前端调 PluginController 的 schema/图标端点做补全和渲染,调 FlowController/ExecutionController 做增删改查与触发。前端本身不含引擎逻辑,只是本章 Web 层 API 的最大消费者。 深入到组件层不在本章范围。


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

  • 用「自定义反序列化器 + 注册表」替代「编译期子类枚举」。 这是让声明式 YAML 能接到运行时才存在的类的关键手法。@JsonSubTypes 需要预先列举,而 Kestra 的插件是开放集合——把类型解析推迟到运行时查表,才撑得起插件生态(PluginDeserializer.java:85)。
  • child-first + 精确排除名单。 既给插件私有依赖隔离,又保证「引擎与插件交换的类型」共享同一加载器。这份 EXCLUDES 名单是这套隔离能正常工作的核心——少列一个 io.kestra.core 就会满盘皆输(PluginClassLoader.java:29)。
  • CRC32 去重 + 类加载器关闭。 注册表用 bundle 的 CRC32 判断「同一个 jar 是否已注册」,替换时主动 close() 旧类加载器释放句柄(DefaultPluginRegistry.java:199:163),这对支持热插拔/重装很关键。
  • 控制器薄、Service 厚。 强约束「控制器不写业务」让同一套业务能被 REST、CLI、内部调用三处复用(FlowController.java:270 那种两行方法就是范例)。

7. 边界与局限(诚实)

  • OSS 注册表不支持插件版本共存。 DefaultPluginRegistry.isVersioningSupported() 直接返回 false(DefaultPluginRegistry.java:346),getAllVersionsForType 只会返回单个 manifest 版本(:85)。多版本并存是企业版能力。
  • 同名类冲突静默取第一个。 图标聚合等处用 toMap(..., (a1, a2) -> a1) 遇冲突留先到的(PluginController.java:257),不告警。
  • 具体插件类必须自带 @JsonDeserialize,否则栈溢出。 见 3.1 的坑,这是设计上的隐性契约,源码只用注释提醒(PluginDeserializer.java:137)。
  • 方法级鉴权在 OSS 控制器不可见。 见 4.3 说明,细粒度权限属 EE。
  • 插件下载依赖 Maven 仓库可达。 MavenPluginDownloader 走 Aether 联网解析,离线或坐标错时抛 KestraRuntimeException(MavenPluginDownloader.java:227-229)。

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

主题文件路径(相对克隆根)符号名
插件顶层接口 / getTypemodel/src/main/java/io/kestra/core/models/Plugin.javaPlugingetTypegetAliases
type→class 反序列化core/src/main/java/io/kestra/core/plugins/serdes/PluginDeserializer.javaPluginDeserializerfromObjectNodeextractPluginRawIdentifier
反序列化器注册到 Jacksoncore/src/main/java/io/kestra/core/plugins/PluginModule.javaPluginModule
全局 ObjectMapper 装配core/src/main/java/io/kestra/core/serializers/JacksonMapper.javaregisterModule(new PluginModule())
注册表接口core/src/main/java/io/kestra/core/plugins/PluginRegistry.javaPluginRegistryfindClassByIdentifier
注册表默认实现core/src/main/java/io/kestra/core/plugins/DefaultPluginRegistry.javaDefaultPluginRegistryregisterClassTypeIdentifierinit
ServiceLoader 扫描core/src/main/java/io/kestra/core/plugins/PluginScanner.javaPluginScannerscanscanClassLoader
隔离类加载器core/src/main/java/io/kestra/core/plugins/PluginClassLoader.javaPluginClassLoaderloadClassEXCLUDESshouldLoadFromUrls
扫描结果容器core/src/main/java/io/kestra/core/plugins/RegisteredPlugin.javaRegisteredPluginisValidbaseClassallClassGrouped
安装/卸载管理core/src/main/java/io/kestra/core/plugins/LocalPluginManager.javaLocalPluginManagerinstall
Maven 下载core/src/main/java/io/kestra/core/plugins/MavenPluginDownloader.javaMavenPluginDownloaderresolveArtifactlistAllVersions
Flow RESTwebserver/src/main/java/io/kestra/webserver/controllers/api/FlowController.javaFlowControllercreateFlow
Execution RESTwebserver/src/main/java/io/kestra/webserver/controllers/api/ExecutionController.javaExecutionControllercreateExecution
Trigger RESTwebserver/src/main/java/io/kestra/webserver/controllers/api/TriggerController.javaTriggerControllersearchTriggers
Plugin REST(两脸交汇)webserver/src/main/java/io/kestra/webserver/controllers/api/PluginController.javaPluginControllergetPluginBySubgroupsgetSchemasFromType

相关章节