ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

LLM 服务防饥饿:用 TypeScript 实现请求调度与并发控制

LLM 服务防饥饿:用 TypeScript 实现请求调度与并发控制 LLM 服务并发一高最先牺牲的往往是交互体验。模型推理本身没变慢但离线批量任务把并发窗口占满之后在线请求只能排队。前端等了 30 秒再重试重试又挤进队列服务整体雪崩。这个问题叫 starvation中文可以理解为“调度饥饿”。今天要看的这个 TypeScript 项目做的正是把 batch 任务和 interactive 请求放进同一个调度体系防止离线任务饿死在线流量。项目的核心不是模型本身而是模型服务上层的请求调度层。它处理三类事情请求分优先级、并发数控制、批处理任务的排队与超时。也就是说后端在把请求转发给 vLLM、Ollama 或 OpenAI 兼容接口之前先经过一个调度器由调度器决定什么时候放行、放行多少、给哪个任务让路。这篇文章会从问题出发拆解 batch 与 interactive 冲突的根因再给出一套 TypeScript 调度器的最小实现说明如何接入现有 LLM 服务、如何做模拟测试验证饥饿是否被消除最后给一份排查清单和可落地的工程建议。项目本身是 Hacker News 上以 Show HN 形式发布的 TypeScript 调度方案下面所有代码示例是通用实现骨架不是原项目源码的快照。适合的读者正在开发 LLM 网关、模型代理、推理平台调度或者需要同时支撑离线评测和在线服务的后端工程师。如果你只是单机单模型自己玩不涉及并发队列可以只读第二、四节了解概念即可。1. 核心能力速览能力项说明项目类型LLM 请求调度与防饥饿控制TypeScript 实现解决的问题批量任务占满并发窗口交互式请求排队饥饿开发语言TypeScript运行环境Node.js可直接用 tsx/ts-node 运行也可编译为 JavaScript 部署核心能力双优先级队列、并发限制、批处理排队、超时终止、指标采集批量任务支持批量提交整体进入低优先级队列接口能力以库/中间件方式集成request 级拦截也可封装成 HTTP 接口实测状态需要按自己的模型服务、请求时长、并发上限做指标验证关键文件scheduler.ts、simulate.ts、server.ts 等按项目结构调整从能力速览可以看到这个项目并不直接做模型推理也不负责图片生成、语音合成这类业务。它的定位是模型服务上层的“交通警察”交互式请求有急事优先通过批量任务不着急但也不能永远堵死。2. 批处理任务为什么会让交互请求“饿死”先说结论不是模型变慢了而是并发窗口被长任务占住了。假设你的模型服务配置了 8 路并发。离线批处理任务一次提交 100 条文本每条生成 2000 个 token。模型服务按顺序执行这些请求一个请求可能耗时 30 秒甚至更久。此时用户在前端点了一下“帮我总结这段文章”这个交互请求进入同一个等待队列排在 100 条批量任务后面。用户在意的不是“队列里有 101 个请求”而是“为什么我的请求 10 秒还没开始推理”。如果前端设置了 5 秒超时第一轮请求超时后自动重试重试又加到队尾等待时间反而更长。服务端同时要处理不断到来的新请求和重试请求队列越来越长这就是典型的饥饿恶化过程。batch 和 interactive 请求的特征差异非常明显维度batch 任务interactive 任务用户感知稍后看结果即可等待几秒就会觉得卡请求时长长可能几十秒到几分钟短期望秒级返回并发数量大可能成百上千少但要求稳定超时容忍高可以排队等待低超时直接失败失败影响重跑即可用户流失、前端报错所以防饥饿的核心不是“拒绝批处理任务”而是“让批处理任务不能独占并发窗口”。调度器要做的是让交互式请求有更高的出队优先级同时保底给批量任务一定的处理机会否则批量任务会反过来被饿死。3. 适用场景与使用边界什么样的服务适合引入这套调度逻辑第一类是 OpenAI 兼容网关前面接用户请求后面统一转发给模型服务。此时请求来源混合了聊天窗口、Agent 工具调用和后台数据清洗任务网关层天然适合做优先级分流。第二类是内部推理平台比如给算法团队做批量评测、给在线业务提供小模型接口。评测任务经常一次提交几千条必须和在线推理接口隔离调度。第三类是 RAG 应用的后端文档解析、向量化、批量索引是后台任务而检索问答是交互任务。两者共用一个 LLM 服务时就会出现排队竞争。不支持、不适合的场景也要说清楚。如果你的服务只是单实例、单模型、单用户使用请求量很小加调度器反而增加复杂度。如果模型负载已经接近物理极限比如 GPU 显存只能支持 1 到 2 路并发调度器解决不了根本问题这时候需要做的是横向扩机器或者拆分服务。另外要注意调度器只能控制“请求何时进入模型服务”不能控制“模型服务内部如何处理已进入的请求”。vLLM 有自己的 continuous batchingOllama 有自己的并发队列调度器的位置是在这些服务之前。合规与安全边界批量任务里可能包含用户上传的文档、个人信息、企业内部数据。引入调度器时需要确认数据的传递链路、日志记录策略、模型服务所在地域是否符合数据合规要求。如果请求需要做脱敏应该在入队之前完成不要把敏感信息留在队列日志里。4. 调度器整体设计思路4.1 双队列分离不要把交互请求和批量任务塞进同一个 FIFO 队列。调度器内部维护两个队列interactiveQueue 和 batchQueue。入队时按任务类型分发出队时优先从 interactiveQueue 取任务。这里有一个关键点单纯的“永远优先交互任务”会饿死批处理任务。生产环境里批处理任务如果长期得不到执行会造成延迟报表、数据管道卡住。所以调度器需要引入“防饿死机制”当 batchQueue 中的任务等待超过一定时间后提升它的出队优先级。4.2 权重轮询一个简单有效的办法是带权重的队列选择。比如默认每轮最多处理 5 个交互任务然后必须拿走 1 个批量任务。这就保证了批量任务在交互请求洪峰时也能缓慢推进。权重可以通过配置文件调节比如interactiveWeight: 5、batchWeight: 1。4.3 并发上限隔离更彻底的方式是给两类任务设置独立的并发上限。交互式请求最多占 6 路并发批量任务最多占 2 路并发。即使批量任务再多也不会挤占交互请求的并发槽位。缺点是模型服务的总并发可能没有被充分利用需要按实际负载调整。4.4 超时与重试控制队列等待超时和任务执行超时要分开处理。交互式请求在队列中等待超过 3 秒可以直接拒绝并返回“服务繁忙”而不是让它继续排队浪费时间。批量任务可以放宽到 10 分钟甚至更长。执行超时则取决于模型服务本身一般建议批量任务的超时时间大于单条请求的极端耗时。重试必须谨慎。交互式请求超时后重试请求应该直接进入队列队首或者携带更高级别的标识避免重试任务再次排到队尾。批量任务重试则建议采用退避策略不要制造重试风暴。4.5 指标输出调度器要能输出运行指标至少要包括当前运行中的请求数、两个队列的长度、任务平均等待时间、交互任务 p95 等待时间、批量任务最后一次执行时间。这些指标可以直接打到 console 或接入 Prometheus方便观察调度是否生效。5. TypeScript 最小实现下面给出一个可运行的调度器骨架。它包含双队列、权重轮询、并发上限、等待超时和任务自动执行。代码是通用示例项目落地时你需要根据自己的模型服务、依赖库和运行环境调整。// scheduler.ts import { EventEmitter } from node:events; export type TaskPriority interactive | batch; export interface TaskT unknown { id: string; priority: TaskPriority; run: () PromiseT; } export interface LlmSchedulerOptions { maxConcurrency: number; interactiveWeight?: number; batchWeight?: number; interactiveQueueTimeoutMs?: number; batchQueueTimeoutMs?: number; } interface QueuedTask { id: string; priority: TaskPriority; run: () Promisevoid; queuedAt: number; resolve: (value: unknown) void; reject: (reason?: unknown) void; timeoutTimer: NodeJS.Timeout; } export class LlmScheduler extends EventEmitter { private readonly interactiveQueue: QueuedTask[] []; private readonly batchQueue: QueuedTask[] []; private running 0; private readonly options: RequiredLlmSchedulerOptions; constructor(options: LlmSchedulerOptions) { super(); this.options { interactiveWeight: 5, batchWeight: 1, interactiveQueueTimeoutMs: 5000, batchQueueTimeoutMs: 300_000, ...options, }; } submitT(id: string, priority: TaskPriority, run: () PromiseT): PromiseT { return new PromiseT((resolve, reject) { const task: QueuedTask { id, priority, run, queuedAt: Date.now(), resolve: resolve as (value: unknown) void, reject, timeoutTimer: setTimeout(() { this.removeTask(id); reject(new Error(task ${id} timeout in queue)); }, priority interactive ? this.options.interactiveQueueTimeoutMs : this.options.batchQueueTimeoutMs), }; if (priority interactive) { this.interactiveQueue.push(task); } else { this.batchQueue.push(task); } this.emit(queued, { id, priority, interactiveWaiting: this.interactiveQueue.length, batchWaiting: this.batchQueue.length, }); this.pump(); }); } private removeTask(id: string): void { this.removeFromQueue(this.interactiveQueue, id); this.removeFromQueue(this.batchQueue, id); } private removeFromQueue(queue: QueuedTask[], id: string): void { const index queue.findIndex((task) task.id id); if (index ! -1) { queue.splice(index, 1); } } private pickNextTask(): QueuedTask | undefined { // 优先取交互任务同时用权重保证批量任务不被饿死 const interactiveReady this.interactiveQueue.length 0; const batchReady this.batchQueue.length 0; if (!interactiveReady !batchReady) { return undefined; } if (interactiveReady !batchReady) { return this.interactiveQueue.shift(); } if (!interactiveReady batchReady) { return this.batchQueue.shift(); } // 两个队列都有任务按权重轮询 const total this.options.interactiveWeight this.options.batchWeight; const interactiveScore this.interactiveQueue.length * (this.options.batchWeight / total); const batchScore this.batchQueue.length * (this.options.interactiveWeight / total); // 如果批量任务等待过久强制放行一个 const oldestBatch this.batchQueue[0]; const batchWaitingMs Date.now() - oldestBatch.queuedAt; if (batchWaitingMs this.options.batchQueueTimeoutMs * 0.8) { return this.batchQueue.shift(); } if (interactiveScore batchScore) { return this.interactiveQueue.shift(); } return this.batchQueue.shift(); } private pump(): void { if (this.running this.options.maxConcurrency) { return; } const next this.pickNextTask(); if (!next) { return; } this.running 1; clearTimeout(next.timeoutTimer); this.emit(started, { id: next.id, priority: next.priority }); Promise.resolve() .then(() next.run()) .then((result) { next.resolve(result); this.emit(completed, { id: next.id, priority: next.priority }); }) .catch((error) { next.reject(error); this.emit(failed, { id: next.id, priority: next.priority, error }); }) .finally(() { this.running - 1; this.emit(drained, { running: this.running, interactiveWaiting: this.interactiveQueue.length, batchWaiting: this.batchQueue.length, }); this.pump(); }); } getMetrics() { return { running: this.running, interactiveWaiting: this.interactiveQueue.length, batchWaiting: this.batchQueue.length, }; } }这个类做的事情可以拆开看submit是唯一入口调用方传入任务 ID、优先级和实际执行函数。任务进入对应对列后会设置一个队列等待超时计时器防止任务一直在队列里空等。pickNextTask负责决定下一个出队的任务既考虑了队列长度也考虑了等待时间。pump在每次任务结束后自动获取下一个任务保证并发数不超过maxConcurrency。这里有一个容易踩坑的点removeTask里如果没有找到任务说明它已经被超时处理不能重复 reject。上面的实现没有处理“超时和正常执行同时触发”的竞态严格一点需要给每个任务增加一个cancelled状态超时置为 cancelled执行结束后检查状态再做 resolve/reject。下面是一个更严谨的片段interface QueuedTask { // ...其他字段 cancelled?: boolean; } // 超时回调里 if (task.cancelled) return; task.cancelled true; this.removeTask(id); reject(new Error(task ${id} timeout in queue)); // 执行完成回调里 if (task.cancelled) { return; } clearTimeout(task.timeoutTimer); task.resolve(result);别小看这个细节。生产环境中任务超时和执行完成往往是竞态关系不处理会导致resolve和reject都被调用外部 Promise 行为不可预期。6. 接入现有 LLM 服务调度器写完之后要接进现有的 LLM 请求链路。假设你有一个 Node.js HTTP 服务请求进来后转发到 OpenAI 兼容接口。改造前的逻辑可能是直接在 handler 里fetch模型服务。接入调度器后请求先入队出队后再执行fetch。下面是一个中间件示例// server.ts import express from express; import { LlmScheduler, type TaskPriority } from ./scheduler; const app express(); app.use(express.json()); const scheduler new LlmScheduler({ maxConcurrency: 8, interactiveWeight: 5, batchWeight: 1, interactiveQueueTimeoutMs: 5000, batchQueueTimeoutMs: 300_000, }); function callLlm(payload: unknown, headers: Recordstring, string) { return fetch(http://127.0.0.1:8000/v1/chat/completions, { method: POST, headers: { Content-Type: application/json, ...headers, }, body: JSON.stringify(payload), }); } app.post(/chat, async (req, res) { const priority: TaskPriority req.headers[x-priority] batch ? batch : interactive; const taskId ${Date.now()}-${Math.random().toString(16).slice(2)}; try { await scheduler.submit(taskId, priority, async () { const response await callLlm(req.body, { Authorization: req.headers.authorization ?? , }); if (!response.ok) { throw new Error(LLM service error: ${response.status}); } return response.json(); }); res.json({ ok: true }); } catch (error) { const message error instanceof Error ? error.message : unknown error; res.status(429).json({ ok: false, error: message }); } }); app.listen(3000, () { console.log(server listen on 3000); });这个例子展示了一个完整链路HTTP 请求到达服务端根据请求头决定优先级交给调度器入队出队后调用真实模型服务。队列等待超过 5 秒的交互请求会被拒绝并返回 429调用方看到 429 后可以采取降级策略而不是无限等待。还应该加一层把调度结果和上游模型服务的结果区分开。调度器返回的是 Promise.then拿到的结果可能是模型服务的完整响应。如果模型服务本身超时或出错错误会从submit的 Promise reject 出来handler 需要统一捕获并返回合适的 HTTP 状态码。7. 功能测试与效果验证调度器接入之后不能靠“感觉变快了”来判断效果要用模拟脚本量化指标。下面这个模拟脚本故意制造压力一次性提交 60 个 batch 任务和 20 个 interactive 任务batch 任务耗时较长interactive 任务耗时较短然后统计 interactive 任务的平均等待时间和 p95 等待时间。// simulate.ts import { LlmScheduler } from ./scheduler; function randomRun(minMs: number, maxMs: number) { const duration Math.floor(Math.random() * (maxMs - minMs 1)) minMs; return new Promisevoid((resolve) { setTimeout(resolve, duration); }); } async function main() { const scheduler new LlmScheduler({ maxConcurrency: 4, interactiveWeight: 4, batchWeight: 1, interactiveQueueTimeoutMs: 10_000, batchQueueTimeoutMs: 120_000, }); const interactiveWaitTimes: number[] []; for (let i 0; i 60; i 1) { scheduler.submit(batch-${i}, batch, () randomRun(800, 1500)).catch(() {}); } for (let i 0; i 20; i 1) { const start Date.now(); scheduler .submit(interactive-${i}, interactive, () randomRun(200, 500)) .then(() { interactiveWaitTimes.push(Date.now() - start); }) .catch(() {}); } setTimeout(() { const sorted interactiveWaitTimes.slice().sort((a, b) a - b); const p50 sorted[Math.floor(sorted.length * 0.5)] ?? 0; const p95 sorted[Math.floor(sorted.length * 0.95)] ?? 0; console.log(interactive tasks:, sorted.length); console.log(p50 wait time:, p50, ms); console.log(p95 wait time:, p95, ms); process.exit(0); }, 30_000); } main();运行方式npx tsx simulate.ts没有接入调度器时同样的模拟脚本里60 个 batch 任务会优先占满 4 路并发interactive 任务可能排在 batch 后面p95 等待时间会明显拉长。接入调度器后interactive 任务会在每 4 个交互任务中间穿插 1 个 batch 任务等待时间会明显下降。你不需要关心脚本执行后具体的毫秒值因为每台机器、每次随机模拟都不同。重点是观察两个指标interactive 任务能否在约定的队列超时时间内完成。batch 任务是否还在持续执行而不是被完全饿死。判断标准可以写成这样如果 20 个 interactive 任务全部成功且 p95 等待时间小于交互式请求的预期阈值说明防饥饿生效。如果 batch 任务在 120 秒内仍然有成功完成记录说明批量任务没有被饿死。如果 interactive 任务大面积超时说明权重配置太低需要提高interactiveWeight或降低maxConcurrency。真实环境测试思路也一样先构造一个混合负载脚本再在调度器前后各跑一轮对比等待时间和完成率。8. 批量任务与接口设计批量任务和交互式请求在接口设计上应该有差异。这里给出一个内部接口拆分的建议。批量任务的提交接口可以这样设计客户端提交一个 batchId包含多条消息服务端把它们全部作为 batch 优先级入队。客户端通过 batchId 查询整体进度。// batchController.ts interface BatchTask { batchId: string; items: unknown[]; } const batchProgress new Mapstring, { done: number; total: number; failed: number }(); async function submitBatch(batch: BatchTask, scheduler: LlmScheduler) { const total batch.items.length; batchProgress.set(batch.batchId, { done: 0, total, failed: 0 }); const promises batch.items.map((item, index) { return scheduler .submit( ${batch.batchId}-${index}, batch, async () { // 这里调用你的 LLM 服务 return { ok: true }; }, ) .then(() { const progress batchProgress.get(batch.batchId); if (progress) progress.done 1; }) .catch(() { const progress batchProgress.get(batch.batchId); if (progress) progress.failed 1; }); }); await Promise.allSettled(promises); } function getBatchProgress(batchId: string) { return batchProgress.get(batchId) ?? { done: 0, total: 0, failed: 0 }; } export { submitBatch, getBatchProgress };批量任务设计上要注意几个问题。第一队列要不要持久化。如果服务重启队列里的任务会全部丢失。对于重要批量任务至少要在入队前把任务信息写入数据库或 Redis 队列调度器只负责内存中这一轮的执行调度。第二批量任务整体超时。一批 500 条消息不可能全部在几分钟内完成。建议给单条任务设置超时整体进度通过 batchId 查询而不是让调用方长时间挂着一个 HTTP 请求。第三失败重试。批处理任务失败后如果调用方拿到 429 或超时错误不要立刻重启整批任务。更稳妥的做法是记录失败明细等一段时间后重新提交失败的那几条。这些设计不依赖具体框架核心是把批量任务的生命周期拆成“提交、执行、查询、重试”四个阶段。9. 性能观察与资源占用调度器本身是 Node.js 进程内的内存逻辑不做网络转发也不代理模型请求所以它的额外开销不会很大。真正需要观察的不是调度器 CPU 占用而是它保护的模型服务的指标。建议在调度器运行过程中持续输出以下指标指标说明正常范围running当前正在执行的任务数不应长期等于 maxConcurrencyinteractiveWaiting交互队列长度应该保持较小偶发波动batchWaiting批量队列长度可以较大但不应该无限增长interactive 平均等待交互任务从入队到开始执行的时间应小于交互超时阈值batch 平均等待批量任务从入队到开始执行的时间如果持续超过 batch 超时阈值说明批量饿死丢弃率队列超时/服务拒绝的比例交互任务应尽量低观察显存和模型并发。如果你的模型服务跑在 GPU 上调度器并发数设置过高会直接导致显存不足设置过低GPU 利用率会下降。调度器解决的是请求排队问题不解决模型实例内部的 batch 优化问题。实际调度中maxConcurrency应该参考模型服务的最大并发数和每条请求的平均耗时来设置。如果要降低调度器带来的风险最简单的方式是队列长度上限。当 interactiveQueue 和 batchQueue 总长度超过阈值时直接把新请求拒绝掉防止内存被队列占满。if (this.interactiveQueue.length this.batchQueue.length maxQueueSize) { throw new Error(queue is full); }这属于流量准入控制和调度策略配合使用。只做优先调度不做容量控制系统仍然可能因为积压任务过多而内存上涨。10. 常见问题与排查方法问题现象可能原因排查方式解决方案interactive 请求还是经常超时interactiveWeight 太低或 maxConcurrency 过大查看调度器指标中 interactiveWaiting 是否持续增长调高 interactiveWeight或为 interactive 设置独立并发上限batch 任务完全无进展权重配置太偏向 interactive查看 batch 任务最后完成时间和 batchQueue 长度降低 interactiveWeight或增加 batch 保底权重服务内存持续上涨队列没有上限任务无限堆积观察进程 heap 和队列长度设置 maxQueueSize添加准入控制同一个任务被重复执行超时与执行完成竞态检查 scheduler 是否有 cancelled 状态处理按示例补齐 cancelled 变量大量 429 返回队列超时时间设得太短查看交互队列超时时间配置权衡前端等待容忍度动态调整超时调度器启动后请求全部失败Node 进程崩溃或未监听端口查看进程日志和端口监听状态检查是否有未捕获异常补充 domain 或 process 级兜底batch 任务批量失败上游模型服务过载或被限流查看调用模型服务的响应码增加重试退避降低 batch 并发压力重试导致负载翻倍客户端超时后立即重试查看请求日志中的重试频率使用指数退避限制最大重试次数比较高频的问题是“交互请求仍然饥饿”。遇到这种情况不要先改代码先看指标。如果interactiveQueue一直处于空的状态说明请求根本没有到达调度器问题在网关或路由层如果interactiveQueue持续增长但running没有达到上限说明pump逻辑没有等到任务执行结束就返回了检查任务 Promise 是否有被 catch。另一个容易忽略的问题是日志量。每次 submit、completed、failed 都打日志在高并发下会产生大量日志影响 Node 进程性能。生产环境建议把事件日志调整成采样输出或者只输出 error 级别指标交给 Prometheus 这类系统采集。11. 最佳实践与使用建议结合上面所有内容落地这个调度器时建议按下面的顺序推进。第一先做压测基线。在接入调度器之前先让模型服务直接处理混合请求记录交互式请求的 p50、p95、p99 延迟。这个基线决定了你能不能量化“调度器是否有效”。第二调度器配置参数化。把maxConcurrency、interactiveWeight、batchWeight、各类超时时间放在配置文件或环境变量里不要硬编码。上线前用低权重、低并发先跑通链路。第三每个任务都要有唯一 ID。任务 ID 用于日志追踪、队列查询和失败重试。如果没有 ID出现问题时你很难定位批量任务里具体是哪条请求失败。第四给入口设置重试策略。交互式请求超时后客户端重试时最好带上原请求 ID 和优先级标记。调度器如果看到同一个请求 ID 已经来过一次可以把它放到队首而不是再次排到队尾。第五保持模型服务本身的退避能力。即使调度器控制了一批请求模型服务仍然可能因为过载而返回 429。上游限流、下游重试这套组合拳要同时设计好。第六关注请求内容合规。LLM 请求可能携带个人数据、企业敏感信息甚至代码片段在入队之前要确认是否允许发送到对应的模型服务。特别是批量任务一次处理的数据量大超时或者失败重试时日志不能随意落盘。第七集群场景下不要每个实例单独一份队列。如果你有多副本运行调度器最好依赖 Redis 或类似的外部队列保证一致性。内存队列适合单实例多实例会造成同样的请求在多个节点重复排队。12. 总结与下一步这个 TypeScript 调度方案解决的问题很具体LLM 服务里批量任务和交互式请求的共存问题。它的价值不在于代码量多而在于把“防饥饿”这个系统设计概念落到了工程里。你应该先验证的是交互式请求在混合负载下的等待时间用模拟脚本或者真实压测工具跑 10 分钟观察 p95 延迟是否明显下降。最容易踩的坑是权重配置失衡要么交互请求仍然饥饿要么批量任务长时间无法执行所以调度器要输出队列长度和等待时间指标用数据说话不要凭感觉调参。后续可以考虑的扩展方向包括把调度器和 Redis 队列结合支持多实例部署。增加带优先级的配额管理按租户或团队隔离。集成 Prometheus 指标采集和 Grafana 面板。在批量任务执行前加入资格检查例如模型路由、内容审计。支持动态调整权重根据当前模型服务负载自动改变 interactive 和 batch 的比例。部署之前先在自己的测试环境里做一轮混合负载验证确认不会影响现有生产流量再逐步灰度上线。这个项目值得关注的点不在于它是不是一个复杂的分布式系统而在于用很小的代码量解决了一个所有 LLM 服务都会遇到的真实问题。
返回列表