跳到主要内容

任务执行:Worker 与 RunContext

30 秒导读: 第 2 章的 Executor 是"大脑",它决定某个任务该跑了,就往队列里丢一个 WorkerJob本章讲的 Worker 是"手脚":它把这个 job 从队列里拉过来,在一个隔离的执行环境(RunContext)里真正调用插件代码 RunnableTask.run(),再把结果发回去。Worker 有一条铁律——它从不直接碰数据库,一切外部状态都走 gRPC 门面(MetaStore/StateStore)。


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

一句话定义: Worker 是 Kestra 里真正执行任务代码的进程

Executor 负责"推进状态机"(第 2 章),但它自己不跑任务——跑一个 Python 脚本、发一个 HTTP 请求、查一次数据库,这些活儿都在 Worker 上发生。这样设计是为了把"编排"和"执行"拆开:Executor 轻量、快、只碰状态;Worker 重、可能长时间阻塞、可能崩,但崩了不影响编排大脑。

给谁用 / 解决什么问题: 想象你写了一个 flow,里面有个任务 runScript。当轮到它跑时:

  • Executor 把它包装成一个 WorkerTask(一种 WorkerJob)扔进队列。
  • 某台 Worker 把它领走,在自己机器上建一个临时工作目录、渲染好所有 {{ }} 表达式、调用插件的 run()
  • 插件返回 outputs,Worker 把 WorkerTaskResult 发回给 Executor,状态机继续推进。

核心直觉/类比: 把 Executor 当餐厅前台(记单、叫号、更新订单状态),Worker 是后厨的厨师(真正炒菜)。前台永远不进厨房,厨师也永远不改订单系统——两者之间只通过"传菜口"(gRPC 流)递条子。

本章覆盖三大块:

讲什么关键类
分发骨架job 怎么从 controller 流到 workerWorkerJobDispatcherWorkerJobFetcher
作业循环worker 内部怎么逐个跑 jobWorkerLoopWorkerJobExecutorWorkerTaskProcessor
执行沙箱插件在什么环境里跑RunContextDefaultRunContext

本章不讲状态机推进(见第 2 章)和触发器如何产生 job(见第 4 章);队列抽象本身见第 5 章


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

先看一次任务从"被派发"到"结果回传"的完整数据流。怎么读这张图:从左到右是一次 job 的生命周期,虚线是控制信号(permits/completion),实线是数据(job/result)。

controller 侧 (worker-controller) worker 侧 (worker)
┌───────────────────────────────────┐ ┌──────────────────────────────────┐
│ WorkerJobDispatcher │ │ WorkerJobFetcher (一条 gRPC 双向流) │
│ · 订阅 job 队列 │◀──permits/completion(虚线)──│ │
│ · 按 permits+bucket 选一个 worker │ │ · 拉 job 放进本地 WorkerQueue │
│ · 持久化到 StateStore 再发 │──job──▶│ │
└───────────────────────────────────┘ └──────────────┬───────────────────┘
▲ │ poll
│ WorkerTaskResult ▼
┌─────────────────┴─────────────────┐ ┌──────────────────────────────────┐
│ GrpcWorkerControllerService │◀─send──│ WorkerJobExecutor │
│ · sendWorkerTaskResults→结果队列 │ │ · N 个 Consumer(虚拟线程)poll │
│ · sendWorkerLogEntries/Metrics │ │ · 每个 job 提交到平台线程池执行 │
└────────────────────────────────────┘ │ │ │
│ ▼ │
│ WorkerTaskProcessor │
│ → WorkerTaskCallable │
│ → task.run(RunContext) ★ │
└──────────────────────────────────┘
★ = 真正的插件代码在这里跑

部件一句话职责:

