AI 后台任务的问题通常不是“能否执行”,而是失败后能否解释任务在哪里、由谁持有以及何时恢复。本文用 Node.js 内置模块实现一个最小队列,完整覆盖入队、租约、心跳、超时回收、有限重试、死信和运维查询。

黑盒从哪里产生

一次 AI 任务可能包含模型调用、文件解析、向量化和结果持久化,执行时间从几秒到数分钟不等。若队列只有“待处理”和“已完成”两个标记,消费者崩溃后,系统很难判断任务仍在执行,还是已经永久失联。

常见故障来自三个边界:

  • 消费者先确认任务、再执行业务,进程中途退出会造成任务丢失。
  • 消费者执行完业务、尚未确认时退出,任务会被再次投递,形成重复执行。
  • 消费者取走任务后没有明确的失效时间,异常任务会永久占用运行状态。

因此,可靠队列通常提供的是“至少一次”投递,而不是天然的“恰好一次”。租约解决永久占用,重试解决暂时失败,幂等性负责吸收重复执行。三者不能互相替代。

状态机和关键不变量

最小实现采用四种状态,并记录每次流转事件:

状态含义允许的后续状态
queued等待消费者领取running
running已被带租约地领取succeededqueueddead
succeeded执行成功,终态
dead达到重试上限,等待人工处理queued

消费者领取任务时生成不可预测的 leaseToken,同时写入 leaseExpiresAt。心跳只有携带当前令牌才能续租,完成和失败操作也必须校验令牌。这样,旧消费者即使在租约过期后恢复,也不能覆盖新消费者的结果。

还需要坚持几个不变量:attempts 在每次领取时递增;只有 running 任务能够续租或完成;扫描器只回收已经过期的租约;达到 maxAttempts 后进入死信,不能无限重试。实际业务仍应使用任务 ID 或业务唯一键实现幂等,例如对结果表建立唯一约束。

一个可运行的 Node.js 最小实现

下面的示例只依赖 Node.js 内置的 node:httpnode:crypto,可直接保存为 app.js,使用 Node.js 18 或更高版本运行。内存存储便于看清协议,不适合跨进程共享或持久化生产数据。

const http = require('node:http');
const { randomUUID } = require('node:crypto');

class JobQueue {
  constructor({ leaseMs = 5000, maxAttempts = 3 } = {}) {
    this.jobs = new Map();
    this.leaseMs = leaseMs;
    this.maxAttempts = maxAttempts;
  }

  enqueue(payload) {
    const now = Date.now();
    const job = {
      id: randomUUID(), payload, state: 'queued', attempts: 0,
      availableAt: now, leaseToken: null, leaseExpiresAt: null,
      workerId: null, lastError: null, createdAt: now, updatedAt: now,
      events: [{ type: 'enqueued', at: now }]
    };
    this.jobs.set(job.id, job);
    return this.view(job);
  }

  lease(workerId) {
    const now = Date.now();
    const job = [...this.jobs.values()].find(
      item => item.state === 'queued' && item.availableAt <= now
    );
    if (!job) return null;

    job.state = 'running';
    job.attempts += 1;
    job.workerId = workerId;
    job.leaseToken = randomUUID();
    job.leaseExpiresAt = now + this.leaseMs;
    job.updatedAt = now;
    job.events.push({ type: 'leased', at: now, workerId, attempt: job.attempts });
    return { ...this.view(job), leaseToken: job.leaseToken };
  }

  heartbeat(id, token) {
    const job = this.jobs.get(id);
    if (!this.owns(job, token)) return false;
    const now = Date.now();
    job.leaseExpiresAt = now + this.leaseMs;
    job.updatedAt = now;
    job.events.push({ type: 'heartbeat', at: now });
    return true;
  }

  complete(id, token) {
    const job = this.jobs.get(id);
    if (!this.owns(job, token)) return false;
    this.release(job, 'succeeded', null);
    return true;
  }

  fail(id, token, error) {
    const job = this.jobs.get(id);
    if (!this.owns(job, token)) return false;
    const message = String(error?.message || error);
    job.lastError = message;
    if (job.attempts >= this.maxAttempts) {
      this.release(job, 'dead', message);
    } else {
      this.release(job, 'queued', message);
      job.availableAt = Date.now() + 1000 * job.attempts;
    }
    return true;
  }

  reapExpired() {
    const now = Date.now();
    let count = 0;
    for (const job of this.jobs.values()) {
      if (job.state !== 'running' || job.leaseExpiresAt > now) continue;
      job.lastError = 'lease expired';
      this.release(job, job.attempts >= this.maxAttempts ? 'dead' : 'queued', 'lease expired');
      job.availableAt = now;
      count += 1;
    }
    return count;
  }

  retryDead(id) {
    const job = this.jobs.get(id);
    if (!job || job.state !== 'dead') return false;
    job.state = 'queued';
    job.attempts = 0;
    job.availableAt = Date.now();
    job.updatedAt = Date.now();
    job.events.push({ type: 'manual_retry', at: job.updatedAt });
    return true;
  }

  owns(job, token) {
    return Boolean(job && job.state === 'running' && job.leaseToken === token);
  }

