如果你正在维护一个需要对外提供 LLM 服务的系统,无论是 RAG 问答、Agent 对话,还是一个简单的 Chatbot 网关,很可能已经遇到过这种情况:用户白天反馈“机器人变笨了、变慢了”,排查一圈发现,不是模型问题,也不是网络问题,而是有人正在后台批量跑数据任务——凌晨回灌向量库、跑一批评测集、批量生成摘要、给一批文档做 Embedding。
这些批量任务往往一下就把 GPU 或第三方大模型 API 的并发额度占满了,交互式请求只能排在后面。最难受的是,你给批量任务加了队列,交互请求还是慢。因为队列只解决了“谁先进来”,没有解决“正在执行的批量任务会不会把资源抢光”。
这篇文章要聊的,就是这个让很多 LLM 工程团队头疼的问题:如何用 TypeScript 设计一个调度器,让批量 LLM 任务在交互式流量面前主动让路,而不是把用户请求饿死。
读完你至少能带走三样东西:一套关于“饥饿”和“优先级调度”的清晰认知,一个可以直接在 Node.js 服务里跑起来的最小实现,以及一套在生产环境落地的工程建议。
1. 这篇文章真正要解决的问题
1.1 什么是批量 LLM 任务与交互式流量
先定义清楚两个角色。
交互式流量(Interactive Traffic)是指用户在线发起的请求,典型特征是:等待结果的人是一个真实用户。比如用户在对话框里提问,在 RAG 系统里查询资料,在前端页面上触发一次总结。这类请求对延迟极其敏感,通常要求 P95 在几百毫秒到几秒以内。
批量 LLM 任务(Batch LLM Jobs)则是指离线、后台运行的大规模任务,典型特征是:不需要用户在线等待结果。比如:
- 用大模型批量给历史文档生成摘要;
- 对一批样本做质量评测,计算准确率;
- 批量调用 Embedding 接口回灌知识库;
- 用 Agent 批量自动回复工单。
这类任务吞吐优先,延迟可以放宽到分钟级甚至小时级。
1.2 饥饿是怎么发生的
在传统的 Web 服务里,一个请求占用的 CPU 时间通常很短,几百毫秒就结束了。数据库连接、线程池这些资源虽然也共享,但一般不至于让一个用户请求等几分钟。
LLM 场景完全不同。一次大模型推理可能持续几秒到几十秒,而且 GPU 显存是独占式的。一个批量任务占住了 GPU,其他请求即便排在队列里,也只能等它算完。如果一台推理服务器同时允许 8 个请求并发,后台一次性塞进 20 个批量任务,那么前 8 个批量任务会占满所有并发槽位,交互请求排在第 9 位以后。在用户眼里,服务就是“卡死了”。
再叠加一个现实:批量任务往往是大文件、长文本,推理时间比普通聊天请求更长。它们一旦开始执行,交互式请求就可能被压制很久。
这就是“饥饿”(starvation)问题。它不只是“排队顺序不对”,而是“资源被长时间占用的任务垄断了”。
1.3 为什么简单加队列不够
很多人遇到这个问题,第一反应是引入一个优先级队列。核心思路是:交互请求标记为高优先级,批量任务标记为低优先级,调度时先处理高优先级。
这个方案能解决一部分问题,但远远不够。原因有两点:
第一,纯优先级队列只控制“谁先进入执行”,控制不了“正在执行的批量任务占着资源不放”。一个已经跑起来的批量任务不会因为来了一个高优先级请求就主动停下来。
第二,即使你严格控制并发上限,比如总并发只有 4,批量任务也限制在 4 并发以内,看似没问题,实际仍然可能把所有并发槽位都占满。交互请求虽然排在高优位置,却永远没有空位可以执行。
所以,这篇文章真正要解决的问题是:在批量任务与交互式流量共存的系统里,如何通过“优先级队列 + 自适应批量并发控制”的组合方案,确保交互请求既排得靠前,又有资源可以立刻跑起来。
小结论:不要迷信单一优先级队列。要让批量任务“主动感知”交互流量的压力,动态调整自身在途并发数,才是治本思路。
2. 基础概念:饥饿、优先级与协作式退让
2.1 饥饿(Starvation)的本质
饥饿在并发编程里是一个经典问题。当一个低优先级或长耗时的任务持续占用共享资源,导致其他任务永远得不到资源时,就发生了饥饿。
LLM 服务里的饥饿有两个层次:
- 队列饥饿:新来的交互请求排不到队头,因为前面持续有批量任务到达。
- 执行饥饿:交互请求到了队头,但没有足够并发槽位,因为后端的批量任务把资源占满了。
表格对比一下:
| 维度 | 批量任务 | 交互式任务 |
|---|---|---|
| 延迟要求 | 宽松(分钟级) | 严格(秒级) |
| 吞吐偏好 | 高吞吐,可慢跑 | 低延迟,快速响应 |
| 资源占用时长 | 长,可能数十秒 | 短,通常几秒内 |
| 典型场景 | 离线评估、批量摘要 | 在线问答、Agent 对话 |
| 饥饿危害 | 晚一点完成 | 用户直接流失 |
2.2 优先级队列与公平队列
优先级队列很好理解:每个任务带一个 priority,数值小的先执行。实现上通常用二叉堆,插入和取出都是 O(log n)。
但只有优先级队列是不够的。如果所有并发执行位置都被低优先级任务占用,高优先级任务依然要等。
公平队列则更关注“所有人都有资源可用”,典型手段是加权轮询、虚拟时钟、令牌桶。调度器不是简单地“先来先服务”,而是按照每个队列的配额比例分配资源。
本文的设计采用了一个折中:交互队列和批量队列分开,交互队列绝对优先,但批量队列的“最大并发数”动态变化。这既保留了优先级的好处,又通过资源配额避免了执行饥饿。
2.3 协作式退让 vs 抢占式调度
操作系统里的抢占式调度,是通过中断把 CPU 从当前任务手里抢过来。但在 LLM 服务里,你很难安全地“杀掉”一个正在 GPU 上跑的推理任务。强行取消一个批量任务,可能丢失已经算到一半的中间状态,也可能让 GPU 显存无法及时释放。
更务实的办法是协作式退让(Cooperative Backoff):批量任务执行之前和过程中,不断观察交互流量的健康状况,如果发现交互延迟升高、RPS 增大,就主动降低自己的并发数,把资源让出来。
这是本文所有设计背后的核心假设:批量任务是可控的,它们可以在任意时刻停止领取新任务,并把并发数降低到一个安全的水平。
3. 方案设计:一个分层调度器的核心架构
3.1 设计目标
我们要做一个 TypeScript 调度器,它满足以下目标:
- 交互请求永远比批量请求拥有更高优先级。
- 批量请求的并发数不是固定的,而是根据交互流量的实时状态动态变化。
- 总并发有硬上限,防止任何一方把系统打垮。
- 实现不依赖复杂外部组件,跑在 Node.js 进程内即可。
3.2 核心架构
整个调度器分成三层:
- 队列层:两个独立的优先级队列,一个存交互任务,一个存批量任务。
- 决策层:根据最近时间窗口内的交互流量指标,计算批量任务当前允许的最大并发数。
- 执行层:一个简单的并发信号量,控制总的任务在途数量。
流程可以简化为:
- 任务进来,按类型进入交互队列或批量队列。
- drain 循环优先从交互队列取任务。
- 只有交互队列为空时,才尝试从批量队列取任务。
- 从批量队列取任务前,检查“当前批量并发数”是否小于“动态允许的批量并发数”。
- 执行每个任务,并在交互任务结束时,把它的端到端耗时写入滑动窗口。
这里最关键的一点是:批量任务执行期间,系统仍然会有不少在途的批量任务。我们要确保“maxConcurrency - batchMaxConcurrency”始终大于 0,也就是永远给交互请求预留并发槽位。
3.3 为什么用滑动窗口
判断“交互流量是否健康”,本质上是在做实时流量感知。常用的手段有两种:指数移动平均(EMA)和滑动窗口。
滑动窗口的优势是直观:统计最近 10 秒内完成的交互任务数量,以及这期间的 P95 延迟。窗口越短,反应越快,但也越容易抖动;窗口越长,越平滑,但反应迟钝。
本文实现里默认使用 10 秒窗口。生产中可以根据业务容忍度调整,通常建议 5 到 30 秒。
小结论:这个方案的本质是“优先级保证排队顺序,动态并发保证资源配额”。两个机制缺一不可。
4. 环境准备与项目初始化
4.1 运行环境
本示例代码基于 Node.js 18+ 和 TypeScript 5。具体版本以你本机实际安装为准,文章重点演示通用思路。
建议提前安装:
- Node.js 18 或更高版本
- npm 或 pnpm
- TypeScript
- tsx(用于直接运行 TypeScript,不需要先编译)
4.2 初始化项目
创建一个新目录,并初始化 TypeScript 项目:
mkdir llm-traffic-scheduler cd llm-traffic-scheduler npm init -y npm install -D typescript tsx @types/node npx tsc --init修改tsconfig.json,重点配置如下:
{ "compilerOptions": { "target": "ES2022", "module": "CommonJS", "moduleResolution": "Node", "strict": true, "outDir": "dist", "esModuleInterop": true, "skipLibCheck": true }, "include": ["src", "scripts"] }然后在package.json里加上运行脚本:
{ "scripts": { "build": "tsc", "simulate": "tsx scripts/simulate.ts" } }这样一个最小环境就准备好了。第三步代码实现里,我们会写出 4 个文件,尽量不依赖第三方库,便于理解原理。
5. 核心代码实现
5.1 模块一:优先级队列
优先级队列是整个调度器的基础数据结构。这里用一个最小二叉堆实现,节点的 priority 数值越小,优先级越高。
文件路径:src/priority-queue.ts
// src/priority-queue.ts export interface QueueTask<T> { id: string; priority: number; value: T; enqueuedAt: number; } export class PriorityQueue<T> { private heap: Array<QueueTask<T>> = []; push(task: QueueTask<T>): void { this.heap.push(task); this.bubbleUp(this.heap.length - 1); } pop(): QueueTask<T> | undefined { if (this.heap.length === 0) { return undefined; } const top = this.heap[0]; const last = this.heap.pop()!; if (this.heap.length > 0) { this.heap[0] = last; this.sinkDown(0); } return top; } get size(): number { return this.heap.length; } private bubbleUp(index: number): void { while (index > 0) { const parent = Math.floor((index - 1) / 2); if (this.heap[parent].priority <= this.heap[index].priority) { break; } [this.heap[parent], this.heap[index]] = [ this.heap[index], this.heap[parent], ]; index = parent; } } private sinkDown(index: number): void { const n = this.heap.length; while (true) { let smallest = index; const left = 2 * index + 1; const right = 2 * index + 2; if (left < n && this.heap[left].priority < this.heap[smallest].priority) { smallest = left; } if (right < n && this.heap[right].priority < this.heap[smallest].priority) { smallest = right; } if (smallest === index) { break; } [this.heap[index], this.heap[smallest]] = [ this.heap[smallest], this.heap[index], ]; index = smallest; } } }这段代码维护的是小顶堆。push 进去后向上冒泡,pop 时把堆顶元素取走,再把堆尾元素放到堆顶向下调整。这样每次取出的都是当前队列里 priority 最小的任务。
5.2 模块二:滑动窗口指标收集器
调度器需要实时掌握交互流量的状态。这个模块用一个环形数组存最近 N 秒的样本,每次交互请求结束时就写入一条延迟样本,并统计窗口内的 RPS 和 P95 延迟。
文件路径:src/interactive-metrics.ts
// src/interactive-metrics.ts export interface TrafficSnapshot { rps: number; p95LatencyMs: number; sampleCount: number; } export class SlidingWindowMetrics { private readonly windowMs: number; private samples: Array<{ ts: number; latencyMs: number }> = []; constructor(windowMs = 10_000) { this.windowMs = windowMs; } addSample(latencyMs: number): void { const now = Date.now(); this.samples.push({ ts: now, latencyMs }); this.trim(now); } snapshot(): TrafficSnapshot { const now = Date.now(); this.trim(now); if (this.samples.length === 0) { return { rps: 0, p95LatencyMs: 0, sampleCount: 0 }; } const rps = (this.samples.length / this.windowMs) * 1000; const latencies = this.samples .map((s) => s.latencyMs) .sort((a, b) => a - b); const p95Index = Math.min( latencies.length - 1, Math.floor(latencies.length * 0.95) ); return { rps, p95LatencyMs: latencies[p95Index], sampleCount: this.samples.length, }; } private trim(now: number): void { while ( this.samples.length > 0 && now - this.samples[0].ts > this.windowMs ) { this.samples.shift(); } } }在真实系统里,RPS 通常统计的是“请求到达率”,也就是请求刚进来的流量。而这里统计的是“交互任务完成率和P95延迟”。这其实是一个工程取舍:完成率和延迟数据在业务侧更容易采集,而且当系统过载时,P95 延迟会迅速上升,这个信号已经足够驱动调度逻辑。
如果你希望更精确地反映“请求到达压力”,可以把 addSample 改成在提交任务时调用,再单独记录延迟。这不会影响架构。
5.3 模块三:自适应并发调度器
这是核心模块。它同时维护两个优先级队列,drain 循环优先从交互队列取任务,只有交互队列为空时才处理批量任务。
文件路径:src/scheduler.ts
// src/scheduler.ts import { PriorityQueue, QueueTask } from './priority-queue'; import { SlidingWindowMetrics } from './interactive-metrics'; export type TaskType = 'interactive' | 'batch'; export interface SchedulerOptions { maxConcurrency: number; batchMaxConcurrency: number; batchMinConcurrency: number; interactiveRpsThreshold: number; interactiveP95LatencyMs: number; windowMs?: number; } interface ScheduledTask<T = unknown> { id: string; type: TaskType; execute: () => Promise<T>; resolve: (value: T) => void; reject: (reason?: unknown) => void; submittedAt: number; } export class LmTrafficScheduler { private readonly interactiveQueue = new PriorityQueue<ScheduledTask>(); private readonly batchQueue = new PriorityQueue<ScheduledTask>(); private readonly metrics: SlidingWindowMetrics; private readonly options: SchedulerOptions; private running = 0; private batchRunning = 0; constructor(options: SchedulerOptions) { this.options = options; this.metrics = new SlidingWindowMetrics(options.windowMs ?? 10_000); } submit<T>( type: TaskType, id: string, execute: () => Promise<T> ): Promise<T> { return new Promise<T>((resolve, reject) => { const task: ScheduledTask<T> = { id, type, execute, resolve, reject, submittedAt: Date.now(), }; const queue = type === 'interactive' ? this.interactiveQueue : this.batchQueue; queue.push({ id, priority: type === 'interactive' ? 0 : 10, value: task, enqueuedAt: task.submittedAt, } as QueueTask<ScheduledTask>); this.drain(); }); } private get allowedBatchConcurrency(): number { const { rps, p95LatencyMs } = this.metrics.snapshot(); const rpsPressure = Math.min( 1, rps / this.options.interactiveRpsThreshold ); const latencyPressure = Math.min( 1, p95LatencyMs / this.options.interactiveP95LatencyMs ); const pressure = Math.max(rpsPressure, latencyPressure); const range = this.options.batchMaxConcurrency - this.options.batchMinConcurrency; return Math.max( this.options.batchMinConcurrency, Math.round(this.options.batchMaxConcurrency - range * pressure) ); } private drain(): void { while (this.running < this.options.maxConcurrency) { const interactiveTask = this.interactiveQueue.pop(); if (interactiveTask) { this.execute(interactiveTask.value); continue; } const batchLimit = this.allowedBatchConcurrency; if (this.batchRunning >= batchLimit) { break; } const batchTask = this.batchQueue.pop(); if (batchTask) { this.execute(batchTask.value); } else { break; } } } private execute(task: ScheduledTask): void { this.running++; if (task.type === 'batch') { this.batchRunning++; } Promise.resolve() .then(() => task.execute()) .then((value) => task.resolve(value)) .catch((reason) => task.reject(reason)) .finally(() => { this.running--; if (task.type === 'batch') { this.batchRunning--; } else { this.metrics.addSample(Date.now() - task.submittedAt); } this.drain(); }); } }这段代码是整套方案的核心,需要重点讲清几个设计:
第一,交互队列和批量队列是分开的。这避免了单一优先级队列里“队头被低优批量任务挡住,高优交互任务取不到”的问题。交互任务永远优先于批量任务出队。
第二,allowedBatchConcurrency 是动态计算的。它取 RPS 压力和 P95 延迟压力中的最大值作为 pressure,pressure 越大,批量并发越往 batchMinConcurrency 收缩。这实现了“交互流量越忙,批量任务越收敛”。
第三,drain 循环里先处理交互队列,再处理批量队列。如果批量队列已经达到动态上限,就直接 break,等待某个任务完成后的 finally 回调再次触发 drain。
要注意:batchMinConcurrency 不要设成 0。如果设成 0,在交互流量持续很高时,批量任务可能永远得不到执行机会,这会导致另一种饥饿——批量任务饿死。这也很重要。
5.4 模块四:模拟运行脚本
为了验证调度器行为,我们写一个模拟脚本。它可以配置几种不同类型的任务,模拟批量任务持续占用资源和交互请求穿插进入的过程。
文件路径:scripts/simulate.ts
// scripts/simulate.ts import { LmTrafficScheduler } from '../src/scheduler'; const scheduler = new LmTrafficScheduler({ maxConcurrency: 8, batchMaxConcurrency: 6, batchMinConcurrency: 1, interactiveRpsThreshold: 20, interactiveP95LatencyMs: 500, }); const sleep = (ms: number) => new Promise((resolve) => setTimeout(resolve, ms)); function runBatchJob(id: number) { return scheduler.submit('batch', `batch-${id}`, async () => { const start = Date.now(); // 模拟一次耗时较长的批量推理 while (Date.now() - start < 200) { await sleep(5); } }); } async function runInteractiveJob(id: number) { const start = Date.now(); await scheduler.submit('interactive', `interactive-${id}`, async () => { // 模拟交互式请求,耗时短 await sleep(20); }); const cost = Date.now() - start; console.log(`[interactive-${id}] end-to-end ${cost}ms`); } async function main() { // 先塞入 20 个批量任务 for (let i = 0; i < 20; i++) { void runBatchJob(i); } // 500ms 后开始持续发送 10 个交互请求 await sleep(500); for (let i = 1; i <= 10; i++) { await runInteractiveJob(i); await sleep(10); } } main().then(() => process.exit(0));这个脚本里,批量任务耗时 200ms,交互任务耗时 20ms。从逻辑上讲,如果没有调度器,20 个批量任务会迅速占满并发槽位,交互请求的端到端耗时可能达到几百毫秒甚至更久。当交互请求进来了,如果只靠固定并发控制,批量任务仍然占着 6 个并发,交互请求最多只能拿到 2 个并发,依然会明显积压。
但有了动态批量并发控制后,一旦检测到交互 RPS 上升或 P95 延迟上升,批量任务的最大并发就会下降,给交互请求腾出更多执行槽位。
6. 运行结果与效果验证
6.1 运行方式
在项目根目录执行:
npm run simulate6.2 预期输出
每次运行的数字会有波动,但趋势应当一致:后出现的交互请求,其端到端耗时应该远小于没有调度控制时的水平。如果打印 batch 任务开始和结束的关键时间点,可以看到批量任务在交互请求密集时段被“压住”了并发数。
更严谨的验证方式是做一个 A/B 对比:
- 方案 A:不使用调度器,直接用
Promise.all并发执行所有任务; - 方案 B:使用本文的 LmTrafficScheduler;
- 对比指标:交互请求的 P95 端到端耗时、批量任务的整体完成时间。
从行为上判断是否有效的标准有三条:
- 交互请求的 P95 延迟没有持续超过阈值。
- 批量队列在交互流量高峰时增长变慢,但仍能推进,没有完全卡死。
- 系统总并发没有超过 maxConcurrency。
6.3 如何判断失败
如果运行后交互请求依然很慢,优先检查两处:
第一,交互 RPS 阈值设置是否合理。如果阈值设成 1000,而实际交互 RPS 只有 50,那么 RPS 压力永远接近 0,批量并发永远收缩不下来。
第二,batchMaxConcurrency 是否接近 maxConcurrency。如果 maxConcurrency 是 8,batchMaxConcurrency 也是 8,那么批量任务可能把资源占满,交互请求没有预留额度。
生产环境调优时,这三个参数一定要联动调整,不能只改一个。
7. 常见问题与排查思路
| 问题现象 | 可能原因 | 排查方式 | 解决方案 |
|---|---|---|---|
| 交互请求 P95 延迟依然很高 | 批量在途任务占住 GPU 显存,调度器无法抢占已执行任务 | 查看 GPU 利用率、推理服务是否长时间排队 | 限制 batchMaxConcurrency;为批量推理设置超时或可抢占执行 |
| 批量任务完全不动 | batchMinConcurrency 设为 0,且交互流量长期高 | 查看 batch 队列积压长度 | 将 batchMinConcurrency 设为 1 或更高 |
| 批量并发频繁抖动 | 滑动窗口太小,RPS 波动被放大 | 观察 allowedBatchConcurrency 变化曲线 | 拉大 windowMs,或对 pressure 做移动平均 |
| 交互延迟没超,但批量任务堆积过多 | 交互 RPS 阈值过高,压力信号失真 | 记录实际 RPS 与 P95 延迟 | 用压测标定阈值,或改为按延迟为主 |
| 多实例部署后,交互请求仍被挤占 | 每个实例各自维护本地计数,总并发超出预期 | 检查实例数 × 单实例并发 | 引入 Redis 分布式信号量,或使用 BullMQ 等服务端队列 |
| 进程重启后批量任务丢失 | 队列只存在内存中 | 查看任务是否有持久化 | 使用 Redis Streams、BullMQ 或数据库任务表,消费时保证幂等 |
这里要特别提醒一点:本文的调度器是“协作式”的,它只能控制“新任务是否开始执行”,不能中断已经开始执行的推理请求。如果批量任务已经跑起来并且占用了 GPU 显存,调度器层面能做的只是不继续放新的批量任务进去。真要抢占显存,需要在推理服务层配合,比如 vLLM 等支持连续批处理的服务可以动态调整最大批大小,或者干脆对批量长任务设置硬超时。
8. 最佳实践与工程建议
8.1 监控指标不能只看平均
交互延迟要重点看 P95 / P99,而不是平均值。平均值很容易被少数极快请求拉低。调度器的效果验证,建议以 P95 端到端延迟为主,配合队列积压长度和批量并发度一起看。
推荐的核心指标:
- 交互请求 P95 / P99 端到端延迟;
- 交互请求到达 RPS;
- 批量任务当前在途并发数;
- 批量任务队列积压长度;
- 总并发使用率。
这些指标建议全部接入 Prometheus 或云监控,通过 Grafana 面板实时展示。
8.2 阈值一定要用压测标定
interactiveRpsThreshold 和 interactiveP95LatencyMs 这两个阈值是整个调度器的大脑。如果拍脑袋设,调度器给出的压力信号就是失真的。
比较靠谱的做法是:压测环境下,逐渐提高交互 RPS,记录延迟拐点。当 P95 延迟开始超过业务的容忍线时,把此时的 RPS 和 P95 作为阈值。注意不同模型、不同输入长度、不同 GPU 配置,阈值都会不同,换模型后要重新标定。
8.3 用平滑策略避免抖动
滑动窗口天然有一定平滑作用,但压力计算里使用Math.max(rpsPressure, latencyPressure)会让批量并发随着最紧张的指标快速变化。
如果生产中发现批量并发来回跳动,可以引入简单的一次指数平滑:
let smoothedPressure = 0; const alpha = 0.3; function updatePressure(rawPressure: number): number { smoothedPressure = alpha * rawPressure + (1 - alpha) * smoothedPressure; return smoothedPressure; }每次计算 allowedBatchConcurrency 时,先对压力做平滑,再计算并发数。这是最省事也最有效的防抖手段。
8.4 分布式环境需要改造
本文示例是单进程方案,适合单体服务或单实例推理网关。如果你部署了多个实例,每个实例都维护自己的 running 和 batchRunning 计数