FastGPT 代码沙盒 queueId 排队机制:按业务维度控制代码执行并发的设计与实现

FastGPT 代码沙盒 queueId 排队机制:按业务维度控制代码执行并发的设计与实现 FastGPT 代码沙盒 queueId 排队机制按业务维度控制代码执行并发的设计与实现【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT本文围绕 FastGPT 代码沙盒code-sandbox为POST /sandbox/js与POST /sandbox/python接口引入的queueId排队能力展开先分析多业务方共享同一语言 worker 池时互相阻塞的根因再结合仓库中的QueueIdLimiter实现、环境变量SANDBOX_QUEUE_ID_CONCURRENCY、API 边界校验和测试用例完整讲解这套“按业务 id 限流”的两层队列机制如何落地。读完本文你可以理解排队控制为何放在 HTTP API 边界而非进程池内部并能自行启用和验证该能力。背景单语言共享队列的并发问题FastGPT 的代码沙盒通过ProcessPoolJS和PythonIsolatedRunnerPython同样基于进程池基类为每种语言维护一个 worker 池。请求进入/sandbox/js或/sandbox/python后会直接调用对应进程池的execute()其行为是有空闲 worker 时立即执行没有空闲 worker 时进入进程池内部的waitQueue等待 worker 释放所有请求共享同一个语言池队列无法按业务维度限制某一类请求的并发。进程池的这条等待队列在 base-process-pool.ts 中定义为一个存放{ resolve, reject }回调的数组acquire()/release()L287-L308负责“有空闲 worker 则分配、无则入队、任务结束则唤醒队头或归还空闲池”的完整生命周期管理。execute()L314-L336拿到 worker 后写入任务、等待结果、最终release期间还叠加了超时强制回收与 RSS 内存监控。问题在于这条队列只对“worker 数量”负责对“谁来排队”一无所知。当某个业务方在短时间内提交大量代码执行请求时会占满同语言 worker 并填满池内等待队列其他业务方的请求只能被动排在同一队列后面延迟被最重的业务方拖高。这就是引入queueId的直接动因。需求为运行接口增加 queueId 排队能力需求定义见 问题分析文档非常收敛为POST /sandbox/js和POST /sandbox/python两个运行接口增加可选queueId字段新增环境变量控制同一个queueId同时允许多少个请求进入执行流程环境变量为空时认为不启用排队能力保持现有行为即默认完全兼容既有调用方。配套的 设计文档 给出了目标行为清单默认不启用且完全兼容现有调用启用后仅对带queueId的请求生效同一queueId超出并发上限时按 FIFO 等待进程池原有的 worker 等待队列继续负责真实的 worker 分配。现状分析可复用能力与插入位置已有的两块积木原始分析文档指出实现可以站在两个已有组件之上进程池的 worker 维度等待队列。base-process-pool.ts 中的waitQueue负责“待运行/待分配 worker”的排队语义是“资源不足时等待资源”。通用 FIFO 信号量。semaphore.ts 提供了一个简单的信号量acquire()在current max时直接放行并计数否则把 resolve 回调压入queuerelease()时若队头有等待者则直接唤醒许可“转移”给下一个等待者否则计数减一。其注释明确说明了用途“超出并发上限的请求排队等待避免子进程数爆炸”。QueueIdLimiter的实现与它的acquire/release语义一致但把单一计数扩展为按queueId分桶的Map。插入位置API 边界与进程池之间排队控制被刻意放在 HTTP API 边界和进程池之间形成如下调用链HTTP request - queueId limiter - process pool waitQueue - worker execution这样做的收益是职责隔离进程池只关心 worker 生命周期预热、健康检查、超时回收、one-shot 销毁重建不把业务queueId的概念扩散到 worker 管理层。从源码看limiter 与池之间只有一行组合关系池的execute()对queueId完全无感知。边界语义五个关键决策问题文档对边界行为做了逐条约定这些约定直接决定了实现里每一处分支场景语义环境变量为空不创建 queueId 队列所有请求走现有进程池逻辑环境变量有值但请求未传queueId不按 queueId 排队避免把所有未标识请求挤到同一个匿名队列同一个queueId内按 FIFO 唤醒不同queueId之间不做额外公平调度实际执行顺序仍由进程池 worker 队列决定队列空闲时queueId队列在无运行请求且无等待请求后清理避免高基数 id 导致内存长期增长最后一条尤其重要多租户场景下queueId的基数可能很高例如按团队 id 传递若 Map 条目永不释放服务长期运行会持续积累状态。实现中通过“运行计数归零即删除 Map entry”保证内存占用与“当前活跃的 queueId 数”成正比。实现详解1. 环境变量SANDBOX_QUEUE_ID_CONCURRENCY环境变量在 env.ts 中声明// 进程池 /** 进程池大小预热 worker 数量 */ SANDBOX_POOL_SIZE: IntSchema.min(1).max(100).default(5), /** 同一 queueId 同时可进入执行流程的请求数为空时不启用 queueId 排队 */ SANDBOX_QUEUE_ID_CONCURRENCY: IntSchema.min(1).max(100).optional(),要点取值范围为1100 的正整数optional()表示可以不设置createEnv配置了emptyStringAsUndefined: trueL33因此把该变量设置为空字符串等价于未设置limiter 不会启用与SANDBOX_POOL_SIZE默认 5共同决定系统总吞吐limiter 只限制“每个 queueId 进入执行流程的并发数”真正同时跑的请求数仍受池大小约束。2. 请求体校验queueId 字段约束在 index.ts 中两个接口共用同一个 zod schemaconst queueIdSchema z.preprocess((value) { if (typeof value ! string) return value; const queueId value.trim(); return queueId || undefined; }, z.string().max(128).optional()); const executeSchema z.object({ code: z.string().min(1).max(5 * 1024 * 1024), // 最大 5MB 代码 variables: z.record(z.string(), z.any()).default({}), queueId: queueIdSchema });字段约束为queueId可选字符串会被trim空字符串或纯空白按未传处理映射为undefined非字符串如数字、对象校验失败返回 400最大长度 128避免异常请求造成队列 key 膨胀。合法请求体示例{ code: async function main() { return {} }, variables: {}, queueId: team-xxx }3. QueueIdLimiter按 queueId 分桶的 FIFO 限流器核心实现位于 queue-id-limiter.ts类注释直接点明了设计意图limiter 只负责同一个 queueId 内的 FIFO 限流真实 worker 分配仍交给 ProcessPool 的等待队列处理避免把业务排队规则耦合到 worker 生命周期管理。其结构是一个Mapstring, QueueState每个QueueState含running当前运行数和waitersFIFO 等待回调数组private acquire(queueId: string): Promisevoid { const state this.getOrCreateQueue(queueId); if (state.running this.maxConcurrency!) { state.running; return Promise.resolve(); } return new Promisevoid((resolve) { state.waiters.push(resolve); }); } private release(queueId: string): void { const state this.queues.get(queueId); if (!state) return; const next state.waiters.shift(); if (next) { // 直接把当前运行名额转交给等待队列头部running 保持不变。 next(); return; } state.running Math.max(0, state.running - 1); if (state.running 0) { this.queues.delete(queueId); } }几个值得注意的实现细节许可直接转交release()唤醒等待队头时不增减running而是把当前运行名额直接移交给下一个任务保证运行数恒不超过上限自动清理运行数归零时删除 Map entry对应边界语义中的内存清理要求入口快速路径run()在“未启用”或“queueId 为空”时直接执行 task零额外开销L41-L52构造校验maxConcurrency必须是正整数否则直接抛错防止配置错误在运行期才暴露可观测性statsgetter 暴露enabled、maxConcurrency、queueCount及每个队列的running/queued明细便于排查排队情况。4. API 接入limiter 包裹池的 executelimiter 与服务实例在 index.ts 中一次性创建const jsPool new ProcessPool(env.SANDBOX_POOL_SIZE); const pythonRunner new PythonIsolatedRunner(); const queueIdLimiter new QueueIdLimiter(env.SANDBOX_QUEUE_ID_CONCURRENCY);启动时若启用了排队会打一条 INFO 日志L104-L108。两个执行接口的处理逻辑对称以 JS 为例L192-L219app.post(/sandbox/js, async (c) { const raw await readLimitedJsonBody(c); const parsed executeSchema.safeParse(raw); if (!parsed.success) { return c.json({ success: false, message: Invalid request: ... }, 400); } const result await queueIdLimiter.run(parsed.data.queueId, () jsPool.execute(parsed.data as ExecuteOptions) ); return c.json(result); });即queueIdLimiter.run(queueId, () pool.execute(options))先通过 limiter 拿到该queueId的执行名额任务完成后无论成功失败finally中释放名额并唤醒同队下一位。注意请求体在进入 zod 解析前还经过readLimitedJsonBody的流式大小限制默认SANDBOX_API_MAX_BODY_MB8排队能力与这一层防滥用校验是叠加关系。Python 接口POST /sandbox/python结构完全一致L222-L249只是内部改为pythonRunner.execute(...)。5. 类型与 SDKqueueId 的可选参数一路透传沙盒侧的执行参数类型 types.ts 将queueId定义为可选字段export type ExecuteOptions { code: string; variables: Recordstring, any; queueId?: string; };主应用侧的 HTTP 客户端 codeSandbox/index.ts 中runCode()同样接受可选queueId并原样放进请求体async runCode({ codeType, code, variables, queueId }: { codeType: string; code: string; variables: Recordstring, any; queueId?: string; }) { const url codeType SandboxCodeTypeEnum.py ? /python : /js; const { data } await this.client.post{ codeReturn: Recordstring, any; log: string; }(url, { code, variables, queueId }); return data; }6. 测试验证单元测试 queue-id-limiter.test.ts 覆盖了设计文档中列出的全部行为未启用时不按 queueId 排队new QueueIdLimiter()无并发参数时3 个同 queueId 任务maxRunning为 3且stats显示enabled: false、queueCount: 0同一 queueId 超并发后按 FIFO 排队并发上限设为 1first 占位时 second、third 依次入队stats显示running: 1, queued: 2first 结束后 second 立即开始顺序断言为first-start → first-end → second-start → second-end → third-start全部结束后queueCount归 0验证了状态清理不同 queueId 互不阻塞team-a 持锁时 team-b 的任务照常完成queueId 为空时不排队limiter.run(undefined, ...)下 3 个任务全并发。此外仓库中还存在 API 层集成测试 api.test.ts可验证接口对queueId的接受与非字符串时的 400 响应。实施约定与业务侧接入方式问题文档最后给出两条实施约定值得强调环境变量名称确定为SANDBOX_QUEUE_ID_CONCURRENCY与源码一致。FastGPT 主应用要对工作流代码节点启用业务排队时必须在调用codeSandbox.runCode()时显式传入 queueId。沙盒侧只提供“接口参数 环境变量”这一最小能力不替业务侧猜测默认 queueId——这与“环境变量有值但请求未传 queueId 时不排队”的边界语义是同一设计哲学的两端宁可让排队能力保持显式也不引入隐式的全局默认队列。小结queueId 排队能力是代码沙盒在“worker 资源层排队”之上叠加的一层“业务身份层排队”第一层queueId 队列QueueIdLimiter限制同一业务 id 同时进入执行流程的请求数FIFO 唤醒、空闲自动清理仅对显式携带queueId的请求生效第二层process pool 队列BaseProcessPool.waitQueue继续负责真实的 worker 分配与生命周期管理对业务概念零感知。两层队列在 HTTP API 边界处通过一行queueIdLimiter.run(queueId, () pool.execute(...))组合。默认未设置SANDBOX_QUEUE_ID_CONCURRENCY时行为与历史版本完全一致启用后多业务方共享沙盒的场景下单个业务方的请求洪峰只会在自己的队列内排队不再挤占其他业务方进入 worker 池的通道。【免费下载链接】FastGPTFastGPT is a knowledge-based platform built on the LLMs, offers a comprehensive suite of out-of-the-box capabilities such as data processing, RAG retrieval, and visual AI workflow orchestration, letting you easily develop and deploy complex question-answering systems without the need for extensive setup or configuration.项目地址: https://gitcode.com/GitHub_Trending/fa/FastGPT创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考