AI 任务队列调度实战:并发限流、幂等重试与持久化

作者:袖梨 2026-09-19

一次 AI 内容生成可能同时触发模型推理、语音合成、字幕处理等多项异步调用。只把这些请求简单并发执行,遇到用户重复提交、服务限流或进程中断时,就可能产生重复扣费、状态错乱和任务丢失。要让流程稳定运行,需要从任务状态、资源配额、失败重试和恢复机制几方面建立完整的调度体系。

AI 任务队列与并发调度:限流、退避、幂等与状态持久化

一个 AI 剪辑任务会引爆一连串下游调用:视频理解、文案生成、TTS 合成、字幕对齐——每个外部服务都有自己的配额(RPM/TPM)、自己的故障率、自己的计费。没有调度的裸调用,在用户连点三次"生成"后就会撞限流、重复计费、状态错乱三连击。本文拆解 AI 任务调度的四块基石:并发限流、退避重试、幂等去重、状态持久化。

一、任务状态机:一切调度的前提

先定义任务生命周期——不是布尔值,是显式状态集:

type TaskStatus =
  | 'queued'       // 已入队,等待调度
  | 'running'      // 执行中
  | 'succeeded'    // 成功终态
  | 'failed'       // 失败终态(重试耗尽)
  | 'cancelled'    // 用户取消
  | 'deduped'      // 命中幂等键,直接复用已有结果

const VALID_TRANSITIONS: Record<TaskStatus, TaskStatus[]> = {
  queued:    ['running', 'cancelled', 'deduped'],
  running:   ['succeeded', 'failed', 'cancelled'],
  succeeded: [],
  failed:    ['queued'],   // 手动重试回到队列
  cancelled: [],
  deduped:   [],
}

function transition(task: Task, next: TaskStatus): Task {
  if (!VALID_TRANSITIONS[task.status].includes(next)) {
    throw new IllegalTransitionError(task.status, next)
  }
  return { ...task, status: next, updatedAt: Date.now() }
}

非法转移直接抛错而不是静默忽略——调度器 bug 要在开发期炸出来,不能在用户任务上静默吞掉。failed → queued 是唯一回边,支撑手动重试;自动重试不回到 queued(避免和排队延迟耦合),而是在 running 内部做。

二、并发限流:令牌桶与分层配额

外部服务限流是多维的:每分钟请求数(RPM)、每分钟 token 数(TPM)、并发连接数。简单 semaphore 只能管并发数,管不了速率。令牌桶:

class TokenBucket {
  private tokens: number
  private lastRefill: number

  constructor(
    private capacity: number,      // 桶容量(突发上限)
    private refillRate: number,    // 每秒补充速率
  ) {
    this.tokens = capacity
    this.lastRefill = Date.now()
  }

  async acquire(cost = 1): Promise<void> {
    for (;;) {
      this.refill()
      if (this.tokens >= cost) {
        this.tokens -= cost
        return
      }
      const waitMs = ((cost - this.tokens) / this.refillRate) * 1000
      await sleep(Math.min(waitMs, 1000))
    }
  }

  private refill() {
    const now = Date.now()
    const elapsed = (now - this.lastRefill) / 1000
    this.tokens = Math.min(this.capacity, this.tokens + elapsed * this.refillRate)
    this.lastRefill = now
  }
}

// 分层配额:LLM 按 token 计费,TTS 按字符计费
const llmBucket = new TokenBucket(capacity: 100_000, refillRate: 80_000)  // TPM 80k
const ttsBucket = new TokenBucket(capacity: 500, refillRate: 40)          // RPM 40

acquire(cost) 按真实成本扣减:LLM 请求按预估 token 扣,TTS 按字符数扣。配额感知的调度顺序:队列里多个任务时,优先执行"剩余配额能完成"的任务,避免一个 100k token 的大任务把 30 个小任务全部堵死——这是短作业优先(SJF)思想在配额场景的应用。

三、退避重试:区分"值得重试"和"重试也白搭"

无脑重试是最常见的错误。先分类错误:

type FailureKind = 'retryable' | 'fatal' | 'poison'

function classifyFailure(err: unknown): FailureKind {
  if (err instanceof RateLimitError) return 'retryable'     // 429,等一等就好
  if (err instanceof TimeoutError) return 'retryable'       // 超时,可能服务端已成功——需幂等
  if (err instanceof ServerError && err.status >= 500) return 'retryable'
  if (err instanceof QuotaExhaustedError) return 'fatal'    // 配额耗尽,重试无意义
  if (err instanceof InvalidRequestError) return 'fatal'    // 请求本身错,重试还是错
  if (err instanceof SchemaViolationError && err.attempts >= 3) return 'poison'  // 模型持续输出非法
  return 'fatal'
}

poison(毒丸)类值得单独说:LLM 对某个输入持续输出非法结构,重试 3 次仍失败——再重试只会烧钱。正确动作是标记输入并跳过,同时降级到规则引擎(见结构化输出那篇的降级链)。

指数退避 + 抖动:

function backoffDelay(attempt: number, baseMs = 1000, maxMs = 60_000): number {
  const exponential = Math.min(baseMs * 2 ** attempt, maxMs)
  // 全抖动(full jitter):在 [0, exponential) 内随机,打散重试风暴
  return Math.random() * exponential
}

