Skip to main content
Task Resource Coordinator 是 BlockX 的控制面进程。它只做一件事:Client 调用 ReserveWorkerSlot(taskId),它从 etcd 心跳重建的 Worker 视图里选一个候选,向该 Worker 调 RequestTaskSlot,然后把 workerAddr + slotId 交还给 Client。之后 Client 直接向 Worker 提交 task,Coordinator 不再参与。它在整体架构中的位置见 架构总览,主链路见 Task 生命周期

职责与边界

负责:
  • watch etcd 中的 Worker 注册与心跳,维护一份本地、可丢弃的 WorkerView
  • 对每次 ReserveWorkerSlot 做校验、就绪判定、候选过滤与排序,然后串行尝试少量候选。
  • 把 Worker 返回的 slotId 原样转交 Client;不解释 slotId 的内部含义。
  • 基于最近失败对 Worker 做本地惩罚和短暂熔断。
不负责:
  • 解析 task payload、生成 taskId、接收 SubmitTaskTaskResult
  • 维护中心化 slot 账本、pending 队列或全局 taskId 去重。
  • 向 etcd 写任何数据;etcd 只承载 Worker 发现与心跳。
核心不变量:
  • Worker 本地 slot 状态机是唯一真相;Coordinator 重启只丢失视图缓存和路由偏置,不影响正确性。
  • CoordinatorCore 是 Sans-IO 的纯决策引擎:不读时钟、不做网络,所有方法显式接收 now(UTC 毫秒)。
  • 首轮全量 Worker 视图建立前,以及 watch 断流超过宽限期后,ReserveWorkerSlot 一律返回 Unavailable
  • 单次请求串行尝试候选,不做 hedged request;结果不确定时只对同一 Worker、同一 taskId 重试,绝不切换 Worker。
  • 同一套代码用不同 app.Profile 起两套进程:cmd/coordinator(block 集群)和 cmd/bundle_coordinator(bundle 集群)。bundle 集群的细节见 Bundle 集群

代码位置

核心类型与接口

core

  • core.Configinternal/coordinator/core/coordinator.go):选址与保护参数,DefaultConfig() 给默认值,Validate() 要求 WatchGracePeriodMs <= HeartbeatTimeoutMs
  • core.CoordinatorCore:持有 map[string]*WorkerViewinitializedwatchDisconnectedAtMs。事件入口:ApplyFullWorkerListApplyWorkerUpdateApplyWorkerRemovalMarkWatchDisconnectedMarkWatchReconnected。决策出口:IsReadySelectCandidate。反馈入口:RecordSlotSuccessRecordSlotRejectionRecordSlotBackpressureRecordSlotTimeout
  • core.WorkerViewinternal/coordinator/core/worker_view.go):内嵌 etcd.WorkerHeartbeat,加三个只在本地存在、不写回 etcd 的字段。
  • WorkerView.UsedTotal() = ReservedSlots + RunningTasksWorkerView.FreeCapacity() = TaskSlots - UsedTotal()
  • core.CandidateResult{WorkerAddr, Err}core.CoordinatorError{Code, Message}SelectCandidate 的输出,Codecommontypes.ErrorCode
  • core.ValidateReserveRequest(taskID)internal/coordinator/core/validate.go):只检查 taskId 非空。

adapters

  • adapters.CoordinatorServergrpc_server.go):实现 coordinatorpb.CoordinatorServiceServer,持有 coreSlotReserver 和一把 sync.RWMutexLock/Unlock 暴露给 watcher 用于写 core。
  • adapters.HandlerConfig{MaxTimeoutRetries}:不确定结果时对同一 Worker 的重试次数,默认 1
  • adapters.SlotReserver 接口:RequestTaskSlot(ctx, workerAddr, req) (*commontypes.TaskSnapshot, *SlotError),测试里用 mock 替换。
  • adapters.WorkerRPCClientworker_rpc.go):SlotReserver 的 gRPC 实现,按 workerAddr 缓存一条 grpc.ClientConngrpcErrToSlotError 决定 SlotError.MayHaveAllocated
  • adapters.SlotError{Code, Message, MayHaveAllocated}MayHaveAllocated=true 表示请求可能已到达 Worker,必须留在同一 Worker 重试。
  • adapters.EtcdWorkerWatcher / EtcdWatcherConfig{Endpoints, KeyPrefixes, DialTimeout, RetryDelay}etcd_watcher.go):Run(ctx) 阻塞运行 list+watch。
  • adapters.BundleCoordinatorServerbundle_grpc_server.go):bundle proto 到 CoordinatorServer 的薄转换层。

