taskId,向 Worker 提交一份 TaskInput,Worker 在本地跑完 Builder → Calls → Plugin 三个阶段后收敛出一份 TaskResult。本页只讲这条主线,各组件细节见 Worker、Coordinator、Call 执行子系统、Plugin 系统;整体架构见 架构总览。
主链路依赖以下约束(另见 blockx 仓库 docs/specs/architecture.md §2.1):
- 以 task 为中心:状态共享、读缓存、超时都锚定到 task。
- 单 Worker 执行:一个 task 只在一个 Worker 上跑,Worker 不感知 Coordinator,只暴露
RequestTaskSlot/SubmitTask。 - task 内 call 全并行:call 之间没有顺序依赖,由 dispatcher 按窗口逐步放量。
- 同一 task 函数版本固定:激活时 pin 住一个
taskCodeEpoch(TaskContext.TaskCodeEpoch),期间的代码更新只影响后续 task。 - 函数只读:用户函数不直接写库,所有写入走 Writer Plugin。
- 不做 task 级自动重试:Worker 只对单个 call 做有界 attempt 重试(
CallRetryPolicy),是否重投 task 由 Client 看TaskResult.executeResult.retryable决定。 - 不支持取消:只有 task 超时(
taskTimeoutMs)和可选的 Watch 断连回收。
Task 模型
对外提交的结构是 gRPCTaskInput(api/grpc/worker/worker.proto),Go 侧对应 commontypes.TaskInput(internal/common/types/types.go):
当前注册的 Builder
type(internal/plugin/callbuilder/types/types.go):BlockTableCallConfig(dbScan)、CallListCallConfig、BlockBundleCallConfig(bundleScan)。当前注册的 Writer Plugin type(internal/plugin/event/types/types.go):BlockWriteResultHandler、ReturnValueResultHandler、BlockBundleWriteResultHandler、TableUpsertsResultHandler。
一个最小的 JSON 形态示例(按 commontypes.TaskInput 和 call_list.CallListCallConfigBuilderDecl 的字段整理,不是仓库内 fixture):
TaskInput 顶层不带区块上下文。区块信息由具体 Builder 的 config 携带(例如 dbscan.DBScanBuilderDecl.Block),再由 Builder 写进每个 call。
Worker 收到请求后,internal/worker/adapters/wire.go 的 submitTaskFromPB 把 proto 转成 commontypes.SubmitTaskRequest(顺带校验 config 是合法 JSON、task_timeout_ms >= 0);internal/worker/adapters/payload.go 的 parseTaskPayload 再包成 commontypes.TaskCtx:
TaskCtx 是 Builder、Writer Plugin 和 IO scope 共享的最小上下文;Stage 在阶段推进时被 adapter 改写为 builder / executor / plugin。WorkerCore 只把它作为不透明指针挂在 TaskContext.AdapterTaskCtx 上。
端到端时序
下面按主链路的 16 步展开(另见docs/specs/architecture.md §3)。
1
Client 构造 task
Client 生成
task_id,准备 TaskInput。BlockX 不负责 task 生成、DAG 编排或定时触发。2
Client 向 Coordinator 申请选址(可选)
调用
CoordinatorService.ReserveWorkerSlot(task_id)(api/grpc/coordinator/coordinator.proto)。Coordinator 从 etcd watch 到的心跳里重建 Worker 视图(CoordinatorCore.ApplyWorkerUpdate),用 SelectCandidate 挑候选。不经过 Coordinator 时,Client 可以直接对 Worker 调 RequestTaskSlot,或直接 SubmitTask 让 Worker 自动补申请。3
Coordinator 向 Worker 申请 slot
CoordinatorServer.ReserveWorkerSlot(internal/coordinator/adapters/grpc_server.go)串行尝试候选 Worker,对每个调 WorkerService.RequestTaskSlot(task_id, ttl_ms),ttl_ms 取 CoordinatorCore.SlotTTLMs()。拒绝就换下一个候选;RPC 超时属于”可能已分配”,会对同一 Worker 同一 task_id 有限次重试。4
Worker 原子分配 slot
WorkerCore.HandleRequestTaskSlot(internal/worker/core/worker.go)在 o.mu 下检查 freeCapacity(),把一个 FREE slot 切到 ALLOCATED,写入 LeaseDeadlineMs = now + ttlMs,返回 slot_id 和 state=ALLOCATED。5
Coordinator 回传 worker_addr + slot_id
Coordinator 不保存 task payload、不保存结果、不写 slot 真相到 etcd;
slot_id 对它只是不透明字符串。6
Client 提交 SubmitTask
Client 直连
worker_addr,调 WorkerService.SubmitTask(task, slot_id)。Orchestrator.SubmitTask(internal/worker/adapters/orchestrator.go)先算有效 deadline(WorkerCore.EffectiveTaskDeadlineMs)、解析 payload,再交给 WorkerCore.HandleSubmitTask。7
Worker 校验并激活
HandleSubmitTask 依次做:ready/draining 检查;同 taskId 已在 w.tasks 则返回 RUNNING,已有保留结果则返回终态;没带 slot_id 则内联执行一次 HandleRequestTaskSlot;校验 slot 存在、处于 ALLOCATED、绑定的是同一 taskId;通过后创建 TaskContext(Phase = Preparing)并返回命令 PrepareTaskActivation。gRPC 层此时就返回 state=RUNNING。8
adapter 完成激活
Orchestrator.activateAndRun(internal/worker/adapters/orchestrator_phases.go)在 goroutine 里解析当前 taskCodeEpoch、pin 住 Function Code View 快照,然后回送 WorkerCore.HandleTaskActivationPrepared:slot 切到 RUNNING,Phase 推进到 Builder,输出 StartBuilderPhase。随后创建 TaskIOScope(WorkerIOScope.NewTaskScope)和 adapter 侧 taskRuntime。任一步失败走 HandleTaskActivationFailed,failureCode 为 ACTIVATION_FAILED 或 IO_SCOPE_FAILED。9
Builder 阶段生成 call list
runBuilderPhase 先拿 builderSem,按 FunctionCallConfig.Type 从 BuilderRegistry 查 Builder,调 CallBuilder.Build(ctx, taskCtx, io) 得到 cbtypes.CallList。结果通过 HandleTaskPhaseFinished(PhaseBuilder, ...) 交回 core,core 输出 StartCallPhase。10
Calls 阶段并行执行
runCallPhase 拿 executorSem,调 DispatcherCore.PrepareCallPhaseWithDigests 建立 TaskDispatchState,之后由 dispatch loop 按窗口把 call 通过 UDS ExecuteCall 投给 Python Executor。Executor 内的 SDK IO 请求以 CallWaiting 回到 Worker,由 IO scope 处理后 ResumeCall;子函数调用同样回到 dispatcher 重新绑定。详见 Call 执行子系统 与 IO 访问子系统。11
汇聚 outputs
所有根 call 终态后
DispatcherCore.checkConvergeTask 产出 ConvergeTaskCalls{Success, FailureCode, Outputs};Outputs 只含成功根 call 的返回值,subcall 返回值不进入。adapter 把它转成 PhaseOutcome 调 HandleTaskPhaseFinished(PhaseCalls, ...)。12
Plugin 阶段落库
core 输出
StartWriterPhase{Outputs};runWriterPhase 拿 pluginSem,若 ResultHandler 非空则从 PluginRegistry 查插件并调 WriterPlugin.Execute(ctx, taskCtx, outputs, io),产出 []PluginResult。13
收敛终态
HandleTaskPhaseFinished(PhasePlugin, ...) 调 convergeTask:Phase = Terminal,taskIndex 从 ActiveSlotID 切成 Result,RetainedUntil = now + ResultRetentionMs,同一临界区内 releaseSlot 把 slot 回收到 FREE,从 w.tasks 删除,输出 ConvergeTask。adapter 的 executePhaseResult 取消未完成的 call、关闭 TaskIOScope、删除 taskRuntime,最后触发 OnTaskTerminal。14
Client 获取结果
SubmitTask 不返回 TaskResult。Client 用 WatchTasks 流收终态 TaskUpdate,或轮询 GetTaskResult。OnTaskTerminal 接到 SubscriptionManager.NotifyTerminal(internal/worker/adapters/stream_subscriber.go)向所有订阅该 taskId 的流推一次终态。15
心跳反映容量
WorkerCore.DeriveHeartbeat 从 slot 表派生 reservedSlots / runningTasks,EtcdPublisher.RunHeartbeatLoop 周期性写 etcd;Coordinator 由此更新视图。slot 与准入
RequestTaskSlot 与 SubmitTask 是同一套 Worker 协议的两步;调用方是 Coordinator 还是 Client 对 Worker 没有区别。
RequestTaskSlot(task_id, ttl_ms)只占位,不执行。ttl_ms必须为正(<= 0返回InvalidArgument),只控制ALLOCATED的自动回收,与 task 超时无关。TickTimers每 500ms 扫一次ExpiredSlots,过期 slot 由HandleSlotLeaseExpired回收,同时清掉还停在Preparing的 task 记录。SubmitTask(task, slot_id?):带slot_id时必须命中一条ALLOCATED且绑定同一taskId的记录,否则InvalidSlot;不带时 Worker 用DefaultTTLMs内联申请一次,没容量返回NoSlot。- 幂等:同一
taskId处于ALLOCATED/RUNNING时,RequestTaskSlot返回同一个slot_id;已进入RUNNING或结果仍在保留窗口内时,SubmitTask返回稳定的当前状态,不会再启动一次执行。旧执行结束且结果过期后,同一taskId可以重新分配 slot。 - 没有 release 接口:拿到 slot 后放弃提交只能等 TTL 回收。
docs/specs/architecture.md §4.3.1 末尾):slot 还没从 ALLOCATED 切到 RUNNING 之前的所有失败都以 gRPC status 返回,不进 TaskResult;一旦 HandleSubmitTask 创建了 TaskContext,之后的激活失败、Builder / Calls / Plugin 失败、超时都收敛为 TaskResult 终态。
slot 表、
taskIndex 和心跳派生细节见 Worker;Coordinator 的候选选择、熔断和不确定态重试见 Coordinator。
三个阶段
WorkerCore 只裁决”是否进入下一阶段、何时写终态”;实际执行都在 adapter。阶段推进靠命令 / 事件往返:core 返回 StartBuilderPhase → adapter 执行 → 回送 HandleTaskPhaseFinished(PhaseBuilder, outcome) → core 返回 StartCallPhase,依此类推(internal/worker/core/types.go 定义了全部命令类型)。
其他要点:
- 每个阶段等准入时如果 task ctx 被取消(例如已超时收敛),以
CANCELLED失败;进入 Calls 前还会先检查 deadline,已过期直接按TIMED_OUT收敛,不再占executorSem。 - Calls 阶段成功后 core 总是输出
StartWriterPhase;ResultHandler为空时runWriterPhase仍会取pluginSem,只是不执行插件、pluginResults为空。 - 流式 Builder(
StreamingCallBuilder,仅BlockBundleCallConfig,StreamBuildConfig.Enabled默认关闭)把 Builder 阶段缩成PrepareStream,扫描放到 Calls 阶段由scanSem(ScanSlots = 8)约束;生产中途失败会导致部分 call 已执行。 - 超时由
TickTimers扫TimedOutTasks后调HandleTaskTimedOut强制收敛为TIMED_OUT(retryable=true),不强求立即中断 Executor 里正在跑的 call;in-flight call 会收到CancelCall。
commontypes.TaskState),内部阶段(core.TaskPhase:Preparing / Builder / Calls / Plugin / Terminal)都折叠在 RUNNING 里:
结果与查询
TaskResult 只表达 task 级结果,不含每个 call 的返回值。proto 与 Go 类型(commontypes.TaskResult)形状一致:
- 终态由
state字段(SUCCEEDED/FAILED)表达,不需要 Client 从success推断;TIMED_OUT、WATCH_DISCONNECTED等都是FAILED下的failureCode,不是独立终态。 - 代码里出现的
failureCode:ACTIVATION_FAILED、IO_SCOPE_FAILED、SLOT_LOST、BUILDER_NOT_FOUND、BUILDER_FAILED、CALL_FAILED、OUTPUT_BYTES_EXCEEDED、PLUGIN_NOT_FOUND、PLUGIN_FAILED、CANCELLED、TIMED_OUT、WATCH_DISCONNECTED(定义在internal/worker/core/worker.go、internal/worker/adapters/orchestrator_phases.go、internal/worker/core/dispatcher_internal.go)。 ReturnValueResultHandler把 call 输出数组放进pluginResults[i].result(gRPC 上是 JSON bytes)。- gRPC
PluginResult只有plugin_name / success / failure_code / result四个字段;Go 类型里的FailureMsg、Retryable不上 wire(taskResultToPB)。
GetTaskResult(task_id):WorkerCore.HandleGetTaskResult。活动 task 返回RUNNING;只预占未提交返回ALLOCATED;终态且未过期返回state + result;否则NotFound。WatchTasks(task_ids) -> stream TaskUpdate:SubscriptionManager.WatchTasks。先对每个task_id发一次当前快照(未知的task_id发is_terminal=true, error="not_found",不让整批失败),再对仍在跑的 task 各推一次终态;所有 task 终态后 stream 自行关闭。TaskUpdate含task_id / state / is_terminal / worker_addr / timestamp_ms / result / error。没有单独的取消订阅方法,取消订阅靠关闭 stream。
HandleWatchDetached,设置 WatchDisconnectDeadlineMs = now + WatchDisconnectGraceMs;到期仍无 Watch 则 HandleWatchDisconnected 收敛为 FAILED / WATCH_DISCONNECTED / retryable=false。WatchDisconnectGraceMs = 0(代码默认)时完全关闭;从未 Watch 过的 task 不受影响。Worker 只在 gRPC 确认 stream 断开后才开始计时,静默断网要先经过 keepalive(默认 30s PING + 10s ACK)。
结果保留:convergeTask 写入 RetainedUntil = now + ResultRetentionMs,TickTimers 里 PurgeExpiredResults 到期清除。slot 在收敛那一刻就已回收,与结果保留无关。
时间语义速查
WorkerCore.EffectiveTaskDeadlineMs 是这几个值的收口:task_timeout_ms 为负或超过 MaxTaskTimeoutMs 都直接返回 InvalidArgument,为 0 时取 TaskDeadlineMs。
不适用场景
BlockX 的 task 模型面向”单行状态、轻量函数、一次性写入”,以下场景不适合(docs/specs/architecture.md §5):
- 依赖同一张表大量行的计算:大范围聚合、窗口统计、K 线。
- 需要任意时间窗口状态或大规模有状态计算。
- 强依赖流式引擎 Exactly Once 状态恢复。
- 普通在线表的任意增量流式计算。
相关文档
blockx 仓库中的 spec:docs/specs/architecture.md:§2 任务模型、§3 主链路、§4.3 结果与协议、§5 不适用场景。docs/specs/worker.md:§2 对外协议、§3 本地状态模型、§4 执行流程、§8 失败语义。docs/specs/task-resource-coordinator.md:选址、不确定态重试、熔断。docs/specs/call-execution-subsystem.md:dispatcher、Executor、子函数调用。docs/specs/plugin-system.md:Call Builder 与 Writer Plugin 的声明结构。
- 架构总览:系统架构总览。
- 协议与接口:gRPC / UDS / etcd 协议细节。
- Worker:slot 表、
WorkerCore事件与命令、配置。 - Coordinator:
ReserveWorkerSlot处理链。 - Call 执行子系统:Calls 阶段内部。
- Plugin 系统:Builder 与 Writer Plugin 扩展。
- IO 访问子系统:
TaskIOScope、retry scope、backend admission。