Skip to main content
task 是 BlockX 的最小调度和执行单元:Client 生成一个 taskId,向 Worker 提交一份 TaskInput,Worker 在本地跑完 Builder → Calls → Plugin 三个阶段后收敛出一份 TaskResult。本页只讲这条主线,各组件细节见 WorkerCoordinatorCall 执行子系统Plugin 系统;整体架构见 架构总览 主链路依赖以下约束(另见 blockx 仓库 docs/specs/architecture.md §2.1):
  • 以 task 为中心:状态共享、读缓存、超时都锚定到 task。
  • 单 Worker 执行:一个 task 只在一个 Worker 上跑,Worker 不感知 Coordinator,只暴露 RequestTaskSlot / SubmitTask
  • task 内 call 全并行:call 之间没有顺序依赖,由 dispatcher 按窗口逐步放量。
  • 同一 task 函数版本固定:激活时 pin 住一个 taskCodeEpochTaskContext.TaskCodeEpoch),期间的代码更新只影响后续 task。
  • 函数只读:用户函数不直接写库,所有写入走 Writer Plugin。
  • 不做 task 级自动重试:Worker 只对单个 call 做有界 attempt 重试(CallRetryPolicy),是否重投 task 由 Client 看 TaskResult.executeResult.retryable 决定。
  • 不支持取消:只有 task 超时(taskTimeoutMs)和可选的 Watch 断连回收。

Task 模型

对外提交的结构是 gRPC TaskInputapi/grpc/worker/worker.proto),Go 侧对应 commontypes.TaskInputinternal/common/types/types.go): 当前注册的 Builder typeinternal/plugin/callbuilder/types/types.go):BlockTableCallConfig(dbScan)、CallListCallConfigBlockBundleCallConfig(bundleScan)。当前注册的 Writer Plugin typeinternal/plugin/event/types/types.go):BlockWriteResultHandlerReturnValueResultHandlerBlockBundleWriteResultHandlerTableUpsertsResultHandler 一个最小的 JSON 形态示例(按 commontypes.TaskInputcall_list.CallListCallConfigBuilderDecl 的字段整理,不是仓库内 fixture):
TaskInput 顶层不带区块上下文。区块信息由具体 Builder 的 config 携带(例如 dbscan.DBScanBuilderDecl.Block),再由 Builder 写进每个 call。 Worker 收到请求后,internal/worker/adapters/wire.gosubmitTaskFromPB 把 proto 转成 commontypes.SubmitTaskRequest(顺带校验 config 是合法 JSON、task_timeout_ms >= 0);internal/worker/adapters/payload.goparseTaskPayload 再包成 commontypes.TaskCtx
TaskCtx 是 Builder、Writer Plugin 和 IO scope 共享的最小上下文;Stage 在阶段推进时被 adapter 改写为 builder / executor / pluginWorkerCore 只把它作为不透明指针挂在 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.ReserveWorkerSlotinternal/coordinator/adapters/grpc_server.go)串行尝试候选 Worker,对每个调 WorkerService.RequestTaskSlot(task_id, ttl_ms)ttl_msCoordinatorCore.SlotTTLMs()。拒绝就换下一个候选;RPC 超时属于”可能已分配”,会对同一 Worker 同一 task_id 有限次重试。
4

Worker 原子分配 slot

WorkerCore.HandleRequestTaskSlotinternal/worker/core/worker.go)在 o.mu 下检查 freeCapacity(),把一个 FREE slot 切到 ALLOCATED,写入 LeaseDeadlineMs = now + ttlMs,返回 slot_idstate=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.SubmitTaskinternal/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;通过后创建 TaskContextPhase = Preparing)并返回命令 PrepareTaskActivation。gRPC 层此时就返回 state=RUNNING
8

adapter 完成激活

Orchestrator.activateAndRuninternal/worker/adapters/orchestrator_phases.go)在 goroutine 里解析当前 taskCodeEpoch、pin 住 Function Code View 快照,然后回送 WorkerCore.HandleTaskActivationPrepared:slot 切到 RUNNINGPhase 推进到 Builder,输出 StartBuilderPhase。随后创建 TaskIOScopeWorkerIOScope.NewTaskScope)和 adapter 侧 taskRuntime。任一步失败走 HandleTaskActivationFailed,failureCode 为 ACTIVATION_FAILEDIO_SCOPE_FAILED
9

Builder 阶段生成 call list

runBuilderPhase 先拿 builderSem,按 FunctionCallConfig.TypeBuilderRegistry 查 Builder,调 CallBuilder.Build(ctx, taskCtx, io) 得到 cbtypes.CallList。结果通过 HandleTaskPhaseFinished(PhaseBuilder, ...) 交回 core,core 输出 StartCallPhase
10

Calls 阶段并行执行

runCallPhaseexecutorSem,调 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 把它转成 PhaseOutcomeHandleTaskPhaseFinished(PhaseCalls, ...)
12

Plugin 阶段落库

core 输出 StartWriterPhase{Outputs}runWriterPhasepluginSem,若 ResultHandler 非空则从 PluginRegistry 查插件并调 WriterPlugin.Execute(ctx, taskCtx, outputs, io),产出 []PluginResult
13

收敛终态

HandleTaskPhaseFinished(PhasePlugin, ...)convergeTaskPhase = TerminaltaskIndexActiveSlotID 切成 ResultRetainedUntil = now + ResultRetentionMs,同一临界区内 releaseSlot 把 slot 回收到 FREE,从 w.tasks 删除,输出 ConvergeTask。adapter 的 executePhaseResult 取消未完成的 call、关闭 TaskIOScope、删除 taskRuntime,最后触发 OnTaskTerminal
14