  release(job, state, detail) {
    const now = Date.now();
    job.state = state;
    job.workerId = null;
    job.leaseToken = null;
    job.leaseExpiresAt = null;
    job.updatedAt = now;
    job.events.push({ type: state, at: now, detail });
  }

  view(job) {
    const { leaseToken, ...safe } = job;
    return structuredClone(safe);
  }

  list(state) {
    return [...this.jobs.values()]
      .filter(job => !state || job.state === state)
      .map(job => this.view(job));
  }
}

const queue = new JobQueue();
let busy = false;

async function runWorker() {
  if (busy) return;
  const job = queue.lease('worker-1');
  if (!job) return;
  busy = true;
  const heartbeat = setInterval(
    () => queue.heartbeat(job.id, job.leaseToken),
    1000
  );

  try {
    await new Promise(resolve => setTimeout(resolve, job.payload.durationMs || 2000));
    if (job.payload.fail) throw new Error('simulated model failure');
    queue.complete(job.id, job.leaseToken);
  } catch (error) {
    queue.fail(job.id, job.leaseToken, error);
  } finally {
    clearInterval(heartbeat);
    busy = false;
  }
}

function readJson(req) {
  return new Promise((resolve, reject) => {
    let body = '';
    req.on('data', chunk => {
      body += chunk;
      if (body.length > 1_000_000) req.destroy();
    });
    req.on('end', () => {
      try { resolve(body ? JSON.parse(body) : {}); }
      catch (error) { reject(error); }
    });
    req.on('error', reject);
  });
}

function send(res, status, data) {
  res.writeHead(status, { 'content-type': 'application/json; charset=utf-8' });
  res.end(JSON.stringify(data));
}

const server = http.createServer(async (req, res) => {
  const url = new URL(req.url, 'http://localhost');
  try {
    if (req.method === 'POST' && url.pathname === '/jobs') {
      return send(res, 201, queue.enqueue(await readJson(req)));
    }
    if (req.method === 'GET' && url.pathname === '/jobs') {
      return send(res, 200, queue.list(url.searchParams.get('state')));
    }
    if (req.method === 'GET' && url.pathname === '/dead') {
      return send(res, 200, queue.list('dead'));
    }
    const match = url.pathname.match(/^\/jobs\/([^/]+)\/retry$/);
    if (req.method === 'POST' && match) {
      return send(res, queue.retryDead(match[1]) ? 202 : 409, { accepted: true });
    }
    send(res, 404, { error: 'not found' });
  } catch (error) {
    send(res, 400, { error: error.message });
  }
});

setInterval(runWorker, 500);
setInterval(() => queue.reapExpired(), 1000);
server.listen(3000, () => console.log('listening on http://localhost:3000'));

验证故障与运维查询

启动服务后,先提交一个正常任务和一个持续失败的任务:

node app.js
curl -X POST http://localhost:3000/jobs \
  -H 'content-type: application/json' \
  -d '{"durationMs":3000}'
curl -X POST http://localhost:3000/jobs \
  -H 'content-type: application/json' \
  -d '{"durationMs":1000,"fail":true}'

通过 GET /jobs 可以查看所有任务及事件历史,GET /jobs?state=running 用于定位正在持有租约的任务,GET /dead 用于检查达到重试上限的任务。死信确认可以重放后,调用 POST /jobs/{id}/retry 将其重新入队。

观察事件时应重点关注 attemptsworkerIdleaseExpiresAtlastError。如果大量任务持续续租,可能是下游模型响应变慢;如果 lease expired 突增,应检查消费者崩溃、网络隔离或事件循环阻塞;如果死信集中出现相同错误,则更可能是输入、权限或模型配置问题,而非短暂抖动。

可以在任务运行期间终止进程来理解租约边界,但这个内存版本重启后数据会消失。生产系统必须把任务和租约放入持久化存储,否则进程崩溃仍会造成任务丢失。

从最小版走向生产

把实现迁移到 PostgreSQL、Redis 或成熟消息队列时,应保留状态机语义,而不是照搬内存数据结构。领取任务需要原子操作:数据库方案可在事务中使用行锁和条件更新;完成任务时应附带任务 ID 与租约令牌,避免过期消费者提交结果。

租约时间不能简单设得越长越好。过短会因网络抖动产生误回收,过长则延迟故障恢复。可根据任务类型设置不同租约,并让心跳间隔明显短于租约时间。对于无法安全重试的外部操作,例如计费或发送通知,应使用幂等键、唯一约束或 outbox 模式协调业务结果与任务状态。

可观测性也不应只依赖任务列表。至少应采集排队时长、执行时长、重试次数、租约过期数、死信数和各状态任务量,并为长期 running、死信增长和队列积压设置告警。事件历史需要保留时间、消费者和错误摘要,同时避免记录提示词、密钥等敏感内容。

总结

一个可恢复的 AI 任务队列,核心不是更多状态,而是明确状态转换的条件:领取时创建有限租约,执行时通过心跳续租,完成时校验持有者,超时后自动回收,超过上限后进入死信。

要点可以归纳为:接受至少一次投递并让业务幂等;用租约令牌阻止旧消费者写入;为重试设置上限和退避;提供状态、事件与死信查询;在生产环境使用支持原子更新的持久化存储。做到这些,后台任务即使失败,也能被定位、解释和恢复。