扩展与接入:插件系统与 Web 层
30 秒导读: Kestra 对外有两张脸。一张朝下,对插件生态——每个 task / trigger / condition 都是一个插件类,YAML 里写
type: io.kestra.plugin.xxx.Foo时,Kestra 靠一张「插件注册表」把这行字符串解析成真实的 Java 类。另一张朝上,对使用者——一组 REST 控制器(FlowController、ExecutionController等)把 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 核心只认识几个基类/接口(Task、AbstractTrigger、TaskRunner……),具体的 BigQueryQuery、SlackExecution 由独立的插件 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 完成——FlowController、ExecutionController、TriggerController、PluginController 就是这些 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 |
PluginDeserializer | Jackson 反序列化器,把 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 API | webserver/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 字段,再去问注册表要类。
第一步,PluginModule 把 PluginDeserializer 绑到每个插件基类上:
// 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-112hasClass/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 层共用的:Executor 和 Worker 反序列化 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 一个 PluginClassLoader | PluginScanner.java:59 |
外部扫描是并行的(parallelStream(),PluginScanner.java:63),每个 bundle 独立建加载器、独立扫描,最后过滤掉空 bundle(isValid(),RegisteredPlugin.java:81)。