Client 获取结果

SubmitTask 不返回 TaskResult。Client 用 WatchTasks 流收终态 TaskUpdate,或轮询 GetTaskResultOnTaskTerminal 接到 SubscriptionManager.NotifyTerminalinternal/worker/adapters/stream_subscriber.go)向所有订阅该 taskId 的流推一次终态。
15

心跳反映容量

WorkerCore.DeriveHeartbeat 从 slot 表派生 reservedSlots / runningTasksEtcdPublisher.RunHeartbeatLoop 周期性写 etcd;Coordinator 由此更新视图。

slot 与准入

RequestTaskSlotSubmitTask 是同一套 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 总是输出 StartWriterPhaseResultHandler 为空时 runWriterPhase 仍会取 pluginSem,只是不执行插件、pluginResults 为空。
  • 流式 Builder(StreamingCallBuilder,仅 BlockBundleCallConfigStreamBuildConfig.Enabled 默认关闭)把 Builder 阶段缩成 PrepareStream,扫描放到 Calls 阶段由 scanSemScanSlots = 8)约束;生产中途失败会导致部分 call 已执行。
  • 超时由 TickTimersTimedOutTasks 后调 HandleTaskTimedOut 强制收敛为 TIMED_OUTretryable=true),不强求立即中断 Executor 里正在跑的 call;in-flight call 会收到 CancelCall
task 对外可见状态只有四个(commontypes.TaskState),内部阶段(core.TaskPhasePreparing / Builder / Calls / Plugin / Terminal)都折叠在 RUNNING 里:

结果与查询

TaskResult 只表达 task 级结果,不含每个 call 的返回值。proto 与 Go 类型(commontypes.TaskResult)形状一致:
  • 终态由 state 字段(SUCCEEDED / FAILED)表达,不需要 Client 从 success 推断;TIMED_OUTWATCH_DISCONNECTED 等都是 FAILED 下的 failureCode,不是独立终态。
  • 代码里出现的 failureCodeACTIVATION_FAILEDIO_SCOPE_FAILEDSLOT_LOSTBUILDER_NOT_FOUNDBUILDER_FAILEDCALL_FAILEDOUTPUT_BYTES_EXCEEDEDPLUGIN_NOT_FOUNDPLUGIN_FAILEDCANCELLEDTIMED_OUTWATCH_DISCONNECTED(定义在 internal/worker/core/worker.gointernal/worker/adapters/orchestrator_phases.gointernal/worker/core/dispatcher_internal.go)。
  • ReturnValueResultHandler 把 call 输出数组放进 pluginResults[i].result(gRPC 上是 JSON bytes)。
  • gRPC PluginResult 只有 plugin_name / success / failure_code / result 四个字段;Go 类型里的 FailureMsgRetryable 不上 wire(taskResultToPB)。
两个查询入口都由 Worker 提供,Coordinator 不参与:
  • GetTaskResult(task_id)WorkerCore.HandleGetTaskResult。活动 task 返回 RUNNING;只预占未提交返回 ALLOCATED;终态且未过期返回 state + result;否则 NotFound
  • WatchTasks(task_ids) -> stream TaskUpdateSubscriptionManager.WatchTasks。先对每个 task_id 发一次当前快照(未知的 task_idis_terminal=true, error="not_found",不让整批失败),再对仍在跑的 task 各推一次终态;所有 task 终态后 stream 自行关闭。TaskUpdatetask_id / state / is_terminal / worker_addr / timestamp_ms / result / error。没有单独的取消订阅方法,取消订阅靠关闭 stream。
Watch 还兼作可选的”运行凭证”:task 首次被 Watch 后,最后一条 Watch 断开会触发 HandleWatchDetached,设置 WatchDisconnectDeadlineMs = now + WatchDisconnectGraceMs;到期仍无 Watch 则 HandleWatchDisconnected 收敛为 FAILED / WATCH_DISCONNECTED / retryable=falseWatchDisconnectGraceMs = 0(代码默认)时完全关闭;从未 Watch 过的 task 不受影响。Worker 只在 gRPC 确认 stream 断开后才开始计时,静默断网要先经过 keepalive(默认 30s PING + 10s ACK)。 结果保留:convergeTask 写入 RetainedUntil = now + ResultRetentionMsTickTimersPurgeExpiredResults 到期清除。slot 在收敛那一刻就已回收,与结果保留无关。

时间语义速查

WorkerCore.EffectiveTaskDeadlineMs 是这几个值的收口:task_timeout_ms 为负或超过 MaxTaskTimeoutMs 都直接返回 InvalidArgument,为 0 时取 TaskDeadlineMs

不适用场景

BlockX 的 task 模型面向”单行状态、轻量函数、一次性写入”,以下场景不适合(docs/specs/architecture.md §5):
  • 依赖同一张表大量行的计算:大范围聚合、窗口统计、K 线。
  • 需要任意时间窗口状态或大规模有状态计算。
  • 强依赖流式引擎 Exactly Once 状态恢复。
  • 普通在线表的任意增量流式计算。
非 onchain 表的流式需求,折中做法是 Writer Plugin 同时写在线表和 Kafka,业务侧监听 topic 后再按需提交新 task。

相关文档

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 的声明结构。
站内页面: