task_id、一份”算什么”的配置、一份可选的”结果怎么处理”的配置和一个可选的超时。Worker 把它展开成若干 call,每个 call 是”一个函数 + 一组位置参数”的一次调用。task 只在一个 Worker 上执行,call 在这个 Worker 的 Executor 池里全并行。
task 的输入
对外提交的结构是 gRPCTaskInput(api/grpc/worker/worker.proto):
JSON 形态的一个最小示例(字段名对应 Go 类型
commontypes.TaskInput 的 json tag):
config 携带,例如 BlockTableCallConfig.block 和 BlockTableWriteHandler.block。
call
call 从哪里来由 Call Builder 决定:
InputsCallConfig 的 callList 每行一个 call;扫描类 Builder 对触发表的每一行生成一个 call,行本身按 params 模板放进 args。
call 的性质:
- 互不依赖、全并行,共享同一时刻的只读世界状态。返回顺序即完成顺序,与 call 列表顺序无关。
- 同一 task 内函数与参数都相同的 call 只执行一次(task 级 result cache 与 singleflight)。
- 单个 call 允许有界重试(
DefaultMaxAttempts = 3);任一根 call 终态失败即整个 task 以CALL_FAILED失败(fastFail)。 - 单个 call 的执行预算是
CallDeadlineMs = 5000,只计 CPU 与 IO 后端时间,不含排队。 - 用户函数里的子函数调用(
function.call)也是 call,但不进入outputs,只把结果交还父调用;嵌套深度上限MaxSubcallDepth = 8。
outputs
Calls 阶段结束后,所有成功根 call 的返回值汇聚成outputs。每个元素是一段原样保留的 JSON:一个对象表示一行,一个对象数组表示多行,null 表示这次调用没有产出。Writer Plugin 只看到 outputs,看不到 call 与行的对应关系。
累计字节数超过 MaxCollectedOutputBytes 时,task 停止派发并以 OUTPUT_BYTES_EXCEEDED 失败。
对外状态
结果
TaskResult 只表达 task 级结果:
- 不含每个 call 的返回值。只有
ReturnValueHandler会把整个outputs数组放进pluginResults[i].result;写表插件只返回written_rows等统计。 - 失败时用
failureCode和retryable表达。retryable=true表示 Client 可以用同一份 task 重投;框架自身不做 task 级重试。 - 终态结果在 Worker 上保留
ResultRetentionMs(默认 300000)后清除。Client 用WatchTasks流订阅终态,或用GetTaskResult轮询。
failureCode:
不变量
- 以 task 为中心:读缓存、状态共享、超时都锚定在 task 上,task 结束即释放。
- 单 Worker 执行:一个 task 只在一个 Worker 上跑;Worker 不感知 Coordinator,只暴露
RequestTaskSlot与SubmitTask。 - 函数版本固定:task 激活时固定一个代码快照(
taskCodeEpoch),之后的所有 call 和子调用都用这份快照;代码热更新只影响新 task。 - 幂等:同一
task_id处于RUNNING或结果仍在保留窗口内时,重复SubmitTask返回当前状态,不会再执行一次。 - 不做 task 级自动重试,不支持取消:框架只提供超时;是否重投由 Client 根据
retryable决定。 - 准入错误不进
TaskResult:slot 无效、无空闲 slot、参数非法、resultHandler.type未注册等都以 gRPC status 返回;一旦 task 进入RUNNING,之后的失败都收敛为FAILED终态。
小结
- task =
task_id+ Call Builder 配置 + 可选 Writer Plugin 配置 + 可选超时;call = 函数 + 位置参数。 - call 全并行、无顺序、可去重;
outputs是扁平的返回值集合。 - 对外只有
ALLOCATED / RUNNING / SUCCEEDED / FAILED四个状态;失败用failureCode+retryable表达。 - 框架不重试、不取消 task;幂等键是
task_id。
- Task 生命周期:端到端时序、slot 状态机、时间语义速查。
- Call Builder 与 Writer Plugin:
functionCallConfig与resultHandler的配置形状。 - 函数代码:call 执行的那段代码长什么样。
- 提交 task 示例:用 blockx-py 或直接调 gRPC 提交 task。