// 尊重 Retry-After 头:429 时服务端说了算
async function withRetry<T>(fn: () => Promise<T>, maxAttempts = 4): Promise<T> {
  let lastError: unknown
  for (let attempt = 0; attempt < maxAttempts; attempt++) {
    try {
      return await fn()
    } catch (err) {
      lastError = err
      if (classifyFailure(err) !== 'retryable') throw err
      if (err instanceof RateLimitError && err.retryAfterMs) {
        await sleep(err.retryAfterMs)      // 服务端指示优先
        continue
      }
      await sleep(backoffDelay(attempt))
    }
  }
  throw lastError
}

Full jitter 的意义:100 个客户端同时失败,无抖动的重试会整齐地再次撞车(重试风暴),抖动把它们摊开。退避参数(base/上限/抖动)属于全局架构资产,不要每个调用点各写一套。

四、幂等:超时的正确处理方式

超时的请求,服务端可能已成功执行(响应丢了)。直接重试 = 可能重复扣费。幂等键:

async function runIdempotent<T>(op: { kind: string; key: string; payload: unknown }, fn: (ctx: { signal: AbortSignal }) => Promise<T>): Promise<T> {
  // 1. 先查:同 key 的任务若已成功,直接返回缓存结果(deduped)
  const existing = await taskStore.findByKey(op.key)
  if (existing?.status === 'succeeded') {
    return existing.result as T
  }

  // 2. 占位:原子写入 running 状态(冲突 = 并发重复提交)
  const claimed = await taskStore.claim(op.key, op)
  if (!claimed) {
    // 已有同 key 任务在跑:订阅它的事件流等结果,而非再起一个
    return await waitForExisting(op.key)
  }

  // 3. 执行 + 结果落库(同一事务内写 status + result)
  try {
    const result = await fn({ signal: claimed.signal })
    await taskStore.complete(op.key, result)
    return result
  } catch (err) {
    await taskStore.fail(op.key, err)
    throw err
  }
}

// 幂等键构造:输入指纹,而非自增 ID
function ttsIdempotencyKey(input: TtsInput): string {
  return sha256(`${input.engineVersion}:${input.voiceId}:${input.speed}:${input.text}`)
}

三个工程细节:

  1. 幂等键是输入指纹sha256(引擎版本+参数+文本),同输入天然同 key。自增 ID 当幂等键是自欺欺人——重复提交时 ID 都不同。
  2. claim 要原子claim 失败说明并发重复,走"等已有任务"而不是再跑一遍。分布式环境下用 INSERT ... ON CONFLICT DO NOTHING 或 Redis SETNX。
  3. 结果与状态同事务:先写 result 再改 status(或同事务),崩溃恢复时不会出现"状态 succeeded 但 result 为空"的僵尸任务。

五、持久化:重启不丢任务

桌面应用的任务队列必须持久化到磁盘——用户导出到一半强杀进程,重开时要么续跑要么明确告知失败,不能"消失"。轻量方案:JSON 状态文件 + 原子写:

class PersistentTaskStore {
  private file: string

  async persist(tasks: Task[]) {
    // 原子写:临时文件 + rename(避免写一半崩溃留下半个 JSON)
    const tmp = `${this.file}.tmp`
    await fs.writeFile(tmp, JSON.stringify({ version: 1, tasks }, null, 2))
    await fs.rename(tmp, this.file)
  }

  async recoverOrphans(): Promise<void> {
    const tasks = await this.load()
    for (const task of tasks) {
      if (task.status === 'running') {
        // 上次进程崩溃时正在跑的任务:状态不可信
        // 结果可能已产出(幂等键查询外部服务/缓存),先尝试对账
        const reconciled = await tryReconcile(task)
        await this.update(task.id, reconciled ? 'succeeded' : 'failed',
          reconciled ? undefined : new InterruptedError())
      }
    }
  }
}

recoverOrphans 的"对账"是关键一步:崩溃前任务可能已在外部服务成功执行。靠幂等键去外部缓存/任务存储查询,查到了就认账(succeeded),查不到才标 failed。状态机的持久化恢复必须假设任何时刻都可能死——包括"写状态文件写到一半"(所以原子写)和"外部调用成功但本地没来得及记录"(所以对账)。

六、优先级与抢占

用户感知的优先级与提交顺序不一致:用户正在编辑第 2 段口播稿,此时点"合成"——这个 TTS 任务的优先级高于后台排队的"全部视频重新分析"。两级队列:

class PriorityTaskQueue {
  private interactive: Task[] = []   // 用户当前操作直接触发
  private background: Task[] = []    // 批处理、预取、缓存预热

  enqueue(task: Task) {
    (task.interactive ? this.interactive : this.background).push(task)
  }

  next(): Task | null {
    return this.interactive.shift() ?? this.background.shift() ?? null
  }
}

不实现抢占(running 任务中断重排)——成本高且 TTS/LLM 半途中断纯浪费。优先级只影响出队顺序,简单且够用。interactive 队列还应该限制深度(如 5),超过提示用户"操作太快"——防连点风暴。

七、小结

机制核心实现
状态机6 状态 + 合法转移表,非法转移抛错
限流令牌桶按真实成本扣减;配额感知调度
重试错误三分类(retryable/fatal/poison);full jitter 退避;尊重 Retry-After
幂等输入指纹做 key;原子 claim;结果与状态同事务
持久化原子写状态文件;恢复时先对账再定生死
优先级interactive/background 双队列,只影响出队不抢占

调度系统的价值在于:用户完全感知不到它的存在。限流让你平滑、重试让你可靠、幂等让你省钱、持久化让你可靠、优先级让你贴心——每一个都是"出问题时才显灵"的隐形基础设施。AI 应用与传统应用在工程上的差距,往往不在模型调用那 10 行代码,在这 600 行调度代码。

相关文章

精彩推荐