部件干什么文件
WorkerJobDispatchercontroller 侧总协调:选 worker、占 bucket、发 jobworker-controller/.../services/WorkerJobDispatcher.java
GrpcWorkerControllerServicegRPC 服务端:开双向流、收 permits、收各类结果worker-controller/.../services/GrpcWorkerControllerService.java
WorkerJobFetcherworker 侧拉 job 的那条流的客户端worker/.../fetchers/WorkerJobFetcher.java
WorkerJobExecutorworker 侧线程池 + N 个消费者,逐个跑 jobworker/.../WorkerJobExecutor.java
WorkerTaskProcessor把一个 WorkerTask 变成真正的执行:建 RunContext、重试、缓存worker/.../processors/WorkerTaskProcessor.java
WorkerTaskCallable最内层:真正调用 RunnableTask.run(),处理超时worker/.../processors/internals/WorkerTaskCallable.java
RunContext执行沙箱:变量渲染、存储、工作目录、日志、密钥core/.../runners/RunContext.java
GrpcWorkerIOSenderworker 侧把 result/log/metric 发回 controllerworker/.../senders/GrpcWorkerIOSender.java

主线走一遍(高层):

  1. Worker 启动时,WorkerJobFetcher 开一条到 controller 的 gRPC 双向流,报上自己的 id、group、最大并发数,以及初始 permits(可接活余量)
  2. WorkerJobDispatcher 一旦从队列拿到 job,就找一个有 permit 的 worker,先把 job 持久化到 StateStore(为了崩溃恢复),再通过流发过去。
  3. Worker 收到后放进本地 WorkerQueue;WorkerJobExecutor 的某个消费者 poll 到它,提交给线程池。
  4. WorkerTaskProcessor 建好 RunContext,WorkerTaskCallable 调用 task.run(runContext)——插件代码在这一刻执行
  5. 结果、日志、指标各走一条 GrpcWorkerIOSender 发回 controller;job 完成时 worker 还会回传一个 completion 信号,让 controller 释放占用的容量。

3. 分发骨架:job 怎么从 controller 流到 worker

这节讲任务与执行器之间的 gRPC 分发。核心难点是:controller 有很多 job,worker 有有限容量,怎么做到既不压垮 worker、又不丢 job。Kestra 的答案是基于 permits 的流控 + pull/ack 模式

3.1 一条长连接,双向流

Worker 侧只开一条 gRPC 双向流(streamWorkerJobs),它同时承载三种信息:

方向内容
worker → controllerpermits(剩余容量)、completion(哪些 job 到终态了)
controller → workerjob payload、广播事件(kill 命令、集群事件、元数据变更)

流的客户端是 WorkerJobFetcher——注意它本身是一个 WorkerLoop(见 §4.1),循环体 doOnLoop() 负责建流、按需重连(指数退避)、周期性上报 permits(worker/.../fetchers/WorkerJobFetcher.java:258 doOnLoop:293 startStream)。服务端在 GrpcWorkerControllerService.streamWorkerJobs 里为每条流建一个 WorkerStreamContext,并把它注册进 dispatcher(worker-controller/.../services/GrpcWorkerControllerService.java:86)。

3.2 permits:worker 说"我还能接几个"

permits 就是 worker 本地队列的剩余容量——一个数字,代表"你还可以再发几个 job 给我"。计算很直白:

// 示意,非源码:calculatePermits 的核心
private int calculatePermits() {
return fetchingPaused.get() ? 0 : workerJobQueue.remainingCapacity(); // 维护模式下报 0
}

真实实现在 WorkerJobFetcher.calculatePermits(worker/.../fetchers/WorkerJobFetcher.java:519)。关键点:permits 是一个电平值(level)而非增量,所以 0 是有意义的——它表示"别再发了"。维护模式/cordon 时 worker 就上报 0,让 controller 停止派发,但流保持存活(还能收 kill 命令)。

3.3 dispatcher 侧:选谁、占坑、发货

controller 从队列拿到一个 job 后,handleIncomingJob 会做三件事(worker-controller/.../services/WorkerJobDispatcher.java:711):

  1. 预检 kill: 如果这个 execution 已经被 kill,直接产出一个 KILLED 结果,不派发。
  2. 找 worker 并占坑: findAndReserveWorker 挑一个"最闲"(in-flight 最少)的候选,先 CAS 消费一个 permit,再占一个 bucket 槽位(:852)。两步都成功才算占到;bucket 占不到就把 permit 还回去,试下一个候选。
  3. 发货: dispatchJobToWorker 先持久化到 WorkerJobRunningStateStore 再 sendResponse(:888)。这个顺序是崩溃恢复的关键——只要 job 要么在队列里、要么在 StateStore 里,controller 崩了也不会丢。

permit 和 bucket 是两个独立的闸门,为什么要两个?

  • permit = worker 的总余量(所有队列共享),防止过量派发
  • bucket = 某个 Worker Queue 的容量槽,由 WorkerCapacityPolicy 管理,支持按队列预留容量(worker-controller/.../services/WorkerCapacityPolicy.java:27)。

默认策略 SinglePoolCapacityPolicy 最简单:一个 maxConcurrency 大小的共享池,占坑就是原子自增(worker-controller/.../services/SinglePoolCapacityPolicy.java:30 tryReserve)。需要"给某队列保底 N 个槽"的部署,替换 WorkerCapacityPolicyFactory 这个 bean 即可。

3.4 completion:占的坑什么时候还

一个微妙点:bucket 槽位从派发一直占到 job 真正跑完,而不是"worker 一收到就还"。job 到终态时,worker 侧的消费者会调 workerJobFetcher.onJobCompleted(jobId)(worker/.../WorkerJobExecutor.java:322),把 jobId 塞进 completion 队列、并立刻 flush 一次(worker/.../fetchers/WorkerJobFetcher.java:556)。controller 收到后在 onCompletionsReceived 里释放对应 bucket,并重新评估该队列的 pause/resume(WorkerJobDispatcher.java:668)。

为什么不等下一个心跳周期? completion 立刻 flush 能省掉最多一个 PERMIT_CHECK_INTERVAL(100ms)的空窗,让"同一 bucket 门控的下一个任务"更快派发出去。


4. Worker 作业循环:job 到了之后怎么跑

这节讲 worker 进程内部:从"本地队列里有个 job"到"插件代码被调用"这一段。这是层层嵌套的四层结构,每层职责单一。

4.1 底座:WorkerLoop —— 一个可暂停/可优雅停的循环

Worker 里所有"后台循环"(拉 job 的 fetcher、发结果的 sender、跑 job 的 consumer)都继承 WorkerLoop。它把"循环"这件事的通用骨架抽出来:

// 示意,非源码:WorkerLoop 的骨架
while (running.get()) {
waitIfPaused(); // 支持暂停(维护模式)
if (!running.get()) continue;
doOnLoop(); // 子类实现:干一次活
// 抛 InterruptedException → 退出;抛别的 → 记日志继续
}
cleanup(); // finally:释放资源

真实实现在 worker/.../WorkerLoop.java:88(runLoop)。子类只需实现 doOnLoop()(:121),外加可选的 signalJobStop()(:238,优雅停机时通知正在跑的活)。一个抽象类,三种用途(fetch / send / consume),是 worker 模块最重要的复用点。

4.2 分发线程:WorkerJobExecutor 与虚拟线程消费者

WorkerJobExecutor.start() 会起 N 个 WorkerJobConsumer,跑在虚拟线程上(worker/.../WorkerJobExecutor.java:82)。为什么是虚拟线程?因为消费者本身只做"poll + 等待",几乎不吃 CPU,虚拟线程廉价。

真正跑任务的却是平台线程池——消费者 poll 到 job 后,把它 submit 给一个 maxCachedThreadPool,然后自己 future.get() 阻塞等它跑完:

// 示意,非源码:WorkerJobConsumer.doOnLoop 的核心
WorkerJob job = workerJobQueue.poll(Duration.ofSeconds(1));
if (job == null) return;
WorkerJobProcessor processor = factory.create(context, job);
Future<?> future = taskExecutorService.submit(() -> processor.process(job)); // 提交到平台线程池
future.get(); // 消费者在此阻塞,直到任务跑完
workerJobFetcher.onJobCompleted(job.uid()); // finally:回传 completion

真实实现在 WorkerJobExecutor.WorkerJobConsumer.doOnLoop(worker/.../WorkerJobExecutor.java:283)。

为什么消费者(虚拟线程)和任务(平台线程)要分开? 这是中断隔离:任务超时/被 kill 时,只中断跑任务的那个平台线程,消费者的虚拟线程不受影响,能干净地继续 poll 下一个。任务里的异常也被 future.get() 完全兜住,消费者永远不会因为任务出错而挂掉(:310)。

4.3 处理器:WorkerJobProcessor 家族

WorkerJobProcessor 是一个接口,只有三个方法:process() / stop() / kill()(worker/.../processors/WorkerJobProcessor.java)。工厂按 job 类型选具体实现:

job 类型处理器干什么
WorkerTask(普通任务)WorkerTaskProcessor本章重点
WorkerTrigger(轮询触发)WorkerTriggerProcessor第 4 章

AbstractWorkerJobProcessor 用一对 CAS 保证一个处理器同时只跑一个 job,并管理 kill/stop 信号的转发(worker/.../processors/AbstractWorkerJobProcessor.java:52)。有个精巧的竞态处理:如果 kill 信号在 callable 被记录之前就到了,callJob 会在设置 callable 后补一次 kill(:69),防止 kill 丢失。

WorkerTaskProcessor.doProcess 分三种情况(worker/.../processors/WorkerTaskProcessor.java:98):

  • RunnableTaskrunTask(),正常执行一个任务。
  • WorkingDirectory(流程控制任务)→ runWorkingDirectory(),在同一个工作目录里顺序跑一串子任务,子任务间用 runIf 决定跳过。
  • 其它 → 报错(worker 只跑 runnable task)。

runTask 里还夹着输出缓存逻辑:如果任务开了 taskCache,先按渲染后的任务定义算 SHA-256 哈希查缓存,命中就跳过执行直接返回缓存的 outputs(:225)。

4.4 最内层:WorkerTaskCallable —— 真正调 run() 的地方

runAttempt 会先发一个"任务进入 RUNNING"的结果(让 Executor 知道任务动了),然后构造 WorkerTaskCallable 并执行(worker/.../processors/WorkerTaskProcessor.java:362:396)。

WorkerTaskCallable.doCall 就是那条"最深的线":

// 示意,非源码:WorkerTaskCallable.doCall 的核心
Duration timeout = runContext.render(workerTask.getTask().getTimeout())...;
if (timeout != null) {
Failsafe.with(Timeout.builder(timeout).withInterrupt().build())
.run(() -> taskOutput = task.run(runContext)); // ★ 插件代码在这里跑
} else {
taskOutput = task.run(runContext); // ★
}
return taskOutput.finalState().orElse(SUCCESS);

真实实现在 worker/.../processors/internals/WorkerTaskCallable.java:60RunnableTask.run(RunContext) 就是所有插件的执行入口(core/.../models/tasks/RunnableTask.java:14)——整章的"手脚"最终落在这一行 task.run(runContext)(:84)。

它的父类 AbstractWorkerCallable.call() 有个关键决定:它 catch Throwable(不只是 Exception),因为插件可能抛 Error(依赖问题、坏行为),必须保证任务无论如何都会以某个状态收场,而不是让 worker 线程裸崩(worker/.../processors/internals/AbstractWorkerCallable.java:63)。超时/kill 通过中断当前线程实现(:109 kill)。


5. RunContext:执行沙箱

task.run(runContext) 里那个 runContext 是插件与 Kestra 世界的唯一接口。插件不能随便碰文件系统、数据库、密钥——它能做的一切都通过 RunContext 暴露。这节讲这个沙箱提供了哪些能力。

RunContext 是抽象类(core/.../runners/RunContext.java:24),实际用的是 DefaultRunContext(core/.../runners/DefaultRunContext.java:53)。它在 worker 侧由 RunContextInitializer.forWorker 构造:从 wire 数据重建变量、装上存储、日志、工作目录、插件配置(core/.../runners/RunContextInitializer.java:85:122)。

四大能力如下。

5.1 变量与 Pebble 渲染

flow 里到处是 {{ execution.id }}{{ inputs.name }} 这种模板。渲染就是把这些占位符替换成真实值。

变量从哪来? RunVariables.DefaultBuilder.build 把一次执行的所有上下文组装成一个大 map:flowexecutiontaskruninputsoutputsvarslabelsenvs… (core/.../runners/RunVariables.java:359)。所有结构化可用路径显式列在 EXECUTION_CONTEXT_PATHS(:49),还有测试专门防止它和实际 build 漂移。

谁来渲染? VariableRenderer 包着 Pebble 模板引擎(core/.../runners/VariableRenderer.java:22)。DefaultRunContext.render(inline) 只是转发给它(DefaultRunContext.java:259)。有个小优化:如果字符串里根本没有 {,直接短路返回,不进 Pebble(VariableRenderer.java:77)。默认还支持递归渲染(渲染结果里若还有 {{ }} 继续渲染),上限 MAX_RENDERING_AMOUNT = 100 防死循环。

5.2 存储与工作目录

每个任务跑在一个独立的临时工作目录里,跑完清掉。这由 LocalWorkingDir 实现(core/.../runners/LocalWorkingDir.java:32)。

沙箱的边界就在 resolve():它保证任何要访问的相对路径解析后仍在工作目录内,否则抛异常——防止插件用 ../../ 逃逸出去(LocalWorkingDir.java:86):

// 示意,非源码:LocalWorkingDir.resolve 的安全检查
Path resolved = baseDir.resolve(path).toAbsolutePath();
if (!resolved.startsWith(baseDir)) {
throw new IllegalArgumentException("必须是工作目录内的相对路径");
}

FilesService 在此之上提供 inputFiles(把 {{ }} 渲染后的内容写进工作目录)和 outputFiles(按 glob 收集产物上传到内部存储)两个静态门面(core/.../runners/FilesService.java:22:74)。内部存储(跨任务传大文件)通过 runContext.storage() 暴露,worker 侧装的是 InternalStorage(RunContextInitializer.java:156)。

5.3 日志

runContext.logger() 返回的不是普通 logger,而是 RunContextLogger 装配的一套 logback appender 链(core/.../runners/RunContextLogger.java:34)。每条日志会被复制到三处:

appender去向
ContextAppender打包成 LogEntry 发到日志队列(UI 上能看的那些行)
FileAppender若任务配了 logToFile,写进可下载的日志文件
ForwardAppender转发到服务器本地日志

5.4 密钥(secret 遮蔽)

这是安全上最要紧的一环:密钥绝不能出现在日志里。机制分两步:

  1. 登记: 用到的密钥通过 RunContextLogger.usedSecret 登记(明文 + Base64 两种形式都存),DefaultRunContext.setLogger 会自动把 SECRET 类型的 input 登记进去(RunContextLogger.java:178DefaultRunContext.java:182)。
  2. 遮蔽: 每条日志经过 appender 时,BaseAppender.replaceSecret 递归扫描消息和参数,把登记过的密钥替换成 ******(RunContextLogger.java:347)。

加解密本身由 Secret 这个包私有小类处理(core/.../runners/Secret.java:12);没配加密 key 时它会告警并返回原文而非崩溃(:22)。runContext.decrypt/encrypt 只是转发给它(DefaultRunContext.java:371)。

任务跑完,runContext.cleanup() 清掉工作目录并重置日志 MDC——此后不应再用这个 logger(DefaultRunContext.java:530)。


6. 结果回传:senders

任务产出的东西不止一种,Kestra 为每种开一条独立的 sender(都是 WorkerLoop 子类,worker/.../senders/GrpcWorkerIOSender.java:36),背后策略不同:

sender内容策略失败重发?
taskResultSenderWorkerTaskResult逐条(PER_ITEM)
logEntrySenderLogEntry批量(BATCH)
metricsSenderMetricEntry批量(BATCH)

差异的道理:任务结果丢了会让 execution 永远卡在 RUNNING,必须重发;日志/指标是高频、尽力而为的,重发反而可能给 worker 造成背压(GrpcWorkerIOSenderFactory.java*Sender 方法)。onSendFailed 只对可重试的传输错误(UNAVAILABLE/DEADLINE_EXCEEDED)且非停机中才把批次重新入队(GrpcWorkerIOSender.java:201)。

结果流回 controller 侧的 GrpcWorkerControllerService.sendWorkerTaskResults,它 emit 到结果队列,并在任务到终态时删掉 StateStore 里的 running 记录(worker-controller/.../services/GrpcWorkerControllerService.java:194)。


7. 一条铁律:Worker 不依赖 repository

Worker 从不直接访问数据库。这是刻意的架构约束——Worker 可能是无状态、动态扩缩、甚至部署在隔离网络里的进程,给它数据库连接既不安全也不现实。

那 Worker 需要读 flow 定义、KV、命名空间文件怎么办?走 gRPC 门面。核心侧定义了 FlowMetaStoreInterface 这类接口,worker 侧用 gRPC 实现替换默认实现:

@Replaces(DefaultFlowMetaStore.class) // 用 gRPC 版顶掉直连 DB 的默认版
public class GrpcWorkerFlowMetaStore implements FlowMetaStoreInterface, ... {
private final WorkerFlowMetaStoreServiceBlockingStub workerFlowMetaStoreStub; // 只有 gRPC stub
}

真实实现在 worker/.../stores/GrpcWorkerFlowMetaStore.java:50,同类还有 GrpcWorkerKVMetadataStateStoreGrpcWorkerNamespaceFileMetadataStateStore。它们全都只持有 gRPC stub,没有任何 repository 依赖——需要数据时向 controller 发一次 gRPC 请求,由 controller 代查。这样 worker 对存储层完全无知,RunContextnamespaceKv() 等能力也顺着这条门面工作(RunContext.java:187)。


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

  • permit + bucket 双闸门,且占坑跨越整个 job 生命周期。 permit 管总量、bucket 管分队列配额;bucket 从派发一直占到 completion 回传才释放,保证"派出去但还没跑完"的 job 也算进容量。占坑失败时先还 permit 再试下一个候选,避免误暂停整个队列订阅(WorkerJobDispatcher.java:852)。

  • 虚拟线程 poll + 平台线程执行的中断隔离。 消费者廉价(虚拟线程),任务可被单独中断(平台线程),任务异常被 future.get() 完全兜住,消费者永不因任务出错而死(WorkerJobExecutor.java:299)。

  • 先持久化再派发。 dispatcher 在 sendResponse 之前必落 StateStore,保证 job 任何时刻要么在队列、要么在 StateStore,controller 崩了也不丢(WorkerJobDispatcher.java:888)。

  • catch Throwable 兜底插件。 插件可能抛 Error,worker 宁可把任务标记为 FAILED 也绝不让执行线程裸崩(AbstractWorkerCallable.java:74)。

  • secret 登记 + 递归遮蔽。 密钥两种编码都登记,日志经 appender 时递归替换成 ******,连嵌套 map/list 里的都不放过(RunContextLogger.java:347)。


9. 边界与局限

  • Worker 不推进状态机。 它只产出 WorkerTaskResult,状态怎么流转是 Executor 的事(第 2 章)。

  • 只跑 runnable task。RunnableTask/WorkingDirectory 的任务 worker 直接报错;真正的"流程控制"(分支、循环)在 Executor 侧展开。

  • 结果丢失的兜底不是无限的。 taskResult 会重发,但过了 liveness 超时 worker 会自断连接,由 executor 重新派发它的任务;重发只在运行中的循环内有界进行,停机时不重发以免 cleanup() 排空变成死循环(GrpcWorkerIOSender.java:201)。

  • 超大 payload 会被直接判失败。 如果一个 job 序列化后超过 worker 声明的 gRPC 入站上限,dispatcher 不发它、而是干净地把任务判 FAILED,避免 worker 反复 RESOURCE_EXHAUSTED、执行永久卡住(WorkerJobDispatcher.java:888rejectOversizedJob)。


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

主题文件关键符号
循环骨架worker/.../WorkerLoop.javaWorkerLooprunLoopdoOnLoopsignalJobStop
默认 workerworker/.../WorkerAgent.javaWorkerAgentresolveWorkerGroupId
生命周期装配worker/.../AbstractWorker.javaAbstractWorkerstart
线程池 + 消费者worker/.../WorkerJobExecutor.javaWorkerJobExecutorWorkerJobConsumer.doOnLoop
处理器接口worker/.../processors/WorkerJobProcessor.javaWorkerJobProcessorprocessstopkill
处理器基类worker/.../processors/AbstractWorkerJobProcessor.javaAbstractWorkerJobProcessorcallJob
任务处理器worker/.../processors/WorkerTaskProcessor.javaWorkerTaskProcessordoProcessrunTaskrunAttempt
最内层执行worker/.../processors/internals/WorkerTaskCallable.javaWorkerTaskCallabledoCall
callable 基类worker/.../processors/internals/AbstractWorkerCallable.javaAbstractWorkerCallablecallkill
拉 job 的流worker/.../fetchers/WorkerJobFetcher.javaWorkerJobFetchercalculatePermitsonJobCompleted
结果回传worker/.../senders/GrpcWorkerIOSender.javaGrpcWorkerIOSendersendonSendFailed
分发协调worker-controller/.../services/WorkerJobDispatcher.javaWorkerJobDispatcherhandleIncomingJobfindAndReserveWorkerdispatchJobToWorker
gRPC 服务端worker-controller/.../services/GrpcWorkerControllerService.javaGrpcWorkerControllerServicestreamWorkerJobssendWorkerTaskResults
容量策略worker-controller/.../services/WorkerCapacityPolicy.javaSinglePoolCapacityPolicy.javaWorkerCapacityPolicytryReservehasCapacity
流上下文worker-controller/.../services/WorkerStreamContext.javaWorkerStreamContexttryReserveBucketcompleteJob
沙箱抽象core/.../runners/RunContext.javaRunContextrenderstorageworkingDirlogger
沙箱实现core/.../runners/DefaultRunContext.javaDefaultRunContextrendercleanup
沙箱构造core/.../runners/RunContextInitializer.javaRunContextInitializerforWorker
变量组装core/.../runners/RunVariables.javaRunVariablesDefaultBuilder.buildEXECUTION_CONTEXT_PATHS
Pebble 渲染core/.../runners/VariableRenderer.javaVariableRendererrenderrenderOnce
工作目录core/.../runners/LocalWorkingDir.javaLocalWorkingDirresolvecreateTempFile
文件门面core/.../runners/FilesService.javaFilesServiceinputFilesoutputFiles
日志/遮蔽core/.../runners/RunContextLogger.javaRunContextLoggerusedSecretBaseAppender.replaceSecret
密钥加解密core/.../runners/Secret.javaSecretdecryptencrypt
插件入口core/.../models/tasks/RunnableTask.javaRunnableTaskrun
MetaStore 门面worker/.../stores/GrpcWorkerFlowMetaStore.javaGrpcWorkerFlowMetaStore(@Replaces DefaultFlowMetaStore)