app 与 api

  • app.Profile{Deployment, Service, WorkerRegistryPrefixes}internal/coordinator/app/app.go):入口身份。Deployment 只能是 "block""bundle",并决定注册哪个 gRPC service、允许哪种 registry kind。
  • app.Run(p Profile):装配 obs、core、WorkerRPCClientCoordinatorServer、etcd client、watcher、gRPC server,阻塞到收到 SIGINT/SIGTERM。
  • app.CoordinatorRuntimeConfigconfig.go):进程级配置聚合,LoadCoordinatorRuntimeConfigFromEnv(p) 从环境变量覆盖并 Validate()
  • etcd.WorkerHeartbeatapi/etcd/types.go):Worker 写入 etcd 的 JSON 值,字段包括 workerAddrtaskSlotsreservedSlotsrunningTaskscpuPercentmemoryPercentmemoryAvailableBytesmemoryTotalBytesprocessRssByteserrorRatelastHeartbeatMs
  • etcd.ParseRegistryPrefix 与常量 LegacyBlockWorkerRegistryPrefix/blockx/workers/)、LegacyBundleWorkerRegistryPrefix/blockx/bundle-workers/)、ProdDefaultBlockWorkerRegistryPrefixProdDefaultBundleWorkerRegistryPrefixTestDefaultBundleWorkerRegistryPrefixapi/etcd/registry.go)。

数据流 / 执行流程

etcd key 是 <prefix> + workerAddr,value 是 WorkerHeartbeat JSON。Worker 侧用带 lease 的 Put 定期刷新,写入逻辑在 internal/worker/adapters/etcd_publisher.go(lease TTL 默认 10s,心跳间隔见 Worker 配置 HeartbeatIntervalMs)。Coordinator 只读。 CoordinatorServer.ReserveWorkerSlot 的编排逻辑按以下顺序执行(internal/coordinator/adapters/grpc_server.go):
1

校验与就绪判定

core.ValidateReserveRequest 拒绝空 taskIdcore.IsReady(now) 为 false 时直接返回 Unavailable,此时未向任何 Worker 发请求。
2

选候选

在读锁下调用 core.SelectCandidate。硬过滤依次为:已在 excludeAddrs、心跳超时(now - LastHeartbeatMs > HeartbeatTimeoutMs)、CPUPercent/MemoryPercent/ErrorRate 超红线、熔断未到期、FreeCapacity() <= 0。剩下的按 candidateRanksBefore 排序:FreeCapacity 大者优先,再比 LastHeartbeatMs 新者、Penalty 小者、最后按 WorkerAddr 字典序。
3

向 Worker 预占

WorkerRPCClient.RequestTaskSlotWorkerRPCTimeout 调 Worker。grpcErrToSlotError 把 gRPC 状态码分成两类:InvalidArgument/ResourceExhausted/FailedPrecondition/NotFound 是明确拒绝;DeadlineExceeded/Canceled/Aborted/Internal/Unknown 以及超时型 Unavailable 标记 MayHaveAllocated=true。返回成功但 SlotID 为空同样按不确定处理。
4

反馈与收敛

成功调用 RecordSlotSuccess 清零失败计数。明确拒绝时:Worker 返回 ResourceExhausted(映射为 ErrNoSlot)走 RecordSlotBackpressure,不计入熔断;其他拒绝走 RecordSlotRejection。不确定结果先在同一 Worker 重试,仍不确定则 RecordSlotTimeout 并返回 Aborted;重试中拿到明确拒绝则换下一个候选。
CoordinatorCore 不做锁保护;CoordinatorServer.mu 是唯一的同步点,watcher 通过 srv.Lock()/Unlock() 写 core,handler 用读锁选候选、写锁记录反馈。

对外错误码

对外协议是 gRPC。业务错误码到 gRPC 状态码的映射在 coordCodeToGRPC
重试后仍无法确认时返回 gRPC AbortedcoordCodeToGRPCErrSlotUncertain 映射为 Aborted,测试工具 internal/testutil/rpc_types.go 反向解释为 ErrSlotUncertain。调用方应把 Aborted 当作”可能已预占,至少等一个 slot TTL 再重试”,把 Unavailable 当作”可立即重试”。

状态与生命周期

CoordinatorCore 只有两个全局状态位,没有显式状态机:
  • initializedApplyFullWorkerList 第一次调用后为 true,之后不再回退。watcher 只有在所有配置 prefix 都完成初始 list 后才调用它。
  • watchDisconnectedAtMs:任一 prefix 的 watch 断开时由 MarkWatchDisconnected(now) 记录首次断开时间;IsReadySelectCandidatenow - watchDisconnectedAtMs > WatchGracePeriodMs 时返回 Unavailable。所有 prefix 重新 list 成功且都处于连接态后,MarkWatchReconnected 清零。
每个 Worker 的本地派生状态:
  • RecentFailureCountPenaltyRecordSlotRejection/RecordSlotTimeout 各加一,Penalty = float64(RecentFailureCount)
  • CircuitBreakerDeadlineMsRecentFailureCount >= CircuitBreakerThreshold(默认 3)时设为 now + CircuitBreakerCooldownMs(默认 10s),期间该 Worker 被硬过滤。
  • RecordSlotSuccess 把三者一起清零;RecordSlotBackpressure 对三者都不动。
  • ApplyFullWorkerList 会保留幸存 Worker 的这三个字段;ApplyWorkerRemoval 删除后重新出现的 Worker 从零开始。
时序参数之间的关系:
  • slot TTL(SlotTTLMs,默认 5s)是 Coordinator 在 RequestTaskSlot 里传给 Worker 的预占窗口,不是 task 执行超时;Coordinator 不再解释它。
  • 单次 ReserveWorkerSlot 没有独立的总时间预算;上界由 MaxCandidates × (1 + MaxTimeoutRetries) × WorkerRPCTimeout 和调用方 ctx 共同决定,默认最坏约 3 × 2 × 3s
  • 心跳新鲜度用 Coordinator 本地时钟与 Worker 写入的 lastHeartbeatMs 相减,机器间时钟偏差会直接影响过滤结果。

配置

所有配置来自 internal/coordinator/app/config.go,默认值取自 DefaultCoordinatorRuntimeConfigcore.DefaultConfigadapters.DefaultHandlerConfigadapters.DefaultEtcdWatcherConfig 可观测性相关的 BLOCKX_PROMETHEUS_LISTEN_ADDRBLOCKX_PROMETHEUS_PATHCHAINTABLE_LOG_BROKERSCHAINTABLE_LOG_TOPICinternal/obs 统一处理,见 可观测性COORDINATOR_CPU_PROFILE / COORDINATOR_TRACEapp.Run 在整个进程生命周期写 pprof / runtime trace 文件,只用于性能排查。 Profile 默认 registry:
  • cmd/coordinator/blockx/workers/
  • cmd/bundle_coordinator/blockx/bundle-workers//blockx/prod/lanes/default/bundle-workers/
  • 测试集群用 COORDINATOR_REGISTRY_PREFIXES=/blockx/test/lanes/default/bundle-workers/ 整体替换。
多 prefix 时 EtcdWorkerWatcherWorkerAddr 合并:LastHeartbeatMs 更新者胜出,相同则 v2 格式胜过 legacy,再相同按 prefix 字符串比较;一个 Worker 只有在所有 prefix 下都消失后才从视图删除。详见 blockx 仓库 docs/specs/2026-07-28-registry-prefix-migration.md

扩展点

  • 改选址策略:只动 internal/coordinator/core/coordinator.goSelectCandidatecandidateRanksBefore,并在 coordinator_test.goTestSelectCandidate_*。core 没有 IO,可以纯函数式测试。当前是确定性排序;心跳里的绝对内存字段(memoryAvailableBytes)只被解析,不参与过滤和排序。
  • 新增一种 Worker 反馈:在 core 加 RecordSlotXxx,在 adapters/grpc_server.gorecordSlotOutcome 里按 SlotError.Code 分派。RecordSlotBackpressure 就是这样加进来的。
  • 改 Worker 错误分类adapters/worker_rpc.gogrpcErrToSlotError 与 Worker 侧 workerCodeToGRPCinternal/worker/adapters/wire.go)必须保持互逆;改一边就要看另一边(见 Worker)。
  • 新增部署形态:加一个 cmd/<name>/main.go 声明 app.Profile,在 app.registerCoordinatorServicevalidateRegistryPrefixes 里加分支,并在 api/etcd/registry.goParseRegistryPrefix 里允许新的 worker kind。
  • 换 Worker RPC 实现:实现 adapters.SlotReserver 接口即可,CoordinatorServer 不依赖具体 client。当前 WorkerRPCClient 对每个 workerAddr 只建一条 grpc.ClientConn,不使用 internal/sdk/grpcclient 的多连接 balancer:指向单个 Worker 的点对点 client 不启用多连接策略。背景见 blockx 仓库 docs/specs/2026-07-23-grpc-client-balancing.md
  • 新增配置项:在 core.Config / adapters.HandlerConfig / adapters.EtcdWatcherConfig 加字段和 Validate 规则,再在 app/config.goLoadCoordinatorRuntimeConfigFromEnv 加环境变量解析,并补 config_test.go

测试

测试组织:
  • internal/coordinator/core/coordinator_test.goSelectCandidate 每个过滤条件与排序规则、熔断触发/恢复、IsReady 的宽限期语义、ApplyFullWorkerList 保留本地字段。
  • internal/coordinator/adapters/handler_test.go:用 mock SlotReserver 覆盖编排循环(超时后换候选、重试仍超时返回 Aborted、全部拒绝返回 ResourceExhausted、空 SlotID 重试),以及真实 gRPC server 的错误码与压缩行为。
  • internal/coordinator/adapters/worker_rpc_test.gogrpcErrToSlotErrorMayHaveAllocated 分类(连接拒绝 vs 超时/取消/Aborted)。
  • internal/coordinator/adapters/etcd_watcher_test.go:内嵌 etcd 上的 list+watch、多 prefix 合并去重、断连/重连的就绪状态、不可解析 value 跳过。
  • internal/coordinator/app/config_test.goprofile_test.go:环境变量解析与 Profile 校验。
  • e2e/contract/coordinator-etcd/:etcd 晚于 Coordinator 启动时先 Unavailable 后恢复;watch 在宽限期内恢复仍可路由;超过宽限期变 Unavailable
  • e2e/contract/coordinator-worker/:明确传输失败切换下一个 Worker、RequestTaskSlot 超时返回 SlotUncertainAborted)、slot 过期后 SubmitTask 被拒、Worker 重启后旧 slotId 失效、统一 slot 池满时拒绝。
E2E 通过 internal/testutilStartCoordinator / StartBundleCoordinator 启动真实进程,用 COORDINATOR_* 环境变量缩短时序参数。整体测试组织见 测试

相关文档

blockx 仓库中的 spec:
  • docs/specs/task-resource-coordinator.md:Coordinator 详细设计,协议、错误语义、选址规则、一致性边界。
  • docs/specs/architecture.md §4.1:控制面在整体架构中的定位与三条边界。
  • docs/specs/2026-07-28-registry-prefix-migration.md:legacy / v2 registry prefix、多 prefix 合并与发布顺序。
  • docs/specs/2026-07-23-bundle-coordinator-api.md:bundle Coordinator 独立 service 身份的由来。
  • docs/specs/2026-07-23-grpc-client-balancing.md:多连接 balancer 的适用范围。
  • docs/specs/worker.mdRequestTaskSlot 与 slot 状态机的权威定义。
站内页面: