
AIBrix Console Planner 架构解析基于策略插件、计划循环与工作池的非阻塞批处理作业调度系统【免费下载链接】aibrixCost-efficient and pluggable Infrastructure components for GenAI inference项目地址: https://gitcode.com/GitHub_Trending/ai/aibrix导读本文深入解析 AIBrix Console 中 Planner 模块的完整设计一个从作业入队Enqueue到执行完成Completed全程异步驱动的批处理作业调度系统。Planner 采用非阻塞架构通过「策略插件Planning Policy Plugin 计划循环Planning Loop 工作池Worker Pool」三大组件协同工作将资源供给Provision、批处理提交MDS CreateBatch与状态轮询从同步请求链路中解耦。读完本文你将掌握 Planner 的作业状态机、策略注册与调度算法、异步轮询的并发模型以及如何结合源码对作业生命周期进行调试与扩展。本文对应的权威设计文档为 apps/console/api/planner/doc/design.md所有源码引用均来自 apps/console/api/planner 目录。1. 总体架构一个非阻塞的批作业调度系统Planner 是 Console 后端BFF与底层资源供给/批处理执行系统之间的调度枢纽。它对外暴露一个精简的接口api/interface.goEnqueue、GetJob、ListJobs、Cancel以及生命周期方法Start、Recover、Close。核心设计哲学是入队即返回——Enqueue只在内存中记录作业并立即返回一个pending占位 batch真正的资源供给Provision、等待资源就绪、提交批处理CreateBatch全部交给后台异步完成。整个系统由三个关键组件构成Planning Policy Plugin调度策略插件可插拔的调度决策组件决定哪些等待作业可以被推进。Planning Loop计划循环周期性的状态机执行器由一个定时器默认 60s和即时触发信号共同驱动。Worker Pool工作池固定数量的 goroutine并发执行资源供给状态查询、MDS 批处理状态查询与清理任务。三者之间的数据流转可用下图概括对应的核心结构体定义在 impl/planner.gopendingQueue/runningQueue是两个基于堆的优先队列utils/priority_queue.gojobs是jobID → *queuedJob的内存索引policy是当前生效的调度策略planningLoop则持有workerPool与触发通道。2. 调度策略插件机制Planning Policy Plugin2.1 接口与注册表策略插件采用经典的注册表模式registry pattern并且注册表以「资源供给类型 × 策略类型」为二级键从而支持为不同资源提供方Kubernetes、AWS、LambdaCloud、RunPod 等挂载不同策略。核心接口定义在 impl/planning_policy.gotype PlanningPolicy[T utils.PriorityQueueItem] interface { Type() PlanningPolicyType Plan(ctx context.Context, input PlanningInput[T]) error } type PlanningInput[T utils.PriorityQueueItem] struct { PlannerBackend plannerBackend // Backend RunningQueue utils.PriorityQueue[T] // Active jobs PendingQueue utils.PriorityQueue[T] // Waiting jobs }策略通过RegisterPlanningPolicy()在包的init()中自注册运行时通过LookupPlanningPolicy()按类型查找工厂函数签名如下后注册者覆盖先注册者last writer winstype Factory[T utils.PriorityQueueItem] func(cfg PolicyConfig) (PlanningPolicy[T], error)从 impl/register.go 可以看到当前仓库中kubernetes、aws、lambda_cloud、runpod四种资源供给类型均注册了simple类型的策略。策略对象在NewPlanner时通过newPlanningPolicy依据Provisioner.Type()与PolicyType构建impl/planner.go。2.2 SimplePolicy默认并发限制策略SimplePolicy是默认且当前唯一内置的策略实现它对供给provisioning施加并发上限防止一次性发起过多资源供给请求。其Plan()算法impl/planning_policy.go分三步统计活跃供给作业数扫描PendingQueue状态为queued且已有scheduledResource已调度但尚未开始供给、或状态为planned的作业均计入扫描RunningQueue状态为planned或resource_preparing供给调用进行中或等待供给就绪的作业计入已过期expiresAt now的作业跳过。限额判断若activeCount MaxConcurrentProvisioning直接返回不推进任何作业。推进等待作业按优先队列顺序遍历PendingQueue跳过终态与cancelling作业对每个未调度的作业调用input.PlannerBackend.Schedule(ctx, job.req)生成ResourceProvisionSpec写入job.scheduledResource直到填满剩余名额。MaxConcurrentProvisioning的默认值为1见 impl/planning_policy.go即默认同一时刻最多只有一个作业处于供给阶段。该值可在构造 Planner 时通过PlannerConfig.MaxConcurrentProvision调整。值得注意的是策略统计的活跃供给作业既包括「已排队待供给」pending 队列中已调度未供给也包括「供给中」running 队列中planned/resource_preparing并在代码注释中明确说明这是为了防止过度调度over-scheduling。3. 计划循环Planning Loop周期性状态机执行器3.1 触发机制与生命周期planningLoopimpl/planning_loop.go的驱动方式有两种定时触发time.NewTicker(w.planInterval)默认间隔 60sdefaultPlanningInterval见 impl/planner.go即时触发Trigger()向容量为 1 的trigger通道发送非阻塞信号——若已有待处理触发则直接丢弃避免信号堆积impl/planning_loop.go。Enqueue、Cancel以及故障恢复完成后都会调用triggerPlanning()立即唤醒一个计划周期而不是傻等下一个 tick从而保证作业在秒级内得到响应。循环具备完整的生命周期管理Start(ctx)启动 goroutine 并等待其就绪通过ready通道同步Stop()取消上下文、停止工作池并Wait()等待 goroutine 退出impl/planning_loop.go。3.2 单次计划周期Plan Once的执行流程每次planOnce()impl/planning_loop.go按固定顺序执行四步StepFunctionPurpose1removeTerminalJobs()从队列和 jobs map 中移除终态作业2policy.Plan()运行调度策略为等待作业分配scheduledResource3processPendingQueue()执行供给处理取消请求4processRunningQueue()提交就绪作业派发异步轮询任务完整时序如下removeTerminalJobs()对jobsmap 做快照克隆后遍历将status.IsTerminal()的作业从对应队列与 map 中移除impl/planning_loop.go。队列处理顺序与公平性planOnce中有一个值得一提的设计——通过runningFirst布尔标志交替先处理RunningQueue还是PendingQueue。注释明确指出这是为了在饱和时防止「活跃作业」与「待供给作业」互相饿死impl/planning_loop.go。processPendingQueue()impl/planning_loop.go收集cancelling作业派发handleCleanup任务源状态cancelling→ 目标状态cancelled收集状态为queued且已调度scheduledResource ! nil的作业派发handleProvisioning任务执行资源供给。processRunningQueue()impl/planning_loop.go第一阶段将状态为resource_preparing且readyToSubmit true、且已到达供给资源起始时间窗provisionStartReached的作业派发submitToMDS提交批处理第二阶段遍历runningQueue中所有非终态作业按状态分流派发任务cancelling→ 清理resource_preparing且尚无allocatedResource→handleResourcePreparing轮询供给状态批处理运行中submitting/scheduling/validating/in_progress/finalizing→handleRunning轮询 MDS 状态过期处理expiresAt now且超出expiryFinalizeGracePeriod5 分钟宽限期后才强制 planner 侧判定过期宽限期内仍轮询 MDS优先采用运行时MDS最终化后的终态 batch保留部分输出与计数详见 impl/planner.go 的注释。防饿死细节被接受派发任务accepted的作业会被移动到 running 队列的队尾Remove后以优先级 0 重新Push确保有界队列饱和时不会每个周期都选中同一批前缀作业。每个计划周期结束后会通过 pkg/metrics 上报console.planner.queue.pending.size/console.planner.queue.running.size队列水位与console.planner.duration周期耗时指标。3.3 任务防重入workInFlight每个queuedJob内嵌一个workInFlight atomic.Boolimpl/job.gotrySubmitJobTask通过CompareAndSwap保证同一作业在同一时刻至多只有一个供给/轮询/提交/清理任务在执行任务完成或TrySubmit失败工作池饱和时都会复位该标志impl/planning_loop.go。这一设计使得计划循环即使被高频触发也不会对同一作业产生重叠的异步操作。4. 工作池Worker Pool有界并发的异步轮询引擎4.1 设计与并发模型工作池实现位于 utils/worker_pool.go关键参数如下parallelismworker goroutine 数量默认10defaultWorkerCountNewWorkerPool在未指定时使用DefaultWorkerParallelism 8而 Planner 构造时显式传入WorkerCount默认 10taskCh容量为queueSize默认1000DefaultWorkerQueueSize的有缓冲任务通道wgWaitGroup 用于追踪已提交任务的完成。工作池提供两种提交方式Submit(fn)阻塞式提交会先登记 WaitGroup 再向通道发送支持 context 取消取消时不做执行并回退计数TrySubmit(fn) bool非阻塞提交通道满时立即返回falseutils/worker_pool.go。计划循环正是使用TrySubmit保证planOnce不被饱和的队列阻塞。Stop()采用优雅停机先取消 context等待所有已受理的Submit完成投递然后**排空drain**通道中已排队的任务并用已取消的 context 执行完毕最后等待所有 worker 退出utils/worker_pool.go。4.2 异步轮询任务的分类TaskFunctionPurposeProvision PollinghandleResourcePreparing()查询供给状态标记readyToSubmitBatch PollinghandleRunning()查询 MDS 批处理状态更新作业状态CleanuphandleCleanup()释放资源、取消批处理ExpiryhandleCleanup()处理过期作业以handleResourcePreparing为例impl/handlers.go它通过provisioner.List按provisionID查询供给状态——状态running记录allocatedResource若后端返回了供给时间窗则把expiresAt设为EndTime并将readyToSubmit置为true状态failed/releasing/released/release_failed将作业标记为resource_failed其他状态等待下一轮。而handleRunningimpl/handlers.go则通过BatchClient.GetBatch拉取 MDS 最新 batch 状态若新状态为终态则转交handleCleanup收敛否则updateStatusUnsafe更新作业状态并把最新 batch 合并持久化到 store。轮询失败不立即判死而是留给下一轮继续注释明确 Dont mark failed, let it be polled again in next iteration。5. 作业状态机Job State Machine5.1 状态分类作业状态定义在 api/status.go按生命周期分为三类CategoryStatesPendingqueued,planned,resource_preparing,submittingRunningvalidating,in_progress,finalizing,cancellingTerminalcompleted,failed,expired,cancelled,resource_failed,submit_failedIsTerminal()用于判断终态ToBatchStatus()将作业状态映射为 OpenAI Batch 状态resource_failed/submit_failed→failedcancelled→cancelledexpired→expired。5.2 状态转移图5.3 关键转移路径FromToHandlerConditionQueuedPlannedSimplePolicy.Plan()Under concurrency limitPlannedResourcePreparinghandleProvisioning()Provision startsResourcePreparingSubmittingsubmitToMDS()readyToSubmit trueSubmittingInProgresssubmitToMDS()CreateBatch returnsInProgressCompletedhandleRunning()MDS reports completionAnyCancellingCancel()User cancelsCancellingCancelledhandleCleanup()Cleanup completesAnyExpiredhandleCleanup()expiresAt now()围绕这些转移源码中有几处细节值得注意handleProvisioningimpl/handlers.go以jobID作为供给请求的幂等键IdempotencyKey供给成功后先把作业推入 running 队列再改状态保证迁移过程中作业在任一队列中可被准入记账发现若此时作业已被取消cancelling立即释放刚申请的供给资源并转交清理。submitToMDSimpl/handlers.go通过backend.BuildRuntime与BuildResourceAllocation构造AIBrixExtraBody再调用BatchClient.CreateBatch。若供给时间窗剩余不足 1 分钟则拒绝提交batchParamsForProvisionDeadline成功后将expiresAt更新为 MDS 返回的 batchExpiresAt并清除readyToSubmit标志。时间戳审计updateStatusUnsafeimpl/handlers.go在每次状态转移时记录对应时间戳queuedAt/plannedAt/resourcePreparingAt/submittingAt/resourceFailedAt/submitFailedAt/canceledAt/expiredAt/completedAt这些字段会通过jobStateSnapshot暴露在GetJob返回的JobState中构成完整的生命周期审计信息。6. 取消、恢复与持久化6.1 取消CancelCancel(jobID)impl/planner.go是异步取消立即将作业状态置为cancelling并返回真正的资源释放与 MDS 批处理取消由计划循环在下一周期派发handleCleanup完成。handleCleanupimpl/handlers.go的行为取决于作业所处阶段尚无batchID仅释放供给资源自然终态expired/completed调用GetBatch读取运行时最终化后的 batch保留 MDS 的输出文件与计数其余情况planner 主动取消调用CancelBatch通知 MDS 取消并采纳返回的 batch 状态——因此最终状态以 MDS 返回为准resolvedStatus JobStatus(batch.Status)。6.2 崩溃恢复RecoverPlanner 具备从持久化存储恢复的能力Recover(ctx)impl/planner.go从 store 读取所有非终态作业并重建内存状态其中有三个精妙的语义处理清理遗留expiresAt完成窗口从 MDS 接受 batch 时才开始计时因此恢复queued/planned作业时清除旧的expiresAtplanned状态回退为queued进程可能在改完内存状态后、Provision返回前崩溃所以恢复的planned作业必须重新调度恢复完成后若 pending 队列非空立即triggerPlanning()触发一个即时计划周期。持久化通过store.UpsertJob完成jobToModel/modelToJobimpl/job.go负责内存对象与存储行的双向映射MDS 拥有的 batch 字段输出文件、计数、Usage、错误信息通过mergeBatchIntoModel合并入库。终态作业从内存逐出后GetJob/ListJobs仍能从 store 与 MDS 补水hydrate出完整视图impl/planner.go。6.3 对外暴露的占位 BatchEnqueue返回的 placeholder batchimpl/planner.go携带endpoint、input_file_id、completion_window、model等字段状态为queued映射的 batch 状态。用户在前端即可立即拿到一个合法的 batch 视图后续由轮询持续刷新——这正是「入队即返回」异步体验的体现。7. 扩展面Per-Provisioner Backend除了策略插件Planner 还通过plannerBackend接口impl/backend.go提供按资源供给类型定制的扩展点每个作业的处理路径为ValidateRequest → Schedule → BuildRuntime → BuildResourceAllocationSchedule产出ResourceProvisionSpec是未来容量感知调度副本数、GPU 类型选择的挂载点默认实现会从ModelTemplateRef.Spec中解析加速器类型与每副本 GPU 数decodeAcceleratorFromTemplate并结合ResourceRequest.Replicas默认 1构建资源组计划BuildRuntime把就绪的ProvisionResult与模型模板的 serving 配置引擎镜像、serve 参数投影为 MDS 的RuntimeRefBuildResourceAllocation将供给结果投影到批处理提交体中的资源分配信息AllocationTimeWindow提取提供方实际分配的时间窗无时间窗时返回 nilplanner 不做起始门控。目前 Kubernetes、AWS、LambdaCloud、RunPod 四类供给方均使用defaultPlannerBackendSimplePolicy的组合impl/register.go说明该架构已为多供应商资源供给预留了清晰的插件化扩展路径。8. 关键配置参数汇总Planner 的全部可调参数集中在PlannerConfigimpl/planner.go默认值由DefaultPlannerConfig()提供参数默认值含义WorkerCount10工作池并发 goroutine 数defaultWorkerCountWorkerQueueSize1000工作池任务通道缓冲大小DefaultWorkerQueueSizePlanningInterval60s计划循环定时触发间隔defaultPlanningIntervalPolicyTypesimple调度策略类型PlanningPolicyTypeSimpleMaxConcurrentProvision1最大并发供给作业数expiryFinalizeGracePeriod5min完成窗口过期后的宽限期之后才强制 planner 侧判定过期配套的单元测试impl/planner_test.go、impl/backend_test.go、utils/worker_pool_test.go、utils/priority_queue_test.go覆盖了状态转移、供给失败、取消竞态、工作池并发与队列优先级排序等关键路径是理解各组件行为的第二手权威资料。9. 总结AIBrix Console Planner 通过「策略插件做决策、计划循环做驱动、工作池做执行」的三段式非阻塞架构把资源供给、批处理提交与状态轮询全部异步化实现了入队即返回、后台收敛的批处理作业调度体验。其核心工程价值在于解耦与可扩展策略插件与 per-provisioner backend 双注册表支持多供应商资源供给与未来容量感知调度并发安全与公平性workInFlight防重入、有界工作池TrySubmit非阻塞提交、队列交替处理与队尾轮转防止饥饿健壮性完整的状态机 时间戳审计、优雅停机排空、崩溃恢复Recover与 MDS 事实源source of truth优先的终态收敛策略。对于希望深入源码的读者建议从 impl/planning_loop.go 的planOnce入手沿handleProvisioning → handleResourcePreparing → submitToMDS → handleRunning → handleCleanup的主链路逐函数阅读即可完整还原设计文档描述的每一个状态转移。【免费下载链接】aibrixCost-efficient and pluggable Infrastructure components for GenAI inference项目地址: https://gitcode.com/GitHub_Trending/ai/aibrix创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考